Skip to main content

mito2/sst/parquet/
writer.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Parquet writer.
16
17use 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
71/// Converts a flat RecordBatch for writing to parquet.
72enum FlatBatchConverter {
73    /// Write as-is in flat format.
74    Flat(FlatWriteFormat),
75    /// Convert flat batch to primary-key format by stripping tag columns.
76    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
93/// Result of splitting a batch at the next series boundary.
94enum SeriesBoundarySplit {
95    /// The whole batch belongs to the current series and stays in the current file.
96    Continue(RecordBatch),
97    /// The batch crosses a series boundary.
98    Split {
99        /// Remaining rows of the current series, appended to the current file.
100        current_file_tail: Option<RecordBatch>,
101        /// Rows of the following series that start the next file.
102        next_file_head: RecordBatch,
103    },
104}
105
106/// Returns the dictionary keys and values of the encoded `__primary_key` column.
107fn 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
131/// Splits `batch` at the first row whose encoded primary key differs from
132/// `current_primary_key`, the last primary key written to the current file.
133///
134/// The writer finishes an oversized file at such a series boundary so that the
135/// primary key ranges of output files never overlap, which allows pickers like
136/// TWCS to detect overlapping files by their time and primary key ranges. A
137/// series is never split in the middle: if no row starts a new series, the whole
138/// batch stays in the current file even if it already exceeds the size limit.
139fn 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
158/// Parquet SST writer.
159pub struct ParquetWriter<'a, F: WriterFactory, I: IndexerBuilder, P: FilePathProvider> {
160    /// Path provider that creates SST and index file paths according to file id.
161    path_provider: P,
162    writer: Option<AsyncArrowWriter<SizeAwareWriter<F::Writer>>>,
163    /// Current active file id.
164    current_file: FileId,
165    writer_factory: F,
166    /// Region metadata of the source and the target SST.
167    metadata: RegionMetadataRef,
168    /// Global index config.
169    index_config: IndexConfig,
170    /// Indexer build that can create indexer for multiple files.
171    indexer_builder: I,
172    /// Current active indexer.
173    current_indexer: Option<Indexer>,
174    bytes_written: Arc<AtomicUsize>,
175    /// Cleaner to remove temp files on failure.
176    file_cleaner: Option<TempFileCleaner>,
177    /// Write metrics
178    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    /// Creates a new parquet SST writer.
245    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    /// Finishes current SST file and index file.
274    async fn finish_current_file(
275        &mut self,
276        ssts: &mut SstInfoArray,
277        stats: &mut SourceStats,
278    ) -> Result<()> {
279        // maybe_init_writer will re-create a new file.
280        if let Some(mut current_writer) = mem::take(&mut self.writer) {
281            let mut stats = mem::take(stats);
282            // At least one row has been written.
283            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            // Finish indexer and writer.
293            // safety: writer and index can only be both present or not.
294            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            // Safety: num rows > 0 so we must have min/max.
312            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    /// Iterates FlatSource and writes all RecordBatch in flat format to Parquet file.
345    ///
346    /// The source must yield batches globally sorted by the encoded primary key.
347    /// `opts.max_file_size` is a soft limit for regions with a primary key: the
348    /// current file is finished at the next series boundary after exceeding the
349    /// limit (see [split_at_next_series]), so a series larger than the limit
350    /// stays in one file. Regions without a primary key are split at batch
351    /// boundaries.
352    ///
353    /// Returns the [SstInfo] if the SST is written.
354    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    /// Iterates FlatSource and writes all RecordBatch in primary-key format to Parquet file.
387    ///
388    /// The source must yield batches globally sorted by the encoded primary key.
389    /// See [Self::write_all_flat] for the `opts.max_file_size` semantics.
390    ///
391    /// Returns the [SstInfo] if the SST is written.
392    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        // object_store.write will make sure all bytes are written or an error is raised.
478        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                // safety: self.current_indexer must be set when first batch has been written.
502                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            // TODO(yingwen): Find and set proper column encoding for internal columns: op type and tsid.
542            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            // safety: self.writer is assigned above
593            Ok(self.writer.as_mut().unwrap())
594        }
595    }
596}
597
598#[derive(Default)]
599struct SourceStats {
600    /// Number of rows fetched.
601    num_rows: usize,
602    /// Time range of fetched batches.
603    time_range: Option<(Timestamp, Timestamp)>,
604    /// Last primary key written to the current file.
605    last_primary_key: Option<Bytes>,
606    /// Series estimator for computing num_series.
607    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        // Get the timestamp column by index
623        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
639/// Gets min and max timestamp from an timestamp array.
640fn 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    // If min timestamp exists, max timestamp should also exist.
692    Ok(min_ts.zip(max_ts))
693}
694
695/// Workaround for [AsyncArrowWriter] does not provide a method to
696/// get total bytes written after close.
697struct 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}