From b9a141f3a96b32169fbfc847133e72aaa15f0c16 Mon Sep 17 00:00:00 2001 From: jianguotian Date: Wed, 16 Sep 2026 11:34:10 +0800 Subject: [PATCH 1/2] perf(manifest): write bucket-first sorted entries --- crates/paimon/src/spec/core_options.rs | 12 ++ crates/paimon/src/table/table_commit.rs | 216 ++++++++++++++++++++++++ 2 files changed, 228 insertions(+) diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index d3114fd4c..3a2adf7d7 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -88,6 +88,7 @@ const MANIFEST_COMPRESSION_OPTION: &str = "manifest.compression"; const MANIFEST_TARGET_FILE_SIZE_OPTION: &str = "manifest.target-file-size"; const MANIFEST_TARGET_SIZE_OPTION: &str = "manifest.target-size"; const MANIFEST_MERGE_MIN_COUNT_OPTION: &str = "manifest.merge-min-count"; +const MANIFEST_SORT_ENABLED_OPTION: &str = "manifest-sort.enabled"; const WRITE_PARQUET_BUFFER_SIZE_OPTION: &str = "write.parquet-buffer-size"; const READ_BATCH_SIZE_OPTION: &str = "read.batch-size"; const PARQUET_ROW_GROUP_PARALLELISM_OPTION: &str = "read.parquet.row-group.parallelism"; @@ -1165,6 +1166,14 @@ impl<'a> CoreOptions<'a> { .unwrap_or(DEFAULT_MANIFEST_MERGE_MIN_COUNT) } + /// Whether manifest compaction writes entries in a pruning-friendly sort + /// order. Disabled by default, matching Java Paimon. + pub fn manifest_sort_enabled(&self) -> bool { + self.options + .get(MANIFEST_SORT_ENABLED_OPTION) + .is_some_and(|value| value.eq_ignore_ascii_case("true")) + } + /// Number of buckets for the table. Default is -1 (dynamic bucket). pub fn bucket(&self) -> i32 { self.options @@ -2562,6 +2571,7 @@ mod tests { assert_eq!(core.manifest_compression(), "zstd"); assert_eq!(core.manifest_target_size(), 8 * 1024 * 1024); assert_eq!(core.manifest_merge_min_count(), 30); + assert!(!core.manifest_sort_enabled()); } #[test] @@ -2579,6 +2589,7 @@ mod tests { ), (MANIFEST_COMPRESSION_OPTION.to_string(), "null".to_string()), (MANIFEST_MERGE_MIN_COUNT_OPTION.to_string(), "3".to_string()), + (MANIFEST_SORT_ENABLED_OPTION.to_string(), "true".to_string()), ]); let core = CoreOptions::new(&options); assert_eq!(core.bucket(), 4); @@ -2590,6 +2601,7 @@ mod tests { assert_eq!(core.manifest_compression(), "null"); assert_eq!(core.manifest_target_size(), 1024); assert_eq!(core.manifest_merge_min_count(), 3); + assert!(core.manifest_sort_enabled()); } #[test] diff --git a/crates/paimon/src/table/table_commit.rs b/crates/paimon/src/table/table_commit.rs index 676d53587..4a5670f9e 100644 --- a/crates/paimon/src/table/table_commit.rs +++ b/crates/paimon/src/table/table_commit.rs @@ -116,6 +116,7 @@ pub struct TableCommit { manifest_compression: String, manifest_target_size: i64, manifest_merge_min_count: usize, + manifest_sort_enabled: bool, row_tracking_enabled: bool, data_evolution_enabled: bool, partition_default_name: String, @@ -140,6 +141,7 @@ impl TableCommit { let manifest_compression = core_options.manifest_compression().to_string(); let manifest_target_size = core_options.manifest_target_size(); let manifest_merge_min_count = core_options.manifest_merge_min_count(); + let manifest_sort_enabled = core_options.manifest_sort_enabled(); let row_tracking_enabled = core_options.row_tracking_enabled(); let data_evolution_enabled = core_options.data_evolution_enabled(); let partition_default_name = core_options.partition_default_name().to_string(); @@ -156,6 +158,7 @@ impl TableCommit { manifest_compression, manifest_target_size, manifest_merge_min_count, + manifest_sort_enabled, row_tracking_enabled, data_evolution_enabled, partition_default_name, @@ -1060,6 +1063,23 @@ impl TableCommit { return Ok(vec![]); } + // Java's bucket-first manifest layout is selected for fixed and + // postponed bucket tables when manifest sorting is enabled. Data + // evolution keeps its row-id-oriented layout and is intentionally not + // reordered here. + let mut sorted_entries = None; + if self.manifest_sort_enabled + && (self.total_buckets > 0 || self.total_buckets == POSTPONE_BUCKET) + && !self.data_evolution_enabled + { + let mut owned = entries.to_vec(); + let partition_fields = self.table.schema().partition_fields(); + let partition_sort_type = partition_fields.first().map(|field| field.data_type()); + sort_manifest_entries_bucket_first(&mut owned, partition_sort_type)?; + sorted_entries = Some(owned); + } + let entries = sorted_entries.as_deref().unwrap_or(entries); + let target_size = self.manifest_target_size.max(1) as usize; let mut result = Vec::new(); let mut chunk_start = 0usize; @@ -2931,6 +2951,54 @@ impl TableCommit { } } +fn sort_manifest_entries_bucket_first( + entries: &mut [ManifestEntry], + partition_sort_type: Option<&DataType>, +) -> Result<()> { + let mut keyed = entries + .iter() + .cloned() + .map(|entry| { + let partition_key = match partition_sort_type { + Some(data_type) if !entry.partition().is_empty() => { + let row = BinaryRow::from_serialized_bytes(entry.partition())?; + extract_datum(&row, 0, data_type)? + } + _ => None, + }; + Ok((entry, partition_key)) + }) + .collect::>>()?; + + keyed.sort_by(|(left, left_partition), (right, right_partition)| { + left.bucket() + .cmp(&right.bucket()) + .then_with(|| compare_partition_sort_keys(left_partition, right_partition)) + .then_with(|| file_kind_order(left.kind()).cmp(&file_kind_order(right.kind()))) + .then_with(|| left.file().file_name.cmp(&right.file().file_name)) + }); + for (target, (entry, _)) in entries.iter_mut().zip(keyed) { + *target = entry; + } + Ok(()) +} + +fn compare_partition_sort_keys(left: &Option, right: &Option) -> std::cmp::Ordering { + match (left, right) { + (None, None) => std::cmp::Ordering::Equal, + (None, Some(_)) => std::cmp::Ordering::Less, + (Some(_), None) => std::cmp::Ordering::Greater, + (Some(left), Some(right)) => left.partial_cmp(right).unwrap_or(std::cmp::Ordering::Equal), + } +} + +fn file_kind_order(kind: &FileKind) -> u8 { + match kind { + FileKind::Add => 0, + FileKind::Delete => 1, + } +} + /// Serialized BinaryRow for partition stats; unlike `datums_to_binary_row`, returns a /// valid arity-N row even when every datum is `None` (the all-null case must still /// decode on the Java side). @@ -3311,6 +3379,84 @@ mod tests { } } + #[test] + fn test_manifest_entries_sort_bucket_first() { + let mut entries = vec![ + ManifestEntry::new( + FileKind::Delete, + vec![2], + 1, + 4, + test_data_file("delete-bucket-1.parquet", 1), + 2, + ), + ManifestEntry::new( + FileKind::Add, + vec![2], + 0, + 4, + test_data_file("add-partition-2.parquet", 1), + 2, + ), + ManifestEntry::new( + FileKind::Delete, + vec![1], + 0, + 4, + test_data_file("delete-partition-1.parquet", 1), + 2, + ), + ManifestEntry::new( + FileKind::Add, + vec![1], + 0, + 4, + test_data_file("add-partition-1.parquet", 1), + 2, + ), + ]; + + sort_manifest_entries_bucket_first(&mut entries, None).unwrap(); + + assert_eq!( + entries + .iter() + .map(|entry| ( + entry.bucket(), + entry.partition().to_vec(), + *entry.kind(), + entry.file().file_name.clone(), + )) + .collect::>(), + vec![ + ( + 0, + vec![1], + FileKind::Add, + "add-partition-1.parquet".to_string(), + ), + ( + 0, + vec![2], + FileKind::Add, + "add-partition-2.parquet".to_string(), + ), + ( + 0, + vec![1], + FileKind::Delete, + "delete-partition-1.parquet".to_string(), + ), + ( + 1, + vec![2], + FileKind::Delete, + "delete-bucket-1.parquet".to_string(), + ), + ] + ); + } + fn test_global_index_file( name: &str, index_field_id: i32, @@ -5926,6 +6072,76 @@ mod tests { assert!(file_names.contains("data-2499.parquet")); } + #[tokio::test] + async fn test_manifest_sort_option_writes_bucket_first() { + let file_io = test_file_io(); + let table_path = "memory:/test_manifest_bucket_sort"; + setup_dirs(&file_io, table_path).await; + + let table = test_table_with_options( + &file_io, + table_path, + HashMap::from([ + ("bucket".to_string(), "2".to_string()), + ("manifest-sort.enabled".to_string(), "true".to_string()), + ]), + ); + let commit = TableCommit::new(table, "test-user".to_string()); + let entries = vec![ + ManifestEntry::new( + FileKind::Add, + vec![], + 1, + 2, + test_data_file("bucket-1.parquet", 1), + 2, + ), + ManifestEntry::new( + FileKind::Add, + vec![], + 0, + 2, + test_data_file("bucket-0-b.parquet", 1), + 2, + ), + ManifestEntry::new( + FileKind::Add, + vec![], + 0, + 2, + test_data_file("bucket-0-a.parquet", 1), + 2, + ), + ]; + + let manifest_dir = format!("{table_path}/manifest"); + let metas = commit + .write_manifest_files(&file_io, &manifest_dir, "manifest-test", &entries) + .await + .unwrap(); + assert_eq!(metas.len(), 1); + assert_eq!(metas[0].min_bucket(), Some(0)); + assert_eq!(metas[0].max_bucket(), Some(1)); + + let written = Manifest::read( + &file_io, + &format!("{manifest_dir}/{}", metas[0].file_name()), + ) + .await + .unwrap(); + assert_eq!( + written + .iter() + .map(|entry| (entry.bucket(), entry.file().file_name.as_str())) + .collect::>(), + vec![ + (0, "bucket-0-a.parquet"), + (0, "bucket-0-b.parquet"), + (1, "bucket-1.parquet"), + ] + ); + } + #[tokio::test] async fn test_manifest_rolling_waits_for_java_check_cadence() { let file_io = test_file_io(); From dfe51233cc62da6d8ec2de23c2c50d262b1a2af4 Mon Sep 17 00:00:00 2001 From: jianguotian Date: Wed, 16 Sep 2026 13:46:39 +0800 Subject: [PATCH 2/2] fix(manifest): move entries through bucket sorting --- crates/paimon/src/table/table_commit.rs | 80 ++++++++++++------------- 1 file changed, 37 insertions(+), 43 deletions(-) diff --git a/crates/paimon/src/table/table_commit.rs b/crates/paimon/src/table/table_commit.rs index 4a5670f9e..f3079d3c6 100644 --- a/crates/paimon/src/table/table_commit.rs +++ b/crates/paimon/src/table/table_commit.rs @@ -906,13 +906,32 @@ impl TableCommit { let delta_manifest_list_path = format!("{manifest_dir}/{delta_manifest_list_name}"); let changelog_manifest_list_path = format!("{manifest_dir}/{changelog_manifest_list_name}"); + // Compute metadata before the owned entries are moved into manifest + // writing. Sorting consumes the vectors so it does not need to clone + // every ManifestEntry and its heap-backed fields. + let mut delta_record_count: i64 = 0; + for entry in &resolved.entries { + match entry.kind() { + FileKind::Add => delta_record_count += entry.file().row_count, + FileKind::Delete => delta_record_count -= entry.file().row_count, + } + } + let statistics = self.generate_partition_statistics(&resolved.entries)?; + let changelog_record_count = (!resolved.changelog_entries.is_empty()).then(|| { + resolved + .changelog_entries + .iter() + .map(|entry| entry.file().row_count) + .sum() + }); + // Write delta manifest files, rolling by target size. let new_manifest_file_metas = self .write_manifest_files( file_io, &manifest_dir, &new_manifest_prefix, - &resolved.entries, + resolved.entries, ) .await?; @@ -926,7 +945,7 @@ impl TableCommit { .await?; let (changelog_record_count, changelog_manifest_list_size) = - if resolved.changelog_entries.is_empty() { + if changelog_record_count.is_none() { (None, None) } else { let changelog_manifest_file_metas = self @@ -934,7 +953,7 @@ impl TableCommit { file_io, &manifest_dir, &changelog_manifest_prefix, - &resolved.changelog_entries, + resolved.changelog_entries, ) .await?; ManifestList::write_with_compression( @@ -945,16 +964,7 @@ impl TableCommit { ) .await?; let status = file_io.get_status(&changelog_manifest_list_path).await?; - ( - Some( - resolved - .changelog_entries - .iter() - .map(|entry| entry.file().row_count) - .sum(), - ), - Some(status.size as i64), - ) + (changelog_record_count, Some(status.size as i64)) }; // Read existing manifests (base + delta from previous snapshot) and write base manifest list @@ -986,14 +996,6 @@ impl TableCommit { ) .await?; - // Calculate delta record count - let mut delta_record_count: i64 = 0; - for entry in &resolved.entries { - match entry.kind() { - FileKind::Add => delta_record_count += entry.file().row_count, - FileKind::Delete => delta_record_count -= entry.file().row_count, - } - } total_record_count += delta_record_count; let snapshot = Snapshot::builder() @@ -1015,8 +1017,6 @@ impl TableCommit { .index_manifest(resolved.index_manifest_name) .build(); - let statistics = self.generate_partition_statistics(&resolved.entries)?; - if self.snapshot_commit.commit(&snapshot, &statistics).await? { Ok(CommitAttemptResult::Success) } else { @@ -1057,7 +1057,7 @@ impl TableCommit { file_io: &FileIO, manifest_dir: &str, name_prefix: &str, - entries: &[ManifestEntry], + entries: Vec, ) -> Result> { if entries.is_empty() { return Ok(vec![]); @@ -1067,18 +1067,16 @@ impl TableCommit { // postponed bucket tables when manifest sorting is enabled. Data // evolution keeps its row-id-oriented layout and is intentionally not // reordered here. - let mut sorted_entries = None; - if self.manifest_sort_enabled + let entries = if self.manifest_sort_enabled && (self.total_buckets > 0 || self.total_buckets == POSTPONE_BUCKET) && !self.data_evolution_enabled { - let mut owned = entries.to_vec(); let partition_fields = self.table.schema().partition_fields(); let partition_sort_type = partition_fields.first().map(|field| field.data_type()); - sort_manifest_entries_bucket_first(&mut owned, partition_sort_type)?; - sorted_entries = Some(owned); - } - let entries = sorted_entries.as_deref().unwrap_or(entries); + sort_manifest_entries_bucket_first(entries, partition_sort_type)? + } else { + entries + }; let target_size = self.manifest_target_size.max(1) as usize; let mut result = Vec::new(); @@ -1217,7 +1215,7 @@ impl TableCommit { let manifest_prefix = format!("manifest-{}", uuid::Uuid::new_v4()); let merged_metas = self - .write_manifest_files(file_io, manifest_dir, &manifest_prefix, &merged_entries) + .write_manifest_files(file_io, manifest_dir, &manifest_prefix, merged_entries) .await?; result.extend(merged_metas.clone()); new_files.extend(merged_metas); @@ -2952,12 +2950,11 @@ impl TableCommit { } fn sort_manifest_entries_bucket_first( - entries: &mut [ManifestEntry], + entries: Vec, partition_sort_type: Option<&DataType>, -) -> Result<()> { +) -> Result> { let mut keyed = entries - .iter() - .cloned() + .into_iter() .map(|entry| { let partition_key = match partition_sort_type { Some(data_type) if !entry.partition().is_empty() => { @@ -2977,10 +2974,7 @@ fn sort_manifest_entries_bucket_first( .then_with(|| file_kind_order(left.kind()).cmp(&file_kind_order(right.kind()))) .then_with(|| left.file().file_name.cmp(&right.file().file_name)) }); - for (target, (entry, _)) in entries.iter_mut().zip(keyed) { - *target = entry; - } - Ok(()) + Ok(keyed.into_iter().map(|(entry, _)| entry).collect()) } fn compare_partition_sort_keys(left: &Option, right: &Option) -> std::cmp::Ordering { @@ -3381,7 +3375,7 @@ mod tests { #[test] fn test_manifest_entries_sort_bucket_first() { - let mut entries = vec![ + let entries = vec![ ManifestEntry::new( FileKind::Delete, vec![2], @@ -3416,7 +3410,7 @@ mod tests { ), ]; - sort_manifest_entries_bucket_first(&mut entries, None).unwrap(); + let entries = sort_manifest_entries_bucket_first(entries, None).unwrap(); assert_eq!( entries @@ -6116,7 +6110,7 @@ mod tests { let manifest_dir = format!("{table_path}/manifest"); let metas = commit - .write_manifest_files(&file_io, &manifest_dir, "manifest-test", &entries) + .write_manifest_files(&file_io, &manifest_dir, "manifest-test", entries) .await .unwrap(); assert_eq!(metas.len(), 1);