From 17840ee233563088381f29af949890a341cd5aa6 Mon Sep 17 00:00:00 2001 From: xiaohongbo Date: Mon, 28 Sep 2026 07:42:38 -0700 Subject: [PATCH] [core] Honor scan.manifest.parallelism during planning --- crates/paimon/src/spec/core_options.rs | 55 ++++++++++++++++++++++++++ crates/paimon/src/table/table_scan.rs | 39 +++++++++++++++--- 2 files changed, 89 insertions(+), 5 deletions(-) diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index d474479b1..9168d4414 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -92,6 +92,7 @@ const MANIFEST_TARGET_FILE_SIZE_OPTION: &str = "manifest.target-file-size"; const MANIFEST_TARGET_SIZE_OPTION: &str = "manifest.target-size"; const MANIFEST_SIDECAR_ENABLED_OPTION: &str = "manifest.sidecar.enabled"; const MANIFEST_SORT_ENABLED_OPTION: &str = "manifest-sort.enabled"; +const SCAN_MANIFEST_PARALLELISM_OPTION: &str = "scan.manifest.parallelism"; const WRITE_PARQUET_BUFFER_SIZE_OPTION: &str = "write.parquet-buffer-size"; const READ_BATCH_SIZE_OPTION: &str = "read.batch-size"; const PARQUET_FILTER_COLUMN_INDEX_ENABLED_OPTION: &str = "parquet.filter.columnindex.enabled"; @@ -414,6 +415,35 @@ impl<'a> CoreOptions<'a> { Ok(value) } + /// Maximum concurrent manifest reads during scan planning. + /// + /// Matches Java Paimon's `scan.manifest.parallelism`: when unset, use the + /// number of processors available to this process. + pub fn scan_manifest_parallelism(&self) -> crate::Result { + let Some(raw) = self.options.get(SCAN_MANIFEST_PARALLELISM_OPTION) else { + return Ok(std::thread::available_parallelism() + .map(|value| value.get()) + .unwrap_or(1)); + }; + let value = raw + .parse::() + .map_err(|error| crate::Error::DataInvalid { + message: format!( + "Option '{SCAN_MANIFEST_PARALLELISM_OPTION}' must be a positive integer, got: {raw}" + ), + source: Some(Box::new(error)), + })?; + if value == 0 { + return Err(crate::Error::DataInvalid { + message: format!( + "Option '{SCAN_MANIFEST_PARALLELISM_OPTION}' must be greater than 0" + ), + source: None, + }); + } + Ok(value) + } + /// Mosaic row groups kept opening ahead of the consumer per file; `0` disables prefetch. pub fn mosaic_read_prefetch_row_groups(&self) -> crate::Result { let Some(raw) = self.options.get(MOSAIC_READ_PREFETCH_ROW_GROUPS_OPTION) else { @@ -2711,6 +2741,31 @@ mod tests { assert_eq!(parallelism(Some("many")), 64); } + #[test] + fn test_scan_manifest_parallelism() { + let parallelism = |value: Option<&str>| { + let options = value + .map(|value| { + HashMap::from([( + SCAN_MANIFEST_PARALLELISM_OPTION.to_string(), + value.to_string(), + )]) + }) + .unwrap_or_default(); + CoreOptions::new(&options).scan_manifest_parallelism() + }; + + assert_eq!( + parallelism(None).unwrap(), + std::thread::available_parallelism() + .map(|value| value.get()) + .unwrap_or(1) + ); + assert_eq!(parallelism(Some("8")).unwrap(), 8); + assert!(parallelism(Some("0")).is_err()); + assert!(parallelism(Some("many")).is_err()); + } + #[test] fn test_scan_timestamp_parses_local_time_and_truncates_to_millis() { let zone = jiff::tz::db().get("Asia/Shanghai").unwrap(); diff --git a/crates/paimon/src/table/table_scan.rs b/crates/paimon/src/table/table_scan.rs index 4165b592e..f4edba05b 100644 --- a/crates/paimon/src/table/table_scan.rs +++ b/crates/paimon/src/table/table_scan.rs @@ -241,6 +241,7 @@ async fn read_all_manifest_entries( bucket_function_type: BucketFunctionType, row_range_index: Option<&RowRangeIndex>, manifest_sidecar_enabled: bool, + manifest_parallelism: usize, trace: Option<&mut ScanTrace>, ) -> crate::Result> { let incremental = matches!(&source, ManifestListSource::AppendDeltas(_)); @@ -256,7 +257,7 @@ async fn read_all_manifest_entries( .collect::>(); let delta = futures::stream::iter(names) .map(|name| async move { read_manifest_list(file_io, table_path, &name).await }) - .buffered(64) + .buffered(manifest_parallelism) .try_fold(Vec::new(), |mut files, next| async move { files.extend(next); Ok(files) @@ -398,10 +399,10 @@ async fn read_all_manifest_entries( Ok::<_, crate::Error>((filtered, counters)) } }) - // Keep manifest read concurrency bounded. `try_fold` releases each - // yielded result after merging it, so peak retained results are the - // accumulator plus at most this bounded set of in-flight reads. - .buffered(64) + // Keep manifest read concurrency bounded by the table option. `try_fold` + // releases each yielded result after merging it, so peak retained results + // are the accumulator plus at most this bounded set of in-flight reads. + .buffered(manifest_parallelism) .try_fold( (Vec::new(), ManifestReadCounters::default()), |(mut all_entries, mut counters), (entries, manifest_counters)| async move { @@ -1672,6 +1673,7 @@ impl<'a> PaimonTableScan<'a> { bucket_function_type, row_range_index, core_options.manifest_sidecar_enabled(), + core_options.scan_manifest_parallelism()?, trace, ) .await?; @@ -6334,6 +6336,33 @@ mod tests { assert_eq!(changelog.snapshot_id(), Some(1)); } + #[tokio::test] + async fn test_plan_rejects_invalid_scan_manifest_parallelism() { + let table = + scan_trace_test_table("memory:/invalid_scan_manifest_parallelism").copy_with_options( + HashMap::from([("scan.manifest.parallelism".to_string(), "0".to_string())]), + ); + setup_scan_trace_dirs(&table).await; + + TableCommit::new(table.clone(), "manifest-parallelism-test".to_string()) + .commit(vec![CommitMessage::new( + BinaryRowBuilder::new(0).build_serialized(), + 0, + vec![stats_trace_file("parallelism.parquet", 1, 1)], + )]) + .await + .unwrap(); + + let error = table + .new_read_builder() + .new_scan() + .plan() + .await + .unwrap_err(); + assert!(matches!(error, Error::DataInvalid { message, .. } + if message.contains("scan.manifest.parallelism"))); + } + #[tokio::test] async fn test_plan_snapshot_metadata_preserves_pruned_historical_snapshot() { let table = scan_trace_small_split_table("memory:/plan_historical_snapshot_metadata");