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<>();