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
| 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)?; |
no test coverage detected