1use std::collections::HashMap;
18use std::future::Future;
19use std::mem;
20use std::pin::Pin;
21use std::sync::Arc;
22use std::sync::atomic::{AtomicUsize, Ordering};
23use std::task::{Context, Poll};
24use std::time::Instant;
25
26use common_telemetry::debug;
27use common_time::Timestamp;
28use datatypes::arrow::array::{
29 ArrayRef, TimestampMicrosecondArray, TimestampMillisecondArray, TimestampNanosecondArray,
30 TimestampSecondArray,
31};
32use datatypes::arrow::compute::{max, min};
33use datatypes::arrow::datatypes::{DataType, SchemaRef, TimeUnit};
34use datatypes::arrow::record_batch::RecordBatch;
35use datatypes::extension::json::is_json2_extension_type;
36use object_store::{FuturesAsyncWriter, ObjectStore};
37use parquet::arrow::AsyncArrowWriter;
38use parquet::basic::{Compression, Encoding, ZstdLevel};
39use parquet::file::metadata::KeyValue;
40use parquet::file::properties::{WriterProperties, WriterPropertiesBuilder};
41use parquet::schema::types::ColumnPath;
42use smallvec::smallvec;
43use snafu::ResultExt;
44use store_api::metadata::RegionMetadataRef;
45use store_api::storage::consts::{OP_TYPE_COLUMN_NAME, SEQUENCE_COLUMN_NAME};
46use store_api::storage::{FileId, SequenceNumber};
47use tokio::io::AsyncWrite;
48use tokio_util::compat::{Compat, FuturesAsyncWriteCompatExt};
49
50use crate::access_layer::{FilePathProvider, Metrics, SstInfoArray, TempFileCleaner};
51use crate::config::{IndexBuildMode, IndexConfig};
52use crate::error::{
53 InvalidMetadataSnafu, OpenDalSnafu, Result, UnexpectedSnafu, WriteParquetSnafu,
54};
55use crate::read::FlatSource;
56use crate::sst::file::RegionFileId;
57use crate::sst::index::{IndexOutput, Indexer, IndexerBuilder};
58use crate::sst::parquet::flat_format::{FlatWriteFormat, time_index_column_index};
59use crate::sst::parquet::format::PrimaryKeyWriteFormat;
60use crate::sst::parquet::{PARQUET_METADATA_KEY, SstInfo, WriteOptions};
61use crate::sst::{
62 DEFAULT_WRITE_BUFFER_SIZE, DEFAULT_WRITE_CONCURRENCY, FlatSchemaOptions, SeriesEstimator,
63 maybe_wrap_schema,
64};
65
66enum FlatBatchConverter {
68 Flat(FlatWriteFormat),
70 PrimaryKey {
72 format: PrimaryKeyWriteFormat,
73 num_fields: usize,
74 },
75}
76
77impl FlatBatchConverter {
78 fn convert_batch(&self, batch: &RecordBatch) -> Result<RecordBatch> {
79 match self {
80 FlatBatchConverter::Flat(f) => f.convert_batch(batch),
81 FlatBatchConverter::PrimaryKey { format, num_fields } => {
82 format.convert_flat_batch(batch, *num_fields)
83 }
84 }
85 }
86}
87
88pub struct ParquetWriter<'a, F: WriterFactory, I: IndexerBuilder, P: FilePathProvider> {
90 path_provider: P,
92 writer: Option<AsyncArrowWriter<SizeAwareWriter<F::Writer>>>,
93 current_file: FileId,
95 writer_factory: F,
96 metadata: RegionMetadataRef,
98 index_config: IndexConfig,
100 indexer_builder: I,
102 current_indexer: Option<Indexer>,
104 bytes_written: Arc<AtomicUsize>,
105 file_cleaner: Option<TempFileCleaner>,
107 metrics: &'a mut Metrics,
109}
110
111pub trait WriterFactory {
112 type Writer: AsyncWrite + Send + Unpin;
113 fn create(&mut self, file_path: &str) -> impl Future<Output = Result<Self::Writer>>;
114}
115
116pub struct ObjectStoreWriterFactory {
117 object_store: ObjectStore,
118}
119
120impl WriterFactory for ObjectStoreWriterFactory {
121 type Writer = Compat<FuturesAsyncWriter>;
122
123 async fn create(&mut self, file_path: &str) -> Result<Self::Writer> {
124 self.object_store
125 .writer_with(file_path)
126 .chunk(DEFAULT_WRITE_BUFFER_SIZE.as_bytes() as usize)
127 .concurrent(DEFAULT_WRITE_CONCURRENCY)
128 .await
129 .map(|v| v.into_futures_async_write().compat_write())
130 .context(OpenDalSnafu)
131 }
132}
133
134impl<'a, I, P> ParquetWriter<'a, ObjectStoreWriterFactory, I, P>
135where
136 P: FilePathProvider,
137 I: IndexerBuilder,
138{
139 pub async fn new_with_object_store(
140 object_store: ObjectStore,
141 metadata: RegionMetadataRef,
142 index_config: IndexConfig,
143 indexer_builder: I,
144 path_provider: P,
145 metrics: &'a mut Metrics,
146 ) -> ParquetWriter<'a, ObjectStoreWriterFactory, I, P> {
147 ParquetWriter::new(
148 ObjectStoreWriterFactory { object_store },
149 metadata,
150 index_config,
151 indexer_builder,
152 path_provider,
153 metrics,
154 )
155 .await
156 }
157
158 pub(crate) fn with_file_cleaner(mut self, cleaner: TempFileCleaner) -> Self {
159 self.file_cleaner = Some(cleaner);
160 self
161 }
162}
163
164impl<'a, F, I, P> ParquetWriter<'a, F, I, P>
165where
166 F: WriterFactory,
167 I: IndexerBuilder,
168 P: FilePathProvider,
169{
170 pub async fn new(
172 factory: F,
173 metadata: RegionMetadataRef,
174 index_config: IndexConfig,
175 indexer_builder: I,
176 path_provider: P,
177 metrics: &'a mut Metrics,
178 ) -> ParquetWriter<'a, F, I, P> {
179 let init_file = FileId::random();
180 let indexer = indexer_builder
181 .build(RegionFileId::new(metadata.region_id, init_file), 0, None)
182 .await;
183
184 ParquetWriter {
185 path_provider,
186 writer: None,
187 current_file: init_file,
188 writer_factory: factory,
189 metadata,
190 index_config,
191 indexer_builder,
192 current_indexer: Some(indexer),
193 bytes_written: Arc::new(AtomicUsize::new(0)),
194 file_cleaner: None,
195 metrics,
196 }
197 }
198
199 async fn finish_current_file(
201 &mut self,
202 ssts: &mut SstInfoArray,
203 stats: &mut SourceStats,
204 ) -> Result<()> {
205 if let Some(mut current_writer) = mem::take(&mut self.writer) {
207 let mut stats = mem::take(stats);
208 assert!(stats.num_rows > 0);
210
211 debug!(
212 "Finishing current file {}, file size: {}, num rows: {}",
213 self.current_file,
214 self.bytes_written.load(Ordering::Relaxed),
215 stats.num_rows
216 );
217
218 let mut index_output = IndexOutput::default();
221 match self.index_config.build_mode {
222 IndexBuildMode::Sync => {
223 index_output = self.current_indexer.as_mut().unwrap().finish().await;
224 }
225 IndexBuildMode::Async => {
226 debug!(
227 "Index for file {} will be built asynchronously later",
228 self.current_file
229 );
230 }
231 }
232 current_writer.flush().await.context(WriteParquetSnafu)?;
233
234 let parquet_metadata = current_writer.close().await.context(WriteParquetSnafu)?;
235 let file_size = self.bytes_written.load(Ordering::Relaxed) as u64;
236
237 let time_range = stats.time_range.unwrap();
239
240 let max_row_group_uncompressed_size: u64 = parquet_metadata
241 .row_groups()
242 .iter()
243 .map(|rg| {
244 rg.columns()
245 .iter()
246 .map(|c| c.uncompressed_size() as u64)
247 .sum::<u64>()
248 })
249 .max()
250 .unwrap_or(0);
251 let num_series = stats.series_estimator.finish();
252 ssts.push(SstInfo {
253 file_id: self.current_file,
254 time_range,
255 file_size,
256 max_row_group_uncompressed_size,
257 num_rows: stats.num_rows,
258 num_row_groups: parquet_metadata.num_row_groups() as u64,
259 file_metadata: Some(Arc::new(parquet_metadata)),
260 index_metadata: index_output,
261 num_series,
262 });
263 self.current_file = FileId::random();
264 self.bytes_written.store(0, Ordering::Relaxed)
265 };
266
267 Ok(())
268 }
269
270 pub async fn write_all_flat(
274 &mut self,
275 source: FlatSource,
276 override_sequence: Option<SequenceNumber>,
277 opts: &WriteOptions,
278 ) -> Result<SstInfoArray> {
279 let mut options = FlatSchemaOptions::from_encoding(self.metadata.primary_key_encoding);
280
281 if source.schema().fields().iter().any(is_json2_extension_type) {
282 options.concretized_json_types = source
283 .schema()
284 .fields()
285 .iter()
286 .filter(|&field| is_json2_extension_type(field))
287 .map(|field| (field.name().clone(), field.data_type().clone()))
288 .collect::<HashMap<_, _>>();
289 }
290
291 let converter = FlatBatchConverter::Flat(
292 FlatWriteFormat::new(self.metadata.clone(), &options)
293 .with_override_sequence(override_sequence),
294 );
295 let res = self.write_all_flat_inner(source, &converter, opts).await;
296 if res.is_err() {
297 let file_id = self.current_file;
298 if let Some(cleaner) = &self.file_cleaner {
299 cleaner.clean_by_file_id(file_id).await;
300 }
301 }
302 res
303 }
304
305 pub async fn write_all_flat_as_primary_key(
309 &mut self,
310 source: FlatSource,
311 override_sequence: Option<SequenceNumber>,
312 opts: &WriteOptions,
313 ) -> Result<SstInfoArray> {
314 let num_fields = self.metadata.field_columns().count();
315 let converter = FlatBatchConverter::PrimaryKey {
316 format: PrimaryKeyWriteFormat::new(self.metadata.clone())
317 .with_override_sequence(override_sequence),
318 num_fields,
319 };
320 let res = self.write_all_flat_inner(source, &converter, opts).await;
321 if res.is_err() {
322 let file_id = self.current_file;
323 if let Some(cleaner) = &self.file_cleaner {
324 cleaner.clean_by_file_id(file_id).await;
325 }
326 }
327 res
328 }
329
330 async fn write_all_flat_inner(
331 &mut self,
332 mut source: FlatSource,
333 converter: &FlatBatchConverter,
334 opts: &WriteOptions,
335 ) -> Result<SstInfoArray> {
336 let mut results = smallvec![];
337 let mut stats = SourceStats::default();
338
339 while let Some(record_batch) = self
340 .write_next_flat_batch(&mut source, converter, opts)
341 .await
342 .transpose()
343 {
344 match record_batch {
345 Ok(batch) => {
346 stats.update_flat(&batch)?;
347 if matches!(self.index_config.build_mode, IndexBuildMode::Sync) {
348 let start = Instant::now();
349 self.current_indexer
351 .as_mut()
352 .unwrap()
353 .update_flat(&batch)
354 .await;
355 self.metrics.update_index += start.elapsed();
356 }
357 if let Some(max_file_size) = opts.max_file_size
358 && self.bytes_written.load(Ordering::Relaxed) > max_file_size
359 {
360 self.finish_current_file(&mut results, &mut stats).await?;
361 }
362 }
363 Err(e) => {
364 if let Some(indexer) = &mut self.current_indexer {
365 indexer.abort().await;
366 }
367 return Err(e);
368 }
369 }
370 }
371
372 self.finish_current_file(&mut results, &mut stats).await?;
373
374 Ok(results)
376 }
377
378 fn customize_column_config(
380 builder: WriterPropertiesBuilder,
381 region_metadata: &RegionMetadataRef,
382 ) -> WriterPropertiesBuilder {
383 let ts_col = ColumnPath::new(vec![
384 region_metadata
385 .time_index_column()
386 .column_schema
387 .name
388 .clone(),
389 ]);
390 let seq_col = ColumnPath::new(vec![SEQUENCE_COLUMN_NAME.to_string()]);
391 let op_type_col = ColumnPath::new(vec![OP_TYPE_COLUMN_NAME.to_string()]);
392
393 builder
394 .set_column_encoding(seq_col.clone(), Encoding::DELTA_BINARY_PACKED)
395 .set_column_dictionary_enabled(seq_col, false)
396 .set_column_encoding(ts_col.clone(), Encoding::DELTA_BINARY_PACKED)
397 .set_column_dictionary_enabled(ts_col, false)
398 .set_column_compression(op_type_col, Compression::UNCOMPRESSED)
399 }
400
401 async fn write_next_flat_batch(
402 &mut self,
403 source: &mut FlatSource,
404 converter: &FlatBatchConverter,
405 opts: &WriteOptions,
406 ) -> Result<Option<RecordBatch>> {
407 let start = Instant::now();
408 let Some(record_batch) = source.next_batch().await? else {
409 return Ok(None);
410 };
411 self.metrics.iter_source += start.elapsed();
412
413 let arrow_batch = converter.convert_batch(&record_batch)?;
414
415 let start = Instant::now();
416 self.maybe_init_writer(arrow_batch.schema_ref(), opts)
417 .await?
418 .write(&arrow_batch)
419 .await
420 .context(WriteParquetSnafu)?;
421 self.metrics.write_batch += start.elapsed();
422 Ok(Some(record_batch))
424 }
425
426 async fn maybe_init_writer(
427 &mut self,
428 schema: &SchemaRef,
429 opts: &WriteOptions,
430 ) -> Result<&mut AsyncArrowWriter<SizeAwareWriter<F::Writer>>> {
431 if let Some(ref mut w) = self.writer {
432 Ok(w)
433 } else {
434 let json = self.metadata.to_json().context(InvalidMetadataSnafu)?;
435 let key_value_meta = KeyValue::new(PARQUET_METADATA_KEY.to_string(), json);
436
437 let props_builder = WriterProperties::builder()
439 .set_key_value_metadata(Some(vec![key_value_meta]))
440 .set_compression(Compression::ZSTD(ZstdLevel::default()))
441 .set_encoding(Encoding::PLAIN)
442 .set_max_row_group_row_count(Some(opts.row_group_size))
443 .set_column_index_truncate_length(None)
444 .set_statistics_truncate_length(None);
445
446 let props_builder = Self::customize_column_config(props_builder, &self.metadata);
447 let writer_props = props_builder.build();
448
449 let sst_file_path = self.path_provider.build_sst_file_path(RegionFileId::new(
450 self.metadata.region_id,
451 self.current_file,
452 ));
453 let writer = SizeAwareWriter::new(
454 self.writer_factory.create(&sst_file_path).await?,
455 self.bytes_written.clone(),
456 );
457 let arrow_writer =
458 AsyncArrowWriter::try_new(writer, maybe_wrap_schema(schema)?, Some(writer_props))
459 .context(WriteParquetSnafu)?;
460 self.writer = Some(arrow_writer);
461
462 let indexer = self
463 .indexer_builder
464 .build(
465 RegionFileId::new(self.metadata.region_id, self.current_file),
466 0,
467 Some(opts.row_group_size),
468 )
469 .await;
470 self.current_indexer = Some(indexer);
471
472 Ok(self.writer.as_mut().unwrap())
474 }
475 }
476}
477
478#[derive(Default)]
479struct SourceStats {
480 num_rows: usize,
482 time_range: Option<(Timestamp, Timestamp)>,
484 series_estimator: SeriesEstimator,
486}
487
488impl SourceStats {
489 fn update_flat(&mut self, record_batch: &RecordBatch) -> Result<()> {
490 if record_batch.num_rows() == 0 {
491 return Ok(());
492 }
493
494 self.num_rows += record_batch.num_rows();
495 self.series_estimator.update_flat(record_batch);
496
497 let time_index_col_idx = time_index_column_index(record_batch.num_columns());
499 let timestamp_array = record_batch.column(time_index_col_idx);
500
501 if let Some((min_in_batch, max_in_batch)) = timestamp_range_from_array(timestamp_array)? {
502 if let Some(time_range) = &mut self.time_range {
503 time_range.0 = time_range.0.min(min_in_batch);
504 time_range.1 = time_range.1.max(max_in_batch);
505 } else {
506 self.time_range = Some((min_in_batch, max_in_batch));
507 }
508 }
509
510 Ok(())
511 }
512}
513
514fn timestamp_range_from_array(
516 timestamp_array: &ArrayRef,
517) -> Result<Option<(Timestamp, Timestamp)>> {
518 let (min_ts, max_ts) = match timestamp_array.data_type() {
519 DataType::Timestamp(TimeUnit::Second, _) => {
520 let array = timestamp_array
521 .as_any()
522 .downcast_ref::<TimestampSecondArray>()
523 .unwrap();
524 let min_val = min(array).map(Timestamp::new_second);
525 let max_val = max(array).map(Timestamp::new_second);
526 (min_val, max_val)
527 }
528 DataType::Timestamp(TimeUnit::Millisecond, _) => {
529 let array = timestamp_array
530 .as_any()
531 .downcast_ref::<TimestampMillisecondArray>()
532 .unwrap();
533 let min_val = min(array).map(Timestamp::new_millisecond);
534 let max_val = max(array).map(Timestamp::new_millisecond);
535 (min_val, max_val)
536 }
537 DataType::Timestamp(TimeUnit::Microsecond, _) => {
538 let array = timestamp_array
539 .as_any()
540 .downcast_ref::<TimestampMicrosecondArray>()
541 .unwrap();
542 let min_val = min(array).map(Timestamp::new_microsecond);
543 let max_val = max(array).map(Timestamp::new_microsecond);
544 (min_val, max_val)
545 }
546 DataType::Timestamp(TimeUnit::Nanosecond, _) => {
547 let array = timestamp_array
548 .as_any()
549 .downcast_ref::<TimestampNanosecondArray>()
550 .unwrap();
551 let min_val = min(array).map(Timestamp::new_nanosecond);
552 let max_val = max(array).map(Timestamp::new_nanosecond);
553 (min_val, max_val)
554 }
555 _ => {
556 return UnexpectedSnafu {
557 reason: format!(
558 "Unexpected data type of time index: {:?}",
559 timestamp_array.data_type()
560 ),
561 }
562 .fail();
563 }
564 };
565
566 Ok(min_ts.zip(max_ts))
568}
569
570struct SizeAwareWriter<W> {
573 inner: W,
574 size: Arc<AtomicUsize>,
575}
576
577impl<W> SizeAwareWriter<W> {
578 fn new(inner: W, size: Arc<AtomicUsize>) -> Self {
579 Self {
580 inner,
581 size: size.clone(),
582 }
583 }
584}
585
586impl<W> AsyncWrite for SizeAwareWriter<W>
587where
588 W: AsyncWrite + Unpin,
589{
590 fn poll_write(
591 mut self: Pin<&mut Self>,
592 cx: &mut Context<'_>,
593 buf: &[u8],
594 ) -> Poll<std::result::Result<usize, std::io::Error>> {
595 let this = self.as_mut().get_mut();
596
597 match Pin::new(&mut this.inner).poll_write(cx, buf) {
598 Poll::Ready(Ok(bytes_written)) => {
599 this.size.fetch_add(bytes_written, Ordering::Relaxed);
600 Poll::Ready(Ok(bytes_written))
601 }
602 other => other,
603 }
604 }
605
606 fn poll_flush(
607 mut self: Pin<&mut Self>,
608 cx: &mut Context<'_>,
609 ) -> Poll<std::result::Result<(), std::io::Error>> {
610 Pin::new(&mut self.inner).poll_flush(cx)
611 }
612
613 fn poll_shutdown(
614 mut self: Pin<&mut Self>,
615 cx: &mut Context<'_>,
616 ) -> Poll<std::result::Result<(), std::io::Error>> {
617 Pin::new(&mut self.inner).poll_shutdown(cx)
618 }
619}