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, 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
70/// Converts a flat RecordBatch for writing to parquet.
71enum FlatBatchConverter {
72    /// Write as-is in flat format.
73    Flat(FlatWriteFormat),
74    /// Convert flat batch to primary-key format by stripping tag columns.
75    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
92/// Result of splitting a batch at the next series boundary.
93enum SeriesBoundarySplit {
94    /// The whole batch belongs to the current series and stays in the current file.
95    Continue(RecordBatch),
96    /// The batch crosses a series boundary.
97    Split {
98        /// Remaining rows of the current series, appended to the current file.
99        current_file_tail: Option<RecordBatch>,
100        /// Rows of the following series that start the next file.
101        next_file_head: RecordBatch,
102    },
103}
104
105/// Returns the dictionary keys and values of the encoded `__primary_key` column.
106fn 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
130/// Splits `batch` at the first row whose encoded primary key differs from
131/// `current_primary_key`, the last primary key written to the current file.
132///
133/// The writer finishes an oversized file at such a series boundary so that the
134/// primary key ranges of output files never overlap, which allows pickers like
135/// TWCS to detect overlapping files by their time and primary key ranges. A
136/// series is never split in the middle: if no row starts a new series, the whole
137/// batch stays in the current file even if it already exceeds the size limit.
138fn 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
157/// Parquet SST writer.
158pub struct ParquetWriter<'a, F: WriterFactory, I: IndexerBuilder, P: FilePathProvider> {
159    /// Path provider that creates SST and index file paths according to file id.
160    path_provider: P,
161    writer: Option<AsyncArrowWriter<SizeAwareWriter<F::Writer>>>,
162    /// Current active file id.
163    current_file: FileId,
164    writer_factory: F,
165    /// Region metadata of the source and the target SST.
166    metadata: RegionMetadataRef,
167    /// Global index config.
168    index_config: IndexConfig,
169    /// Indexer build that can create indexer for multiple files.
170    indexer_builder: I,
171    /// Current active indexer.
172    current_indexer: Option<Indexer>,
173    bytes_written: Arc<AtomicUsize>,
174    /// Cleaner to remove temp files on failure.
175    file_cleaner: Option<TempFileCleaner>,
176    /// Write metrics
177    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    /// Creates a new parquet SST writer.
240    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    /// Finishes current SST file and index file.
269    async fn finish_current_file(
270        &mut self,
271        ssts: &mut SstInfoArray,
272        stats: &mut SourceStats,
273    ) -> Result<()> {
274        // maybe_init_writer will re-create a new file.
275        if let Some(mut current_writer) = mem::take(&mut self.writer) {
276            let mut stats = mem::take(stats);
277            // At least one row has been written.
278            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            // Finish indexer and writer.
288            // safety: writer and index can only be both present or not.
289            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            // Safety: num rows > 0 so we must have min/max.
307            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    /// Iterates FlatSource and writes all RecordBatch in flat format to Parquet file.
340    ///
341    /// The source must yield batches globally sorted by the encoded primary key.
342    /// `opts.max_file_size` is a soft limit for regions with a primary key: the
343    /// current file is finished at the next series boundary after exceeding the
344    /// limit (see [split_at_next_series]), so a series larger than the limit
345    /// stays in one file. Regions without a primary key are split at batch
346    /// boundaries.
347    ///
348    /// Returns the [SstInfo] if the SST is written.
349    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    /// Iterates FlatSource and writes all RecordBatch in primary-key format to Parquet file.
382    ///
383    /// The source must yield batches globally sorted by the encoded primary key.
384    /// See [Self::write_all_flat] for the `opts.max_file_size` semantics.
385    ///
386    /// Returns the [SstInfo] if the SST is written.
387    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        // object_store.write will make sure all bytes are written or an error is raised.
473        Ok(results)
474    }
475
476    /// Customizes per-column config according to schema and maybe column cardinality.
477    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                // safety: self.current_indexer must be set when first batch has been written.
520                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            // TODO(yingwen): Find and set proper column encoding for internal columns: op type and tsid.
560            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            // safety: self.writer is assigned above
595            Ok(self.writer.as_mut().unwrap())
596        }
597    }
598}
599
600#[derive(Default)]
601struct SourceStats {
602    /// Number of rows fetched.
603    num_rows: usize,
604    /// Time range of fetched batches.
605    time_range: Option<(Timestamp, Timestamp)>,
606    /// Last primary key written to the current file.
607    last_primary_key: Option<Bytes>,
608    /// Series estimator for computing num_series.
609    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        // Get the timestamp column by index
625        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
641/// Gets min and max timestamp from an timestamp array.
642fn 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    // If min timestamp exists, max timestamp should also exist.
694    Ok(min_ts.zip(max_ts))
695}
696
697/// Workaround for [AsyncArrowWriter] does not provide a method to
698/// get total bytes written after close.
699struct 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}