This is an automated email from the ASF dual-hosted git repository. JingsongLi pushed a commit to branch release-2.0 in repository https://gitbox.apache.org/repos/asf/paimon.git
commit ddf4f7c8777a61ad206b4e132702a05e2cfdb7f8 Author: Jingsong Lee <[email protected]> AuthorDate: Mon Aug 3 13:13:29 2026 +0800 [core] Fix forced L0 compaction level selection (#8993) --- .../mergetree/compact/UniversalCompaction.java | 15 +++- .../compact/ForceUpLevel0CompactionTest.java | 19 ++++++ .../paimon/table/PrimaryKeySimpleTableTest.java | 79 ++++++++++++++++++++++ 3 files changed, 110 insertions(+), 3 deletions(-) diff --git a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/UniversalCompaction.java b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/UniversalCompaction.java index a40e985e9d..396f3c89a2 100644 --- a/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/UniversalCompaction.java +++ b/paimon-core/src/main/java/org/apache/paimon/mergetree/compact/UniversalCompaction.java @@ -116,9 +116,18 @@ public class UniversalCompaction implements CompactStrategy { candidateCount++; } - return candidateCount == 0 - ? Optional.empty() - : Optional.of(pickForSizeRatio(numLevels - 1, runs, candidateCount, true)); + if (candidateCount == 0) { + return Optional.empty(); + } + + // Level 1 must be compacted with level 0 to avoid producing level 0. Include it in the + // initial candidates so that the size-ratio check also considers level 2 based on the + // combined size of level 0 and level 1. + if (candidateCount < runs.size() && runs.get(candidateCount).level() == 1) { + candidateCount++; + } + + return Optional.of(pickForSizeRatio(numLevels - 1, runs, candidateCount, true)); } @VisibleForTesting diff --git a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/ForceUpLevel0CompactionTest.java b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/ForceUpLevel0CompactionTest.java index 126ce000f0..5011b1734b 100644 --- a/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/ForceUpLevel0CompactionTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/mergetree/compact/ForceUpLevel0CompactionTest.java @@ -60,6 +60,25 @@ public class ForceUpLevel0CompactionTest { assertThat(result.get().outputLevel()).isEqualTo(2); } + @Test + public void testForceCompaction0ConsidersLevel2() { + ForceUpLevel0Compaction compaction = + new ForceUpLevel0Compaction(ofTesting(200, 1, 5), null); + + Optional<CompactUnit> result = + compaction.pick(3, Arrays.asList(run(0, 1), run(1, 99), run(2, 100))); + + assertThat(result).isPresent(); + assertThat(result.get().files()).hasSize(3); + assertThat(result.get().outputLevel()).isEqualTo(2); + + result = compaction.pick(3, Arrays.asList(run(0, 1), run(1, 99), run(2, 102))); + + assertThat(result).isPresent(); + assertThat(result.get().files()).hasSize(2); + assertThat(result.get().outputLevel()).isEqualTo(1); + } + private LevelSortedRun run(int level, int size) { return new LevelSortedRun(level, SortedRun.fromSingle(file(size))); } diff --git a/paimon-core/src/test/java/org/apache/paimon/table/PrimaryKeySimpleTableTest.java b/paimon-core/src/test/java/org/apache/paimon/table/PrimaryKeySimpleTableTest.java index ea8f096a53..4bda9faf3e 100644 --- a/paimon-core/src/test/java/org/apache/paimon/table/PrimaryKeySimpleTableTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/table/PrimaryKeySimpleTableTest.java @@ -2133,6 +2133,85 @@ public class PrimaryKeySimpleTableTest extends SimpleTableTestBase { commit.close(); } + @Test + public void testForceUpLevel0CompactionConsidersLevel2() throws Exception { + FileStoreTable table = + createFileStoreTable( + options -> { + options.set(CoreOptions.NUM_LEVELS, 3); + options.set(CoreOptions.NUM_SORTED_RUNS_COMPACTION_TRIGGER, 10); + options.set(CoreOptions.COMPACTION_FORCE_UP_LEVEL_0, true); + }); + FileStoreTable writeOnlyTable = + table.copy(singletonMap(CoreOptions.WRITE_ONLY.key(), "true")); + + // Build a large level 2 file. + writeRowsWithoutCompaction(writeOnlyTable, 0, 0, 1000); + compactPartition(table, 1, true); + + List<DataFileMeta> files = currentDataFiles(table); + assertThat(files).singleElement().satisfies(file -> assertThat(file.level()).isEqualTo(2)); + + // Build a slightly smaller level 1 file without including level 2. + writeRowsWithoutCompaction(writeOnlyTable, 2, 1000, 900); + compactPartition(table, 3, false); + + files = currentDataFiles(table); + assertThat(files).extracting(DataFileMeta::level).containsExactlyInAnyOrder(1, 2); + + // Add a small level 0 file. Level 0 alone cannot pick level 1, but level 0 and level 1 + // together are large enough to pick level 2. + writeRowsWithoutCompaction(writeOnlyTable, 4, 1900, 200); + + files = currentDataFiles(table); + assertThat(files).extracting(DataFileMeta::level).containsExactlyInAnyOrder(0, 1, 2); + Map<Integer, Long> fileSizeByLevel = + files.stream() + .collect(Collectors.toMap(DataFileMeta::level, DataFileMeta::fileSize)); + long level0Size = fileSizeByLevel.get(0); + long level1Size = fileSizeByLevel.get(1); + long level2Size = fileSizeByLevel.get(2); + assertThat(level0Size * 101).isLessThan(level1Size * 100); + assertThat((level0Size + level1Size) * 101).isGreaterThanOrEqualTo(level2Size * 100); + + compactPartition(table, 5, false); + + files = currentDataFiles(table); + assertThat(files) + .singleElement() + .satisfies( + file -> { + assertThat(file.level()).isEqualTo(2); + assertThat(file.rowCount()).isEqualTo(2100); + }); + } + + private void writeRowsWithoutCompaction( + FileStoreTable table, long commitIdentifier, int start, int count) throws Exception { + try (StreamTableWrite write = table.newWrite(commitUser); + StreamTableCommit commit = table.newCommit(commitUser)) { + for (int i = start; i < start + count; i++) { + write.write(rowData(1, i, (long) i)); + } + commit.commit(commitIdentifier, write.prepareCommit(true, commitIdentifier)); + } + } + + private void compactPartition( + FileStoreTable table, long commitIdentifier, boolean fullCompaction) throws Exception { + try (StreamTableWrite write = table.newWrite(commitUser); + StreamTableCommit commit = table.newCommit(commitUser)) { + write.compact(binaryRow(1), 0, fullCompaction); + commit.commit(commitIdentifier, write.prepareCommit(true, commitIdentifier)); + } + } + + private List<DataFileMeta> currentDataFiles(FileStoreTable table) throws Exception { + return table.newSnapshotReader().read().dataSplits().stream() + .flatMap(split -> split.dataFiles().stream()) + .collect(Collectors.toList()); + } + @Test public void testStreamingReadOptimizedTable() throws Exception { FileStoreTable table =
