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

Method read_single_file_stream

crates/paimon/src/table/data_file_reader.rs:142–296  ·  view source on GitHub ↗

Read a single parquet file from a split, returning a lazy stream of batches. Optionally applies a deletion vector. Handles schema evolution using field-ID-based index mapping: - `data_fields`: if `Some`, the fields from the data file's schema (loaded via SchemaManager). Used to compute index mapping between `read_type` and data fields by field ID. - Columns missing from the file are filled with n

(
        &self,
        split: &DataSplit,
        file_meta: DataFileMeta,
        data_fields: Option<Vec<DataField>>,
        dv: Option<Arc<DeletionVector>>,
        row_ranges: Option<Vec<RowRan

Source from the content-addressed store, hash-verified

140 ///
141 /// Reference: [RawFileSplitRead.createFileReader](https://github.com/apache/paimon/blob/release-1.3/paimon-core/src/main/java/org/apache/paimon/operation/RawFileSplitRead.java)
142 pub(super) fn read_single_file_stream(
143 &self,
144 split: &DataSplit,
145 file_meta: DataFileMeta,
146 data_fields: Option<Vec<DataField>>,
147 dv: Option<Arc<DeletionVector>>,
148 row_ranges: Option<Vec<RowRange>>,
149 ) -> crate::Result<ArrowRecordBatchStream> {
150 let read_type = self.read_type.clone();
151 let table_fields = self.table_fields.clone();
152 let predicates = self.predicates.clone();
153 let file_io = self.file_io.clone();
154 let split = split.clone();
155 let blob_as_descriptor = self.blob_as_descriptor;
156
157 let target_schema = build_target_arrow_schema(&read_type)?;
158 let file_fields = data_fields.clone().unwrap_or_else(|| table_fields.clone());
159
160 // Compute index mapping and determine which columns to read from the file.
161 let (projected_read_fields, index_mapping) = if let Some(ref df) = data_fields {
162 let mapping = create_index_mapping(&read_type, df);
163 match mapping {
164 Some(ref idx_map) => {
165 let mut seen = std::collections::HashSet::new();
166 let fields_to_read: Vec<DataField> = idx_map
167 .iter()
168 .filter(|&&idx| idx != NULL_FIELD_INDEX && seen.insert(idx))
169 .map(|&idx| df[idx as usize].clone())
170 .collect();
171 (fields_to_read, Some(idx_map.clone()))
172 }
173 None => (df.clone(), None),
174 }
175 } else {
176 (read_type.clone(), None)
177 };
178
179 // Remap predicates from table-level to file-level indices.
180 let file_predicates = {
181 let remapped = crate::arrow::filtering::remap_predicates_to_file(
182 &predicates,
183 &table_fields,
184 &file_fields,
185 );
186 if remapped.is_empty() {
187 None
188 } else {
189 Some(crate::arrow::format::FilePredicates {
190 predicates: remapped,
191 file_fields: file_fields.clone(),
192 })
193 }
194 };
195
196 Ok(try_stream! {
197 let path_to_read = split.data_file_path(&file_meta);
198 let format_reader = create_format_reader(&path_to_read, blob_as_descriptor)?;
199 let input_file = file_io.new_input(&path_to_read)?;

Callers 1

open_source_streamFunction · 0.80

Calls 6

create_index_mappingFunction · 0.85
remap_predicates_to_fileFunction · 0.85
insertMethod · 0.80
iterMethod · 0.45
is_emptyMethod · 0.45

Tested by

no test coverage detected