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