MCPcopy Create free account
hub / github.com/apache/paimon-rust / read_batch_stream

Method read_batch_stream

crates/paimon/src/arrow/format/parquet.rs:138–200  ·  view source on GitHub ↗
(
        &self,
        reader: Box<dyn FileRead>,
        file_size: u64,
        read_fields: &[DataField],
        predicates: Option<&FilePredicates>,
        batch_size: Option<usize>,
        r

Source from the content-addressed store, hash-verified

136#[async_trait]
137impl FormatFileReader for ParquetFormatReader {
138 async fn read_batch_stream(
139 &self,
140 reader: Box<dyn FileRead>,
141 file_size: u64,
142 read_fields: &[DataField],
143 predicates: Option<&FilePredicates>,
144 batch_size: Option<usize>,
145 row_selection: Option<Vec<RowRange>>,
146 ) -> crate::Result<ArrowRecordBatchStream> {
147 let arrow_file_reader = ArrowFileReader::new(file_size, reader);
148
149 let mut batch_stream_builder =
150 ParquetRecordBatchStreamBuilder::new(arrow_file_reader).await?;
151
152 let parquet_schema = batch_stream_builder.parquet_schema().clone();
153 let root_schema = parquet_schema.root_schema();
154 let root_indices: Vec<usize> = read_fields
155 .iter()
156 .filter_map(|f| {
157 root_schema
158 .get_fields()
159 .iter()
160 .position(|pf| pf.name() == f.name())
161 })
162 .collect();
163
164 let mask = ProjectionMask::roots(&parquet_schema, root_indices);
165 batch_stream_builder = batch_stream_builder.with_projection(mask);
166
167 let empty_predicates = Vec::new();
168 let (preds, file_fields): (&[Predicate], &[DataField]) = match predicates {
169 Some(fp) => (&fp.predicates, &fp.file_fields),
170 None => (&empty_predicates, &[]),
171 };
172
173 let parquet_row_filter = build_parquet_row_filter(&parquet_schema, preds, file_fields)?;
174 if let Some(f) = parquet_row_filter {
175 batch_stream_builder = batch_stream_builder.with_row_filter(f);
176 }
177
178 let predicate_row_selection = build_predicate_row_selection(
179 batch_stream_builder.metadata().row_groups(),
180 preds,
181 file_fields,
182 )?;
183 let mut combined_selection = predicate_row_selection;
184
185 if let Some(ref ranges) = row_selection {
186 let range_selection =
187 build_row_ranges_selection(batch_stream_builder.metadata().row_groups(), ranges);
188 combined_selection =
189 intersect_optional_row_selections(combined_selection, Some(range_selection));
190 }
191 if let Some(sel) = combined_selection {
192 batch_stream_builder = batch_stream_builder.with_row_selection(sel);
193 }
194 if let Some(size) = batch_size {
195 batch_stream_builder = batch_stream_builder.with_batch_size(size);

Callers

nothing calls this directly

Calls 11

build_parquet_row_filterFunction · 0.85
with_projectionMethod · 0.80
metadataMethod · 0.80
with_batch_sizeMethod · 0.80
iterMethod · 0.45
positionMethod · 0.45
nameMethod · 0.45
buildMethod · 0.45

Tested by

no test coverage detected