Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
55 changes: 55 additions & 0 deletions crates/paimon/src/spec/core_options.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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<usize> {
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::<usize>()
.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<usize> {
let Some(raw) = self.options.get(MOSAIC_READ_PREFETCH_ROW_GROUPS_OPTION) else {
Expand Down Expand Up @@ -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();
Expand Down
39 changes: 34 additions & 5 deletions crates/paimon/src/table/table_scan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Vec<ManifestEntry>> {
let incremental = matches!(&source, ManifestListSource::AppendDeltas(_));
Expand All @@ -256,7 +257,7 @@ async fn read_all_manifest_entries(
.collect::<Vec<_>>();
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)
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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?;
Expand Down Expand Up @@ -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");
Expand Down
Loading