mito2/memtable/bulk/
row_group_reader.rs1use std::sync::Arc;
16
17use bytes::Bytes;
18use datatypes::extension::json::is_json2_extension_type;
19use parquet::arrow::ProjectionMask;
20use parquet::arrow::arrow_reader::{
21 ArrowReaderMetadata, ArrowReaderOptions, ParquetRecordBatchReader,
22 ParquetRecordBatchReaderBuilder, RowSelection,
23};
24use parquet::file::metadata::ParquetMetaData;
25use snafu::ResultExt;
26
27use crate::error;
28use crate::error::ReadDataPartSnafu;
29use crate::memtable::bulk::chunk_reader::MemtableChunkReader;
30use crate::memtable::bulk::context::BulkIterContextRef;
31
32pub(crate) struct MemtableRowGroupReaderBuilder {
33 projection: ProjectionMask,
34 arrow_metadata: ArrowReaderMetadata,
35 data: Bytes,
36 batch_size: usize,
37}
38
39impl MemtableRowGroupReaderBuilder {
40 pub(crate) fn try_new(
41 context: &BulkIterContextRef,
42 projection: ProjectionMask,
43 parquet_metadata: Arc<ParquetMetaData>,
44 data: Bytes,
45 ) -> error::Result<Self> {
46 let mut arrow_reader_options = ArrowReaderOptions::new();
48
49 if !context
56 .read_format()
57 .arrow_schema()
58 .fields()
59 .iter()
60 .any(is_json2_extension_type)
61 {
62 arrow_reader_options =
63 arrow_reader_options.with_schema(context.read_format().arrow_schema().clone());
64 }
65
66 let arrow_metadata =
67 ArrowReaderMetadata::try_new(parquet_metadata.clone(), arrow_reader_options)
68 .context(ReadDataPartSnafu)?;
69 Ok(Self {
70 projection,
71 arrow_metadata,
72 data,
73 batch_size: context.batch_size(),
74 })
75 }
76
77 pub(crate) fn build_row_group_reader(
79 &self,
80 row_group_idx: usize,
81 row_selection: Option<RowSelection>,
82 ) -> error::Result<ParquetRecordBatchReader> {
83 let chunk_reader = MemtableChunkReader::new(self.data.clone());
84
85 let mut builder = ParquetRecordBatchReaderBuilder::new_with_metadata(
86 chunk_reader,
87 self.arrow_metadata.clone(),
88 )
89 .with_row_groups(vec![row_group_idx])
90 .with_projection(self.projection.clone())
91 .with_batch_size(self.batch_size);
92
93 if let Some(selection) = row_selection {
94 builder = builder.with_row_selection(selection);
95 }
96
97 builder.build().context(ReadDataPartSnafu)
98 }
99}