This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new 7c820f3db3 [core] Compact adjacent overlapping blob ranges together 
(#9175)
7c820f3db3 is described below

commit 7c820f3db37900e54ab3af145196fecfe682240f
Author: YeJunHao <[email protected]>
AuthorDate: Tue Aug 11 21:43:38 2026 +0800

    [core] Compact adjacent overlapping blob ranges together (#9175)
---
 .../DataEvolutionCompactCoordinator.java           | 78 ++++++++++------------
 .../DataEvolutionCompactCoordinatorTest.java       | 44 ++++++++++++
 2 files changed, 79 insertions(+), 43 deletions(-)

diff --git 
a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactCoordinator.java
 
b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactCoordinator.java
index b32af63268..f8a0dea7fe 100644
--- 
a/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactCoordinator.java
+++ 
b/paimon-core/src/main/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactCoordinator.java
@@ -570,64 +570,56 @@ public class DataEvolutionCompactCoordinator {
 
             RangeHelper<DataFileMeta> rangeHelper =
                     new RangeHelper<>(DataFileMeta::nonNullRowIdRange);
-            List<DataFileMeta> smallFileCandidates = new ArrayList<>();
+            List<List<DataFileMeta>> continuousOrOverlapFiles = new 
ArrayList<>();
+            long expectedFirstRowId = -1L;
             for (List<DataFileMeta> rowRangeGroup :
                     rangeHelper.mergeOverlappingRanges(sortedFiles)) {
-                if (rowRangeGroup.size() >= BLOB_COMPACT_MIN_FILE_NUM) {
-                    rowRangeGroup.sort(
-                            comparingLong(DataFileMeta::nonNullFirstRowId)
-                                    
.thenComparingLong(DataFileMeta::maxSequenceNumber));
-                    result.add(rowRangeGroup);
-                } else {
-                    smallFileCandidates.add(rowRangeGroup.get(0));
-                }
-            }
-
-            result.addAll(smallFileGroupsToCompact(smallFileCandidates));
-            result.sort(comparingLong(group -> 
group.get(0).nonNullFirstRowId()));
-            return result;
-        }
-
-        private List<List<DataFileMeta>> 
smallFileGroupsToCompact(List<DataFileMeta> files) {
-            List<List<DataFileMeta>> result = new ArrayList<>();
-
-            List<DataFileMeta> continuousFiles = new ArrayList<>();
-            long expectedFirstRowId = -1;
-            for (DataFileMeta file : files) {
-                if (file.fileSize() >= blobTargetFileSize) {
-                    addFileGroupsToCompact(result, continuousFiles);
-                    continuousFiles.clear();
-                    expectedFirstRowId = -1;
-                    continue;
+                List<Range> rowRanges =
+                        rowRangeGroup.stream()
+                                .map(DataFileMeta::nonNullRowIdRange)
+                                .collect(Collectors.toList());
+                Range rowRange = Range.sortAndMergeOverlap(rowRanges).get(0);
+                long firstRowId = rowRange.from;
+                if (!continuousOrOverlapFiles.isEmpty() && firstRowId != 
expectedFirstRowId) {
+                    addFileGroupsToCompact(result, continuousOrOverlapFiles);
+                    continuousOrOverlapFiles.clear();
                 }
 
-                long firstRowId = file.nonNullFirstRowId();
-                if (!continuousFiles.isEmpty() && firstRowId != 
expectedFirstRowId) {
-                    addFileGroupsToCompact(result, continuousFiles);
-                    continuousFiles.clear();
-                }
-                continuousFiles.add(file);
-                expectedFirstRowId = firstRowId + file.rowCount();
+                continuousOrOverlapFiles.add(rowRangeGroup);
+                expectedFirstRowId = rowRange.to + 1;
             }
-            addFileGroupsToCompact(result, continuousFiles);
+            addFileGroupsToCompact(result, continuousOrOverlapFiles);
+            result.sort(comparingLong(group -> 
group.get(0).nonNullFirstRowId()));
             return result;
         }
 
         private void addFileGroupsToCompact(
-                List<List<DataFileMeta>> result, List<DataFileMeta> 
continuousFiles) {
-            if (continuousFiles.size() < BLOB_COMPACT_MIN_FILE_NUM) {
+                List<List<DataFileMeta>> result,
+                List<List<DataFileMeta>> continuousOrOverlapFiles) {
+            int compactFileCount = 
continuousOrOverlapFiles.stream().mapToInt(List::size).sum();
+            if (compactFileCount < BLOB_COMPACT_MIN_FILE_NUM) {
                 return;
             }
+
             List<DataFileMeta> taskFiles = new ArrayList<>();
-            long fileSize = 0L;
-            for (DataFileMeta file : continuousFiles) {
-                taskFiles.add(file);
-                fileSize += file.fileSize();
-                if (fileSize >= blobTargetFileSize
+            long taskFileSize = 0L;
+            for (List<DataFileMeta> fileGroup : continuousOrOverlapFiles) {
+                if (fileGroup.size() == 1 && fileGroup.get(0).fileSize() >= 
blobTargetFileSize) {
+                    if (taskFiles.size() >= BLOB_COMPACT_MIN_FILE_NUM) {
+                        result.add(taskFiles);
+                    }
+                    taskFiles = new ArrayList<>();
+                    taskFileSize = 0L;
+                    continue;
+                }
+
+                taskFiles.addAll(fileGroup);
+                taskFileSize += 
fileGroup.stream().mapToLong(DataFileMeta::fileSize).sum();
+                if (taskFileSize >= blobTargetFileSize
                         && taskFiles.size() >= BLOB_COMPACT_MIN_FILE_NUM) {
                     result.add(taskFiles);
                     taskFiles = new ArrayList<>();
-                    fileSize = 0L;
+                    taskFileSize = 0L;
                 }
             }
 
diff --git 
a/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactCoordinatorTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactCoordinatorTest.java
index 641390be9b..4f21fcbe0b 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactCoordinatorTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/append/dataevolution/DataEvolutionCompactCoordinatorTest.java
@@ -211,6 +211,50 @@ public class DataEvolutionCompactCoordinatorTest {
                 .containsExactly(entries.get(1).file(), entries.get(2).file());
     }
 
+    @Test
+    public void testCompactPlannerMergesAdjacentOverlappingBlobGroups() {
+        List<ManifestEntry> entries = new ArrayList<>();
+        entries.add(makeEntry("file1.parquet", 0L, 2L, 100));
+        entries.add(makeBlobEntry("old-prefix.blob", 0L, 1L, 100, 0, "pic"));
+        entries.add(makeBlobEntry("updated-prefix.blob", 0L, 1L, 100, 1, 
"pic"));
+        entries.add(makeBlobEntry("old-suffix.blob", 1L, 1L, 100, 0, "pic"));
+        entries.add(makeBlobEntry("updated-suffix.blob", 1L, 1L, 100, 1, 
"pic"));
+
+        DataEvolutionCompactCoordinator.CompactPlanner planner =
+                blobPlanner(250, 1, 2, rowType(new DataField(1, "pic", 
DataTypes.BLOB())));
+
+        List<DataEvolutionCompactTask> tasks = planner.compactPlan(entries);
+
+        assertThat(tasks).hasSize(1);
+        
assertThat(tasks.get(0).type()).isEqualTo(DataEvolutionCompactTask.TaskType.BLOB);
+        assertThat(tasks.get(0).compactBefore())
+                .containsExactly(
+                        entries.get(1).file(),
+                        entries.get(2).file(),
+                        entries.get(3).file(),
+                        entries.get(4).file());
+    }
+
+    @Test
+    public void testCompactPlannerUsesMergedOverlappingBlobRangeBoundary() {
+        List<ManifestEntry> entries = new ArrayList<>();
+        entries.add(makeEntry("file1.parquet", 0L, 11L, 100));
+        entries.add(makeBlobEntry("prefix.blob", 0L, 5L, 100, 0, "pic"));
+        entries.add(makeBlobEntry("overlap.blob", 3L, 7L, 100, 1, "pic"));
+        entries.add(makeBlobEntry("suffix.blob", 10L, 1L, 100, 0, "pic"));
+
+        DataEvolutionCompactCoordinator.CompactPlanner planner =
+                blobPlanner(1024, 1, 2, rowType(new DataField(1, "pic", 
DataTypes.BLOB())));
+
+        List<DataEvolutionCompactTask> tasks = planner.compactPlan(entries);
+
+        assertThat(tasks).hasSize(1);
+        
assertThat(tasks.get(0).type()).isEqualTo(DataEvolutionCompactTask.TaskType.BLOB);
+        assertThat(tasks.get(0).compactBefore())
+                .containsExactly(
+                        entries.get(1).file(), entries.get(2).file(), 
entries.get(3).file());
+    }
+
     @Test
     public void testCompactPlannerDoesNotCompactBlobFilesAcrossDataFiles() {
         List<ManifestEntry> entries = new ArrayList<>();

Reply via email to