Merge multiple logical sources column-wise for data evolution. Normal partial-column files remain one source per file. Rolling `.blob` files are first grouped into a logical BlobBunch source per field, then source streams are merged by projected field position.
(
&self,
split: &DataSplit,
prepared_group: &PreparedMergeGroup,
row_ranges: Option<Vec<RowRange>>,
expected_output_rows: usize,
)
| 240 | /// files are first grouped into a logical BlobBunch source per field, then |
| 241 | /// source streams are merged by projected field position. |
| 242 | fn merge_files_by_columns( |
| 243 | &self, |
| 244 | split: &DataSplit, |
| 245 | prepared_group: &PreparedMergeGroup, |
| 246 | row_ranges: Option<Vec<RowRange>>, |
| 247 | expected_output_rows: usize, |
| 248 | ) -> crate::Result<ArrowRecordBatchStream> { |
| 249 | if prepared_group.files.is_empty() { |
| 250 | return Ok(futures::stream::empty().boxed()); |
| 251 | } |
| 252 | |
| 253 | let file_io = self.file_io.clone(); |
| 254 | let schema_manager = self.schema_manager.clone(); |
| 255 | let table_schema_id = self.table_schema_id; |
| 256 | let split = split.clone(); |
| 257 | let prepared_group = prepared_group.clone(); |
| 258 | let read_type = self.file_read_type.clone(); |
| 259 | let table_fields = self.table_fields.clone(); |
| 260 | let blob_descriptor_fields = self.blob_descriptor_fields.clone(); |
| 261 | let blob_as_descriptor = self.blob_as_descriptor; |
| 262 | // Batch size for column-merge output. Matches the default Parquet reader batch size. |
| 263 | const MERGE_BATCH_SIZE: usize = 1024; |
| 264 | let target_schema = build_target_arrow_schema(&read_type)?; |
| 265 | |
| 266 | Ok(try_stream! { |
| 267 | let file_infos = load_file_infos( |
| 268 | &schema_manager, |
| 269 | table_schema_id, |
| 270 | &table_fields, |
| 271 | &prepared_group.files, |
| 272 | ) |
| 273 | .await?; |
| 274 | let source_plan = build_source_plan(&prepared_group, &file_infos, &read_type, &blob_descriptor_fields)?; |
| 275 | |
| 276 | let active_source_indices: Vec<usize> = source_plan |
| 277 | .sources |
| 278 | .iter() |
| 279 | .enumerate() |
| 280 | .filter_map(|(idx, source)| (!source.read_fields().is_empty()).then_some(idx)) |
| 281 | .collect(); |
| 282 | |
| 283 | // Edge case: no file provides any projected column. |
| 284 | if active_source_indices.is_empty() { |
| 285 | let mut emitted = 0usize; |
| 286 | while emitted < expected_output_rows { |
| 287 | let rows_to_emit = (expected_output_rows - emitted).min(MERGE_BATCH_SIZE); |
| 288 | let columns: Vec<Arc<dyn arrow_array::Array>> = target_schema |
| 289 | .fields() |
| 290 | .iter() |
| 291 | .map(|f| arrow_array::new_null_array(f.data_type(), rows_to_emit)) |
| 292 | .collect(); |
| 293 | let batch = if columns.is_empty() { |
| 294 | RecordBatch::try_new_with_options( |
| 295 | target_schema.clone(), |
| 296 | columns, |
| 297 | &arrow_array::RecordBatchOptions::new().with_row_count(Some(rows_to_emit)), |
| 298 | ) |
| 299 | } else { |
nothing calls this directly
no test coverage detected