Skip to content

Commit 7768824

Browse files
committed
feat: support manifest compaction in table commits
1 parent b94a1fa commit 7768824

2 files changed

Lines changed: 559 additions & 4 deletions

File tree

crates/paimon/src/spec/core_options.rs

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,10 @@ const CHANGELOG_FILE_PREFIX_OPTION: &str = "changelog-file.prefix";
4545
const CHANGELOG_FILE_FORMAT_OPTION: &str = "changelog-file.format";
4646
const CHANGELOG_FILE_COMPRESSION_OPTION: &str = "changelog-file.compression";
4747
const CHANGELOG_FILE_STATS_MODE_OPTION: &str = "changelog-file.stats-mode";
48+
const MANIFEST_TARGET_FILE_SIZE_OPTION: &str = "manifest.target-file-size";
49+
const MANIFEST_FULL_COMPACTION_THRESHOLD_SIZE_OPTION: &str =
50+
"manifest.full-compaction-threshold-size";
51+
const MANIFEST_MERGE_MIN_COUNT_OPTION: &str = "manifest.merge-min-count";
4852
const ROW_TRACKING_ENABLED_OPTION: &str = "row-tracking.enabled";
4953
const WRITE_PARQUET_BUFFER_SIZE_OPTION: &str = "write.parquet-buffer-size";
5054
pub(crate) const SEQUENCE_FIELD_OPTION: &str = "sequence.field";
@@ -65,6 +69,9 @@ const DEFAULT_SOURCE_SPLIT_OPEN_FILE_COST: i64 = 4 * 1024 * 1024;
6569
const DEFAULT_PARTITION_DEFAULT_NAME: &str = "__DEFAULT_PARTITION__";
6670
const DEFAULT_CHANGELOG_FILE_PREFIX: &str = "changelog-";
6771
const DEFAULT_TARGET_FILE_SIZE: i64 = 256 * 1024 * 1024;
72+
const DEFAULT_MANIFEST_TARGET_FILE_SIZE: i64 = 8 * 1024 * 1024;
73+
const DEFAULT_MANIFEST_FULL_COMPACTION_THRESHOLD_SIZE: i64 = 16 * 1024 * 1024;
74+
const DEFAULT_MANIFEST_MERGE_MIN_COUNT: usize = 30;
6875
const DEFAULT_WRITE_PARQUET_BUFFER_SIZE: i64 = 256 * 1024 * 1024;
6976
const DYNAMIC_BUCKET_TARGET_ROW_NUM_OPTION: &str = "dynamic-bucket.target-row-num";
7077
const DEFAULT_DYNAMIC_BUCKET_TARGET_ROW_NUM: i64 = 200_000;
@@ -414,6 +421,37 @@ impl<'a> CoreOptions<'a> {
414421
.unwrap_or_else(|| self.target_file_size())
415422
}
416423

424+
/// Suggested file size of a manifest file.
425+
///
426+
/// Corresponds to Java `CoreOptions.MANIFEST_TARGET_FILE_SIZE`.
427+
pub fn manifest_target_file_size(&self) -> i64 {
428+
self.options
429+
.get(MANIFEST_TARGET_FILE_SIZE_OPTION)
430+
.and_then(|v| parse_memory_size(v))
431+
.unwrap_or(DEFAULT_MANIFEST_TARGET_FILE_SIZE)
432+
}
433+
434+
/// Size threshold for triggering full compaction of manifests.
435+
///
436+
/// Corresponds to Java `CoreOptions.MANIFEST_FULL_COMPACTION_FILE_SIZE`.
437+
pub fn manifest_full_compaction_threshold_size(&self) -> i64 {
438+
self.options
439+
.get(MANIFEST_FULL_COMPACTION_THRESHOLD_SIZE_OPTION)
440+
.and_then(|v| parse_memory_size(v))
441+
.unwrap_or(DEFAULT_MANIFEST_FULL_COMPACTION_THRESHOLD_SIZE)
442+
}
443+
444+
/// Minimum number of trailing manifest files to merge.
445+
///
446+
/// Corresponds to Java `CoreOptions.MANIFEST_MERGE_MIN_COUNT`.
447+
pub fn manifest_merge_min_count(&self) -> usize {
448+
self.options
449+
.get(MANIFEST_MERGE_MIN_COUNT_OPTION)
450+
.and_then(|v| v.parse::<usize>().ok())
451+
.filter(|count| *count > 0)
452+
.unwrap_or(DEFAULT_MANIFEST_MERGE_MIN_COUNT)
453+
}
454+
417455
/// File format for data files (e.g. "parquet", "orc", "avro", "vortex").
418456
/// Default is "parquet".
419457
pub fn file_format(&self) -> &str {
@@ -619,6 +657,34 @@ mod tests {
619657
assert_eq!(parse_memory_size("abc"), None);
620658
}
621659

660+
#[test]
661+
fn test_manifest_compaction_options_match_java_defaults_and_overrides() {
662+
let options = HashMap::new();
663+
let core = CoreOptions::new(&options);
664+
assert_eq!(core.manifest_target_file_size(), 8 * 1024 * 1024);
665+
assert_eq!(
666+
core.manifest_full_compaction_threshold_size(),
667+
16 * 1024 * 1024
668+
);
669+
assert_eq!(core.manifest_merge_min_count(), 30);
670+
671+
let options = HashMap::from([
672+
(
673+
MANIFEST_TARGET_FILE_SIZE_OPTION.to_string(),
674+
"500B".to_string(),
675+
),
676+
(
677+
MANIFEST_FULL_COMPACTION_THRESHOLD_SIZE_OPTION.to_string(),
678+
"200B".to_string(),
679+
),
680+
(MANIFEST_MERGE_MIN_COUNT_OPTION.to_string(), "3".to_string()),
681+
]);
682+
let core = CoreOptions::new(&options);
683+
assert_eq!(core.manifest_target_file_size(), 500);
684+
assert_eq!(core.manifest_full_compaction_threshold_size(), 200);
685+
assert_eq!(core.manifest_merge_min_count(), 3);
686+
}
687+
622688
#[test]
623689
fn test_partition_options_defaults() {
624690
let options = HashMap::new();

0 commit comments

Comments
 (0)