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