This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git
The following commit(s) were added to refs/heads/main by this push:
new ded9e59f fix: remove manifest minor compaction (#899)
ded9e59f is described below
commit ded9e59f8c0299643cc6d8d6993ba21e4ba45704
Author: XiaoHongbo <[email protected]>
AuthorDate: Mon Sep 21 17:22:14 2026 +0800
fix: remove manifest minor compaction (#899)
---
crates/paimon/src/spec/core_options.rs | 3 +-
crates/paimon/src/spec/manifest.rs | 3 +-
crates/paimon/src/table/table_commit.rs | 181 ++++++++++----------------------
3 files changed, 56 insertions(+), 131 deletions(-)
diff --git a/crates/paimon/src/spec/core_options.rs
b/crates/paimon/src/spec/core_options.rs
index 58131b6c..0073f172 100644
--- a/crates/paimon/src/spec/core_options.rs
+++ b/crates/paimon/src/spec/core_options.rs
@@ -1223,8 +1223,7 @@ impl<'a> CoreOptions<'a> {
.unwrap_or(DEFAULT_MANIFEST_TARGET_FILE_SIZE)
}
- /// Minimum number of small manifest files required before minor manifest
- /// compaction rewrites them into a new rolling manifest set.
+ /// Compatibility option; Rust commits do not currently compact manifests.
pub fn manifest_merge_min_count(&self) -> usize {
self.options
.get(MANIFEST_MERGE_MIN_COUNT_OPTION)
diff --git a/crates/paimon/src/spec/manifest.rs
b/crates/paimon/src/spec/manifest.rs
index ed2fa162..27485db5 100644
--- a/crates/paimon/src/spec/manifest.rs
+++ b/crates/paimon/src/spec/manifest.rs
@@ -88,8 +88,7 @@ impl Manifest {
}
}
-/// Merge ADD/DELETE entries by file identifier, returning only the active ADD
set.
-/// Mirrors Java
[FileEntry.mergeEntries](https://github.com/apache/paimon/blob/release-1.4/paimon-core/src/main/java/org/apache/paimon/manifest/FileEntry.java).
+/// Resolve a complete manifest set to its active ADD entries.
/// Return order is unspecified.
pub(crate) fn merge_active_entries(entries: Vec<ManifestEntry>) ->
Vec<ManifestEntry> {
use std::collections::HashMap;
diff --git a/crates/paimon/src/table/table_commit.rs
b/crates/paimon/src/table/table_commit.rs
index 2b747483..6d3f7738 100644
--- a/crates/paimon/src/table/table_commit.rs
+++ b/crates/paimon/src/table/table_commit.rs
@@ -115,7 +115,6 @@ pub struct TableCommit {
commit_max_retry_wait_ms: u64,
manifest_compression: String,
manifest_target_size: i64,
- manifest_merge_min_count: usize,
manifest_sidecar_enabled: bool,
row_tracking_enabled: bool,
data_evolution_enabled: bool,
@@ -140,7 +139,6 @@ impl TableCommit {
let commit_max_retry_wait_ms = core_options.commit_max_retry_wait_ms();
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_sidecar_enabled = core_options.manifest_sidecar_enabled();
let row_tracking_enabled = core_options.row_tracking_enabled();
let data_evolution_enabled = core_options.data_evolution_enabled();
@@ -157,7 +155,6 @@ impl TableCommit {
commit_max_retry_wait_ms,
manifest_compression,
manifest_target_size,
- manifest_merge_min_count,
manifest_sidecar_enabled,
row_tracking_enabled,
data_evolution_enabled,
@@ -991,14 +988,10 @@ impl TableCommit {
vec![]
};
- let (base_manifest_files, _merge_new_files) = self
- .merge_manifest_files(file_io, &manifest_dir,
existing_manifest_files)
- .await?;
-
ManifestList::write_with_compression(
file_io,
&base_manifest_list_path,
- &base_manifest_files,
+ &existing_manifest_files,
&self.manifest_compression,
)
.await?;
@@ -1139,91 +1132,6 @@ impl TableCommit {
Ok(result)
}
- /// Minor-compact existing manifest files before writing the base manifest
list.
- async fn merge_manifest_files(
- &self,
- file_io: &FileIO,
- manifest_dir: &str,
- manifest_files: Vec<ManifestFileMeta>,
- ) -> Result<(Vec<ManifestFileMeta>, Vec<ManifestFileMeta>)> {
- if manifest_files.is_empty() {
- return Ok((vec![], vec![]));
- }
-
- let target_size = self.manifest_target_size.max(1);
- let mut result = Vec::new();
- let mut new_files = Vec::new();
- let mut candidates = Vec::new();
- let mut total_size = 0i64;
-
- for manifest in manifest_files {
- total_size += manifest.file_size();
- candidates.push(manifest);
- if total_size >= target_size {
- self.merge_manifest_candidates(
- file_io,
- manifest_dir,
- &mut candidates,
- &mut result,
- &mut new_files,
- )
- .await?;
- total_size = 0;
- }
- }
-
- if candidates.len() >= self.manifest_merge_min_count {
- self.merge_manifest_candidates(
- file_io,
- manifest_dir,
- &mut candidates,
- &mut result,
- &mut new_files,
- )
- .await?;
- } else {
- result.append(&mut candidates);
- }
-
- Ok((result, new_files))
- }
-
- async fn merge_manifest_candidates(
- &self,
- file_io: &FileIO,
- manifest_dir: &str,
- candidates: &mut Vec<ManifestFileMeta>,
- result: &mut Vec<ManifestFileMeta>,
- new_files: &mut Vec<ManifestFileMeta>,
- ) -> Result<()> {
- if candidates.is_empty() {
- return Ok(());
- }
- if candidates.len() == 1 {
- result.append(candidates);
- return Ok(());
- }
-
- let mut entries = Vec::new();
- for manifest in candidates.drain(..) {
- let path = format!("{manifest_dir}/{}", manifest.file_name());
- entries.extend(Manifest::read(file_io, &path).await?);
- }
-
- let merged_entries = merge_active_entries(entries);
- if merged_entries.is_empty() {
- return Ok(());
- }
-
- let manifest_prefix = format!("manifest-{}", uuid::Uuid::new_v4());
- let merged_metas = self
- .write_manifest_files(file_io, manifest_dir, &manifest_prefix,
&merged_entries)
- .await?;
- result.extend(merged_metas.clone());
- new_files.extend(merged_metas);
- Ok(())
- }
-
/// Write already-encoded manifest bytes and return metadata for the
corresponding entries.
async fn write_manifest_file_bytes(
&self,
@@ -6068,9 +5976,9 @@ mod tests {
}
#[tokio::test]
- async fn test_minor_compaction_nets_add_delete_manifest_entries() {
+ async fn test_commit_preserves_delete_outside_retained_manifest() {
let file_io = test_file_io();
- let table_path = "memory:/test_minor_manifest_compaction";
+ let table_path = "memory:/test_commit_preserves_delete";
setup_dirs(&file_io, table_path).await;
let table = test_table_with_options(
@@ -6078,40 +5986,38 @@ mod tests {
table_path,
HashMap::from([("manifest.merge-min-count".to_string(),
"2".to_string())]),
);
- let commit = TableCommit::new(table, "test-user".to_string());
+ let mut commit = TableCommit::new(table, "test-user".to_string());
+ let partition = vec![0, 0, 0, 0];
+ let old_file = test_data_file("old.parquet", 1);
+ let mut initial_files = vec![old_file.clone()];
+ initial_files.extend(
+ (0..128)
+ .map(|_| test_data_file(&format!("filler-{}.parquet",
uuid::Uuid::new_v4()), 1)),
+ );
commit
.commit(vec![CommitMessage::new(
- vec![],
+ partition.clone(),
0,
- vec![test_data_file("data-0.parquet", 100)],
+ initial_files,
)])
.await
.unwrap();
- commit
- .overwrite(
- vec![CommitMessage::new(
- vec![],
- 0,
- vec![test_data_file("data-1.parquet", 50)],
- )],
- None,
- )
- .await
- .unwrap();
+ let mut deletion = CommitMessage::new(partition.clone(), 0, vec![]);
+ deletion.deleted_files.push(old_file);
+ commit.commit(vec![deletion]).await.unwrap();
commit
.commit(vec![CommitMessage::new(
- vec![],
+ partition.clone(),
0,
- vec![test_data_file("data-2.parquet", 25)],
+ vec![test_data_file("replacement.parquet", 1)],
)])
.await
.unwrap();
let snapshot = latest_snapshot(&file_io, table_path).await.unwrap();
- assert_eq!(snapshot.id(), 3);
let manifest_dir = format!("{table_path}/manifest");
let base_metas = ManifestList::read(
&file_io,
@@ -6119,31 +6025,52 @@ mod tests {
)
.await
.unwrap();
- assert_eq!(
- base_metas.len(),
- 1,
- "two previous manifest files should be minor-compacted into one
base manifest"
- );
-
- let base_entries = Manifest::read(
+ assert_eq!(base_metas.len(), 2);
+ let delta_metas = ManifestList::read(
&file_io,
- &format!("{manifest_dir}/{}", base_metas[0].file_name()),
+ &format!("{manifest_dir}/{}", snapshot.delta_manifest_list()),
)
.await
.unwrap();
- assert_eq!(base_entries.len(), 1);
- assert_eq!(*base_entries[0].kind(), FileKind::Add);
- assert_eq!(base_entries[0].file().file_name, "data-1.parquet");
+ let retained = base_metas
+ .iter()
+ .find(|meta| meta.num_added_files() > 1)
+ .unwrap();
+ let deletion = base_metas
+ .iter()
+ .find(|meta| meta.num_deleted_files() == 1)
+ .unwrap();
+ let merge_target = deletion.file_size() + delta_metas[0].file_size() +
1;
+ assert!(retained.file_size() >= merge_target);
+
+ commit.manifest_target_size = merge_target;
+ commit
+ .commit(vec![CommitMessage::new(
+ partition,
+ 0,
+ vec![test_data_file("unrelated.parquet", 1)],
+ )])
+ .await
+ .unwrap();
+ let snapshot = latest_snapshot(&file_io, table_path).await.unwrap();
+ let base_metas = ManifestList::read(
+ &file_io,
+ &format!("{manifest_dir}/{}", snapshot.base_manifest_list()),
+ )
+ .await
+ .unwrap();
+ assert_eq!(base_metas.len(), 3);
+ assert!(base_metas.iter().any(|meta| meta.num_deleted_files() == 1));
let active_file_names = active_entries(&file_io, table_path, &snapshot)
.await
.into_iter()
.map(|entry| entry.file().file_name.clone())
.collect::<HashSet<_>>();
- assert_eq!(
- active_file_names,
- HashSet::from(["data-1.parquet".to_string(),
"data-2.parquet".to_string()])
- );
+ assert_eq!(active_file_names.len(), 130);
+ assert!(!active_file_names.contains("old.parquet"));
+ assert!(active_file_names.contains("replacement.parquet"));
+ assert!(active_file_names.contains("unrelated.parquet"));
}
/// `write_manifest_file` must aggregate min/max bucket and level across
entries so the