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

Method merge_files_by_columns

crates/paimon/src/table/data_evolution_reader.rs:242–440  ·  view source on GitHub ↗

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,
    )

Source from the content-addressed store, hash-verified

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 {

Callers

nothing calls this directly

Calls 2

is_emptyMethod · 0.45

Tested by

no test coverage detected