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 bytes::Bytes;
27use common_telemetry::debug;
28use common_time::Timestamp;
29use datatypes::arrow::array::{
30 ArrayRef, BinaryArray, TimestampMicrosecondArray, TimestampMillisecondArray,
31 TimestampNanosecondArray, TimestampSecondArray, UInt32Array,
32};
33use datatypes::arrow::compute::{max, min};
34use datatypes::arrow::datatypes::{DataType, SchemaRef, TimeUnit};
35use datatypes::arrow::record_batch::RecordBatch;
36use datatypes::extension::json::is_json2_extension_type;
37use object_store::{FuturesAsyncWriter, ObjectStore};
38use parquet::arrow::AsyncArrowWriter;
39use parquet::basic::{Compression, Encoding, ZstdLevel};
40use parquet::file::metadata::KeyValue;
41use parquet::file::properties::{WriterProperties, WriterPropertiesBuilder};
42use parquet::schema::types::ColumnPath;
43use smallvec::smallvec;
44use snafu::{OptionExt, ResultExt};
45use store_api::metadata::RegionMetadataRef;
46use store_api::storage::consts::{OP_TYPE_COLUMN_NAME, SEQUENCE_COLUMN_NAME};
47use store_api::storage::{FileId, SequenceNumber};
48use tokio::io::AsyncWrite;
49use tokio_util::compat::{Compat, FuturesAsyncWriteCompatExt};
50
51use crate::access_layer::{FilePathProvider, Metrics, SstInfoArray, TempFileCleaner};
52use crate::config::{IndexBuildMode, IndexConfig};
53use crate::error::{
54 InvalidMetadataSnafu, InvalidRecordBatchSnafu, OpenDalSnafu, Result, UnexpectedSnafu,
55 WriteParquetSnafu,
56};
57use crate::read::FlatSource;
58use crate::sst::file::RegionFileId;
59use crate::sst::index::{IndexOutput, Indexer, IndexerBuilder};
60use crate::sst::parquet::flat_format::{
61 FlatWriteFormat, primary_key_column_index, time_index_column_index,
62};
63use crate::sst::parquet::format::{PrimaryKeyArray, PrimaryKeyWriteFormat};
64use crate::sst::parquet::{PARQUET_METADATA_KEY, SstInfo, WriteOptions};
65use crate::sst::{
66 DEFAULT_WRITE_BUFFER_SIZE, DEFAULT_WRITE_CONCURRENCY, FlatSchemaOptions, SeriesEstimator,
67 maybe_wrap_schema,
68};
69
70enum FlatBatchConverter {
72 Flat(FlatWriteFormat),
74 PrimaryKey {
76 format: PrimaryKeyWriteFormat,
77 num_fields: usize,
78 },
79}
80
81impl FlatBatchConverter {
82 fn convert_batch(&self, batch: &RecordBatch) -> Result<RecordBatch> {
83 match self {
84 FlatBatchConverter::Flat(f) => f.convert_batch(batch),
85 FlatBatchConverter::PrimaryKey { format, num_fields } => {
86 format.convert_flat_batch(batch, *num_fields)
87 }
88 }
89 }
90}
91
92enum SeriesBoundarySplit {
94 Continue(RecordBatch),
96 Split {
98 current_file_tail: Option<RecordBatch>,
100 next_file_head: RecordBatch,
102 },
103}
104
105fn encoded_primary_keys(batch: &RecordBatch) -> Result<(&UInt32Array, &BinaryArray)> {
107 let column = batch.column(primary_key_column_index(batch.num_columns()));
108 let primary_keys = column
109 .as_any()
110 .downcast_ref::<PrimaryKeyArray>()
111 .with_context(|| InvalidRecordBatchSnafu {
112 reason: format!(
113 "expected dictionary primary key column, got {:?}",
114 column.data_type()
115 ),
116 })?;
117 let values = primary_keys
118 .values()
119 .as_any()
120 .downcast_ref::<BinaryArray>()
121 .with_context(|| InvalidRecordBatchSnafu {
122 reason: format!(
123 "expected binary primary key values, got {:?}",
124 primary_keys.values().data_type()
125 ),
126 })?;
127 Ok((primary_keys.keys(), values))
128}
129
130fn split_at_next_series(
139 batch: RecordBatch,
140 current_primary_key: &[u8],
141) -> Result<SeriesBoundarySplit> {
142 let (keys, values) = encoded_primary_keys(&batch)?;
143 let Some(offset) = (0..batch.num_rows())
144 .find(|&row| values.value(keys.value(row) as usize) != current_primary_key)
145 else {
146 return Ok(SeriesBoundarySplit::Continue(batch));
147 };
148
149 let next_file_head = batch.slice(offset, batch.num_rows() - offset);
150 let current_file_tail = (offset > 0).then(|| batch.slice(0, offset));
151 Ok(SeriesBoundarySplit::Split {
152 current_file_tail,
153 next_file_head,
154 })
155}
156
157pub struct ParquetWriter<'a, F: WriterFactory, I: IndexerBuilder, P: FilePathProvider> {
159 path_provider: P,
161 writer: Option<AsyncArrowWriter<SizeAwareWriter<F::Writer>>>,
162 current_file: FileId,
164 writer_factory: F,
165 metadata: RegionMetadataRef,
167 index_config: IndexConfig,
169 indexer_builder: I,
171 current_indexer: Option<Indexer>,
173 bytes_written: Arc<AtomicUsize>,
174 file_cleaner: Option<TempFileCleaner>,
176 metrics: &'a mut Metrics,
178}
179
180pub trait WriterFactory {
181 type Writer: AsyncWrite + Send + Unpin;
182 fn create(&mut self, file_path: &str) -> impl Future<Output = Result<Self::Writer>>;
183}
184
185pub struct ObjectStoreWriterFactory {
186 object_store: ObjectStore,
187}
188
189impl WriterFactory for ObjectStoreWriterFactory {
190 type Writer = Compat<FuturesAsyncWriter>;
191
192 async fn create(&mut self, file_path: &str) -> Result<Self::Writer> {
193 self.object_store
194 .writer_with(file_path)
195 .chunk(DEFAULT_WRITE_BUFFER_SIZE.as_bytes() as usize)
196 .concurrent(DEFAULT_WRITE_CONCURRENCY)
197 .await
198 .map(|v| v.into_futures_async_write().compat_write())
199 .context(OpenDalSnafu)
200 }
201}
202
203impl<'a, I, P> ParquetWriter<'a, ObjectStoreWriterFactory, I, P>
204where
205 P: FilePathProvider,
206 I: IndexerBuilder,
207{
208 pub async fn new_with_object_store(
209 object_store: ObjectStore,
210 metadata: RegionMetadataRef,
211 index_config: IndexConfig,
212 indexer_builder: I,
213 path_provider: P,
214 metrics: &'a mut Metrics,
215 ) -> ParquetWriter<'a, ObjectStoreWriterFactory, I, P> {
216 ParquetWriter::new(
217 ObjectStoreWriterFactory { object_store },
218 metadata,
219 index_config,
220 indexer_builder,
221 path_provider,
222 metrics,
223 )
224 .await
225 }
226
227 pub(crate) fn with_file_cleaner(mut self, cleaner: TempFileCleaner) -> Self {
228 self.file_cleaner = Some(cleaner);
229 self
230 }
231}
232
233impl<'a, F, I, P> ParquetWriter<'a, F, I, P>
234where
235 F: WriterFactory,
236 I: IndexerBuilder,
237 P: FilePathProvider,
238{
239 pub async fn new(
241 factory: F,
242 metadata: RegionMetadataRef,
243 index_config: IndexConfig,
244 indexer_builder: I,
245 path_provider: P,
246 metrics: &'a mut Metrics,
247 ) -> ParquetWriter<'a, F, I, P> {
248 let init_file = FileId::random();
249 let indexer = indexer_builder
250 .build(RegionFileId::new(metadata.region_id, init_file), 0, None)
251 .await;
252
253 ParquetWriter {
254 path_provider,
255 writer: None,
256 current_file: init_file,
257 writer_factory: factory,
258 metadata,
259 index_config,
260 indexer_builder,
261 current_indexer: Some(indexer),
262 bytes_written: Arc::new(AtomicUsize::new(0)),
263 file_cleaner: None,
264 metrics,
265 }
266 }
267
268 async fn finish_current_file(
270 &mut self,
271 ssts: &mut SstInfoArray,
272 stats: &mut SourceStats,
273 ) -> Result<()> {
274 if let Some(mut current_writer) = mem::take(&mut self.writer) {
276 let mut stats = mem::take(stats);
277 assert!(stats.num_rows > 0);
279
280 debug!(
281 "Finishing current file {}, file size: {}, num rows: {}",
282 self.current_file,
283 self.bytes_written.load(Ordering::Relaxed),
284 stats.num_rows
285 );
286
287 let mut index_output = IndexOutput::default();
290 match self.index_config.build_mode {
291 IndexBuildMode::Sync => {
292 index_output = self.current_indexer.as_mut().unwrap().finish().await;
293 }
294 IndexBuildMode::Async => {
295 debug!(
296 "Index for file {} will be built asynchronously later",
297 self.current_file
298 );
299 }
300 }
301 current_writer.flush().await.context(WriteParquetSnafu)?;
302
303 let parquet_metadata = current_writer.close().await.context(WriteParquetSnafu)?;
304 let file_size = self.bytes_written.load(Ordering::Relaxed) as u64;
305
306 let time_range = stats.time_range.unwrap();
308
309 let max_row_group_uncompressed_size: u64 = parquet_metadata
310 .row_groups()
311 .iter()
312 .map(|rg| {
313 rg.columns()
314 .iter()
315 .map(|c| c.uncompressed_size() as u64)
316 .sum::<u64>()
317 })
318 .max()
319 .unwrap_or(0);
320 let num_series = stats.series_estimator.finish();
321 ssts.push(SstInfo {
322 file_id: self.current_file,
323 time_range,
324 file_size,
325 max_row_group_uncompressed_size,
326 num_rows: stats.num_rows,
327 num_row_groups: parquet_metadata.num_row_groups() as u64,
328 file_metadata: Some(Arc::new(parquet_metadata)),
329 index_metadata: index_output,
330 num_series,
331 });
332 self.current_file = FileId::random();
333 self.bytes_written.store(0, Ordering::Relaxed)
334 };
335
336 Ok(())
337 }
338
339 pub async fn write_all_flat(
350 &mut self,
351 source: FlatSource,
352 override_sequence: Option<SequenceNumber>,
353 opts: &WriteOptions,
354 ) -> Result<SstInfoArray> {
355 let mut options = FlatSchemaOptions::from_encoding(self.metadata.primary_key_encoding);
356
357 if source.schema().fields().iter().any(is_json2_extension_type) {
358 options.concretized_json_types = source
359 .schema()
360 .fields()
361 .iter()
362 .filter(|&field| is_json2_extension_type(field))
363 .map(|field| (field.name().clone(), field.data_type().clone()))
364 .collect::<HashMap<_, _>>();
365 }
366
367 let converter = FlatBatchConverter::Flat(
368 FlatWriteFormat::new(self.metadata.clone(), &options)
369 .with_override_sequence(override_sequence),
370 );
371 let res = self.write_all_flat_inner(source, &converter, opts).await;
372 if res.is_err() {
373 let file_id = self.current_file;
374 if let Some(cleaner) = &self.file_cleaner {
375 cleaner.clean_by_file_id(file_id).await;
376 }
377 }
378 res
379 }
380
381 pub async fn write_all_flat_as_primary_key(
388 &mut self,
389 source: FlatSource,
390 override_sequence: Option<SequenceNumber>,
391 opts: &WriteOptions,
392 ) -> Result<SstInfoArray> {
393 let num_fields = self.metadata.field_columns().count();
394 let converter = FlatBatchConverter::PrimaryKey {
395 format: PrimaryKeyWriteFormat::new(self.metadata.clone())
396 .with_override_sequence(override_sequence),
397 num_fields,
398 };
399 let res = self.write_all_flat_inner(source, &converter, opts).await;
400 if res.is_err() {
401 let file_id = self.current_file;
402 if let Some(cleaner) = &self.file_cleaner {
403 cleaner.clean_by_file_id(file_id).await;
404 }
405 }
406 res
407 }
408
409 async fn write_all_flat_inner(
410 &mut self,
411 mut source: FlatSource,
412 converter: &FlatBatchConverter,
413 opts: &WriteOptions,
414 ) -> Result<SstInfoArray> {
415 let mut results = smallvec![];
416 let mut stats = SourceStats::default();
417
418 loop {
419 let start = Instant::now();
420 let batch = match source.next_batch().await {
421 Ok(Some(batch)) => batch,
422 Ok(None) => break,
423 Err(e) => {
424 self.abort_current_indexer().await;
425 return Err(e);
426 }
427 };
428 self.metrics.iter_source += start.elapsed();
429
430 if self.metadata.primary_key.is_empty() {
431 self.append_flat_batch(&batch, converter, opts, &mut stats)
432 .await?;
433 if self.exceeds_max_file_size(opts) {
434 self.finish_current_file(&mut results, &mut stats).await?;
435 }
436 } else if self.exceeds_max_file_size(opts)
437 && let Some(current_primary_key) = stats.last_primary_key.as_deref()
438 {
439 let series_split = match split_at_next_series(batch, current_primary_key) {
440 Ok(series_split) => series_split,
441 Err(e) => {
442 self.abort_current_indexer().await;
443 return Err(e);
444 }
445 };
446 match series_split {
447 SeriesBoundarySplit::Continue(batch) => {
448 self.append_flat_batch(&batch, converter, opts, &mut stats)
449 .await?;
450 }
451 SeriesBoundarySplit::Split {
452 current_file_tail,
453 next_file_head,
454 } => {
455 if let Some(tail) = current_file_tail {
456 self.append_flat_batch(&tail, converter, opts, &mut stats)
457 .await?;
458 }
459 self.finish_current_file(&mut results, &mut stats).await?;
460 self.append_flat_batch(&next_file_head, converter, opts, &mut stats)
461 .await?;
462 }
463 }
464 } else {
465 self.append_flat_batch(&batch, converter, opts, &mut stats)
466 .await?;
467 }
468 }
469
470 self.finish_current_file(&mut results, &mut stats).await?;
471
472 Ok(results)
474 }
475
476 fn customize_column_config(
478 builder: WriterPropertiesBuilder,
479 region_metadata: &RegionMetadataRef,
480 ) -> WriterPropertiesBuilder {
481 let ts_col = ColumnPath::new(vec![
482 region_metadata
483 .time_index_column()
484 .column_schema
485 .name
486 .clone(),
487 ]);
488 let seq_col = ColumnPath::new(vec![SEQUENCE_COLUMN_NAME.to_string()]);
489 let op_type_col = ColumnPath::new(vec![OP_TYPE_COLUMN_NAME.to_string()]);
490
491 builder
492 .set_column_encoding(seq_col.clone(), Encoding::DELTA_BINARY_PACKED)
493 .set_column_dictionary_enabled(seq_col, false)
494 .set_column_encoding(ts_col.clone(), Encoding::DELTA_BINARY_PACKED)
495 .set_column_dictionary_enabled(ts_col, false)
496 .set_column_compression(op_type_col, Compression::UNCOMPRESSED)
497 }
498
499 async fn append_flat_batch(
500 &mut self,
501 batch: &RecordBatch,
502 converter: &FlatBatchConverter,
503 opts: &WriteOptions,
504 stats: &mut SourceStats,
505 ) -> Result<()> {
506 let result = async {
507 let arrow_batch = converter.convert_batch(batch)?;
508 let start = Instant::now();
509 self.maybe_init_writer(arrow_batch.schema_ref(), opts)
510 .await?
511 .write(&arrow_batch)
512 .await
513 .context(WriteParquetSnafu)?;
514 self.metrics.write_batch += start.elapsed();
515
516 stats.update_flat(batch)?;
517 if matches!(self.index_config.build_mode, IndexBuildMode::Sync) {
518 let start = Instant::now();
519 self.current_indexer
521 .as_mut()
522 .unwrap()
523 .update_flat(batch)
524 .await;
525 self.metrics.update_index += start.elapsed();
526 }
527 Ok(())
528 }
529 .await;
530
531 if result.is_err() {
532 self.abort_current_indexer().await;
533 }
534 result
535 }
536
537 fn exceeds_max_file_size(&self, opts: &WriteOptions) -> bool {
538 opts.max_file_size
539 .is_some_and(|max_size| self.bytes_written.load(Ordering::Relaxed) >= max_size)
540 }
541
542 async fn abort_current_indexer(&mut self) {
543 if let Some(indexer) = &mut self.current_indexer {
544 indexer.abort().await;
545 }
546 }
547
548 async fn maybe_init_writer(
549 &mut self,
550 schema: &SchemaRef,
551 opts: &WriteOptions,
552 ) -> Result<&mut AsyncArrowWriter<SizeAwareWriter<F::Writer>>> {
553 if let Some(ref mut w) = self.writer {
554 Ok(w)
555 } else {
556 let json = self.metadata.to_json().context(InvalidMetadataSnafu)?;
557 let key_value_meta = KeyValue::new(PARQUET_METADATA_KEY.to_string(), json);
558
559 let props_builder = WriterProperties::builder()
561 .set_key_value_metadata(Some(vec![key_value_meta]))
562 .set_compression(Compression::ZSTD(ZstdLevel::default()))
563 .set_encoding(Encoding::PLAIN)
564 .set_max_row_group_row_count(Some(opts.row_group_size))
565 .set_column_index_truncate_length(None)
566 .set_statistics_truncate_length(None);
567
568 let props_builder = Self::customize_column_config(props_builder, &self.metadata);
569 let writer_props = props_builder.build();
570
571 let sst_file_path = self.path_provider.build_sst_file_path(RegionFileId::new(
572 self.metadata.region_id,
573 self.current_file,
574 ));
575 let writer = SizeAwareWriter::new(
576 self.writer_factory.create(&sst_file_path).await?,
577 self.bytes_written.clone(),
578 );
579 let arrow_writer =
580 AsyncArrowWriter::try_new(writer, maybe_wrap_schema(schema)?, Some(writer_props))
581 .context(WriteParquetSnafu)?;
582 self.writer = Some(arrow_writer);
583
584 let indexer = self
585 .indexer_builder
586 .build(
587 RegionFileId::new(self.metadata.region_id, self.current_file),
588 0,
589 Some(opts.row_group_size),
590 )
591 .await;
592 self.current_indexer = Some(indexer);
593
594 Ok(self.writer.as_mut().unwrap())
596 }
597 }
598}
599
600#[derive(Default)]
601struct SourceStats {
602 num_rows: usize,
604 time_range: Option<(Timestamp, Timestamp)>,
606 last_primary_key: Option<Bytes>,
608 series_estimator: SeriesEstimator,
610}
611
612impl SourceStats {
613 fn update_flat(&mut self, record_batch: &RecordBatch) -> Result<()> {
614 if record_batch.num_rows() == 0 {
615 return Ok(());
616 }
617
618 self.num_rows += record_batch.num_rows();
619 self.series_estimator.update_flat(record_batch);
620 let (keys, values) = encoded_primary_keys(record_batch)?;
621 let key = keys.value(record_batch.num_rows() - 1);
622 self.last_primary_key = Some(Bytes::copy_from_slice(values.value(key as usize)));
623
624 let time_index_col_idx = time_index_column_index(record_batch.num_columns());
626 let timestamp_array = record_batch.column(time_index_col_idx);
627
628 if let Some((min_in_batch, max_in_batch)) = timestamp_range_from_array(timestamp_array)? {
629 if let Some(time_range) = &mut self.time_range {
630 time_range.0 = time_range.0.min(min_in_batch);
631 time_range.1 = time_range.1.max(max_in_batch);
632 } else {
633 self.time_range = Some((min_in_batch, max_in_batch));
634 }
635 }
636
637 Ok(())
638 }
639}
640
641fn timestamp_range_from_array(
643 timestamp_array: &ArrayRef,
644) -> Result<Option<(Timestamp, Timestamp)>> {
645 let (min_ts, max_ts) = match timestamp_array.data_type() {
646 DataType::Timestamp(TimeUnit::Second, _) => {
647 let array = timestamp_array
648 .as_any()
649 .downcast_ref::<TimestampSecondArray>()
650 .unwrap();
651 let min_val = min(array).map(Timestamp::new_second);
652 let max_val = max(array).map(Timestamp::new_second);
653 (min_val, max_val)
654 }
655 DataType::Timestamp(TimeUnit::Millisecond, _) => {
656 let array = timestamp_array
657 .as_any()
658 .downcast_ref::<TimestampMillisecondArray>()
659 .unwrap();
660 let min_val = min(array).map(Timestamp::new_millisecond);
661 let max_val = max(array).map(Timestamp::new_millisecond);
662 (min_val, max_val)
663 }
664 DataType::Timestamp(TimeUnit::Microsecond, _) => {
665 let array = timestamp_array
666 .as_any()
667 .downcast_ref::<TimestampMicrosecondArray>()
668 .unwrap();
669 let min_val = min(array).map(Timestamp::new_microsecond);
670 let max_val = max(array).map(Timestamp::new_microsecond);
671 (min_val, max_val)
672 }
673 DataType::Timestamp(TimeUnit::Nanosecond, _) => {
674 let array = timestamp_array
675 .as_any()
676 .downcast_ref::<TimestampNanosecondArray>()
677 .unwrap();
678 let min_val = min(array).map(Timestamp::new_nanosecond);
679 let max_val = max(array).map(Timestamp::new_nanosecond);
680 (min_val, max_val)
681 }
682 _ => {
683 return UnexpectedSnafu {
684 reason: format!(
685 "Unexpected data type of time index: {:?}",
686 timestamp_array.data_type()
687 ),
688 }
689 .fail();
690 }
691 };
692
693 Ok(min_ts.zip(max_ts))
695}
696
697struct SizeAwareWriter<W> {
700 inner: W,
701 size: Arc<AtomicUsize>,
702}
703
704impl<W> SizeAwareWriter<W> {
705 fn new(inner: W, size: Arc<AtomicUsize>) -> Self {
706 Self {
707 inner,
708 size: size.clone(),
709 }
710 }
711}
712
713impl<W> AsyncWrite for SizeAwareWriter<W>
714where
715 W: AsyncWrite + Unpin,
716{
717 fn poll_write(
718 mut self: Pin<&mut Self>,
719 cx: &mut Context<'_>,
720 buf: &[u8],
721 ) -> Poll<std::result::Result<usize, std::io::Error>> {
722 let this = self.as_mut().get_mut();
723
724 match Pin::new(&mut this.inner).poll_write(cx, buf) {
725 Poll::Ready(Ok(bytes_written)) => {
726 this.size.fetch_add(bytes_written, Ordering::Relaxed);
727 Poll::Ready(Ok(bytes_written))
728 }
729 other => other,
730 }
731 }
732
733 fn poll_flush(
734 mut self: Pin<&mut Self>,
735 cx: &mut Context<'_>,
736 ) -> Poll<std::result::Result<(), std::io::Error>> {
737 Pin::new(&mut self.inner).poll_flush(cx)
738 }
739
740 fn poll_shutdown(
741 mut self: Pin<&mut Self>,
742 cx: &mut Context<'_>,
743 ) -> Poll<std::result::Result<(), std::io::Error>> {
744 Pin::new(&mut self.inner).poll_shutdown(cx)
745 }
746}