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 061e021ea8 [Flink] Add configurable initial scan mode for streaming
compaction (#9548)
061e021ea8 is described below
commit 061e021ea83b5b752635a071e693a478a3bc71f5
Author: sanshi <[email protected]>
AuthorDate: Mon Sep 7 22:45:36 2026 +0800
[Flink] Add configurable initial scan mode for streaming compaction (#9548)
---
docs/generated/core_configuration.html | 6 +
.../main/java/org/apache/paimon/CoreOptions.java | 38 +++++++
.../paimon/table/source/AbstractDataTableScan.java | 5 +-
.../ContinuousCompactorStartingScanner.java | 19 ++++
.../ContinuousCompactorStartingScannerTest.java | 123 +++++++++++++++++++++
5 files changed, 190 insertions(+), 1 deletion(-)
diff --git a/docs/generated/core_configuration.html
b/docs/generated/core_configuration.html
index 270a363ec0..6a4955280d 100644
--- a/docs/generated/core_configuration.html
+++ b/docs/generated/core_configuration.html
@@ -392,6 +392,12 @@ under the License.
<td>MemorySize</td>
<td>When incremental size is bigger than this threshold, force a
full compaction.</td>
</tr>
+ <tr>
+ <td><h5>continuous-compaction.initial-scan-mode</h5></td>
+ <td style="word-wrap: break-word;">earliest</td>
+ <td><p>Enum</p></td>
+ <td>Initial snapshot mode for dedicated streaming compaction. When
set to 'earliest' (the default), compaction starts from the earliest available
snapshot if no COMPACT snapshot exists; when a COMPACT snapshot exists,
compaction always resumes from the snapshot after it. When set to 'latest', the
latest snapshot is read in ALL mode as the initial baseline and subsequent
scans start from the next snapshot. The 'latest' mode skips historical snapshot
changes and should only be [...]
+ </tr>
<tr>
<td><h5>compaction.max-size-amplification-percent</h5></td>
<td style="word-wrap: break-word;">200</td>
diff --git a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
index c46ef4fb79..298126f3c6 100644
--- a/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
+++ b/paimon-api/src/main/java/org/apache/paimon/CoreOptions.java
@@ -1655,6 +1655,20 @@ public class CoreOptions implements Serializable {
"Only used to force TableScan to construct
suitable 'StartingUpScanner' and 'FollowUpScanner' "
+ "dedicated internal streaming scan.");
+ public static final ConfigOption<CompactionInitialScanMode>
+ CONTINUOUS_COMPACTION_INITIAL_SCAN_MODE =
+ key("continuous-compaction.initial-scan-mode")
+ .enumType(CompactionInitialScanMode.class)
+ .defaultValue(CompactionInitialScanMode.EARLIEST)
+ .withDescription(
+ "Initial snapshot mode for dedicated
streaming compaction. "
+ + "When set to 'earliest' (the
default), compaction starts from the earliest available snapshot "
+ + "if no COMPACT snapshot exists;
when a COMPACT snapshot exists, compaction always resumes from the snapshot
after it. "
+ + "When set to 'latest', the
latest snapshot is read in ALL mode "
+ + "as the initial baseline and
subsequent scans start from the next snapshot. "
+ + "The 'latest' mode skips
historical snapshot changes and should only be used when historical "
+ + "changelog replay is not
required.");
+
@ExcludeFromDocumentation("Internal use only")
public static final ConfigOption<BatchScanMode> BATCH_SCAN_MODE =
key("batch-scan-mode")
@@ -5194,6 +5208,30 @@ public class CoreOptions implements Serializable {
}
}
+ /** Initial snapshot mode for dedicated streaming compaction. */
+ public enum CompactionInitialScanMode implements DescribedEnum {
+ EARLIEST("earliest", "Read snapshots from the earliest available
snapshot."),
+ LATEST("latest", "Read the latest snapshot as the initial full
baseline.");
+
+ private final String value;
+ private final String description;
+
+ CompactionInitialScanMode(String value, String description) {
+ this.value = value;
+ this.description = description;
+ }
+
+ @Override
+ public String toString() {
+ return value;
+ }
+
+ @Override
+ public InlineElement getDescription() {
+ return text(description);
+ }
+ }
+
/** Inner stream scan mode for some internal requirements. */
public enum StreamScanMode implements DescribedEnum {
NONE("none", "No requirement."),
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataTableScan.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataTableScan.java
index 24bf32e50c..0cdef986fe 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataTableScan.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/AbstractDataTableScan.java
@@ -434,7 +434,10 @@ abstract class AbstractDataTableScan implements
DataTableScan {
case COMPACT_BUCKET_TABLE:
checkArgument(
isStreaming, "Set 'streaming-compact' in batch mode.
This is unexpected.");
- return new ContinuousCompactorStartingScanner(snapshotManager);
+ return new ContinuousCompactorStartingScanner(
+ snapshotManager,
+ options.toConfiguration()
+
.get(CoreOptions.CONTINUOUS_COMPACTION_INITIAL_SCAN_MODE));
case FILE_MONITOR:
return new FullStartingScanner(snapshotManager);
}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/ContinuousCompactorStartingScanner.java
b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/ContinuousCompactorStartingScanner.java
index 5f24366ca1..0dd095b936 100644
---
a/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/ContinuousCompactorStartingScanner.java
+++
b/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/ContinuousCompactorStartingScanner.java
@@ -18,7 +18,9 @@
package org.apache.paimon.table.source.snapshot;
+import org.apache.paimon.CoreOptions;
import org.apache.paimon.Snapshot;
+import org.apache.paimon.table.source.ScanMode;
import org.apache.paimon.utils.SnapshotManager;
import org.slf4j.Logger;
@@ -27,11 +29,21 @@ import org.slf4j.LoggerFactory;
/** {@link StartingScanner} used internally for stand-alone streaming compact
job sources. */
public class ContinuousCompactorStartingScanner extends
AbstractStartingScanner {
+ private final boolean latestInitialSnapshot;
+
private static final Logger LOG =
LoggerFactory.getLogger(ContinuousCompactorStartingScanner.class);
public ContinuousCompactorStartingScanner(SnapshotManager snapshotManager)
{
+ this(snapshotManager, CoreOptions.CompactionInitialScanMode.EARLIEST);
+ }
+
+ public ContinuousCompactorStartingScanner(
+ SnapshotManager snapshotManager,
+ CoreOptions.CompactionInitialScanMode initialScanMode) {
super(snapshotManager);
+ this.latestInitialSnapshot =
+ initialScanMode ==
CoreOptions.CompactionInitialScanMode.LATEST;
this.startingSnapshotId = snapshotManager.earliestSnapshotId();
}
@@ -52,6 +64,13 @@ public class ContinuousCompactorStartingScanner extends
AbstractStartingScanner
}
}
+ if (latestInitialSnapshot) {
+ LOG.debug(
+ "No compact snapshot found, reading the latest snapshot {}
as the initial compaction baseline.",
+ latestSnapshotId);
+ return StartingScanner.fromPlan(
+
snapshotReader.withMode(ScanMode.ALL).withSnapshot(latestSnapshotId).read());
+ }
LOG.debug(
"No compact snapshot found, reading from the earliest snapshot
{}.",
earliestSnapshotId);
diff --git
a/paimon-core/src/test/java/org/apache/paimon/table/source/snapshot/ContinuousCompactorStartingScannerTest.java
b/paimon-core/src/test/java/org/apache/paimon/table/source/snapshot/ContinuousCompactorStartingScannerTest.java
index a55a36f014..ecd798ec97 100644
---
a/paimon-core/src/test/java/org/apache/paimon/table/source/snapshot/ContinuousCompactorStartingScannerTest.java
+++
b/paimon-core/src/test/java/org/apache/paimon/table/source/snapshot/ContinuousCompactorStartingScannerTest.java
@@ -18,13 +18,24 @@
package org.apache.paimon.table.source.snapshot;
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.Snapshot;
+import org.apache.paimon.options.Options;
import org.apache.paimon.table.sink.StreamTableCommit;
import org.apache.paimon.table.sink.StreamTableWrite;
+import org.apache.paimon.table.source.DataSplit;
+import org.apache.paimon.table.source.Split;
+import org.apache.paimon.table.source.StreamTableScan;
+import org.apache.paimon.table.source.TableScan;
import org.apache.paimon.types.RowKind;
import org.apache.paimon.utils.SnapshotManager;
import org.junit.jupiter.api.Test;
+import java.util.Arrays;
+import java.util.HashSet;
+import java.util.Set;
+
import static org.assertj.core.api.Assertions.assertThat;
/** Tests for {@link ContinuousCompactorStartingScanner}. */
@@ -75,4 +86,116 @@ public class ContinuousCompactorStartingScannerTest extends
ScannerTestBase {
new ContinuousCompactorStartingScanner(snapshotManager);
assertThat(scanner.scan(snapshotReader)).isInstanceOf(StartingScanner.NoSnapshot.class);
}
+
+ @Test
+ public void testLatestBaselineIsFollowedByDeltaScan() throws Exception {
+ Options options = new Options();
+ options.set(CoreOptions.WRITE_ONLY, true);
+ options.set(CoreOptions.STREAM_SCAN_MODE,
CoreOptions.StreamScanMode.COMPACT_BUCKET_TABLE);
+ options.set(
+ CoreOptions.CONTINUOUS_COMPACTION_INITIAL_SCAN_MODE,
+ CoreOptions.CompactionInitialScanMode.LATEST);
+ createAppendOnlyTable(options);
+ StreamTableWrite write = table.newWrite(commitUser);
+ StreamTableCommit commit = table.newCommit(commitUser);
+
+ write.write(rowData(1, 10, 100L));
+ commit.commit(0, write.prepareCommit(true, 0));
+ write.write(rowData(1, 11, 101L));
+ commit.commit(1, write.prepareCommit(true, 1));
+
+ StreamTableScan scan = table.newStreamScan();
+ TableScan.Plan baseline = scan.plan();
+ assertThat(baseline.splits()).allMatch(split -> ((DataSplit)
split).snapshotId() == 2L);
+ assertThat(getResult(table.newRead(), baseline.splits()))
+ .hasSameElementsAs(Arrays.asList("+I 1|10|100", "+I
1|11|101"));
+ assertThat(scan.checkpoint()).isEqualTo(3L);
+
+ write.write(rowData(1, 12, 102L));
+ commit.commit(2, write.prepareCommit(true, 2));
+
+ TableScan.Plan delta = scan.plan();
+ assertThat(delta.splits()).allMatch(split -> ((DataSplit)
split).snapshotId() == 3L);
+ assertThat(delta.splits()).isNotEmpty();
+ assertThat(getResult(table.newRead(),
delta.splits())).containsExactly("+I 1|12|102");
+ assertThat(scan.checkpoint()).isEqualTo(4L);
+
+ write.close();
+ commit.close();
+ }
+
+ @Test
+ public void
testNoCompactSnapshotLatestBaselineContainsAllPartitionsAndBuckets()
+ throws Exception {
+ Options options = new Options();
+ options.set(CoreOptions.WRITE_ONLY, true);
+ options.set(CoreOptions.BUCKET, 2);
+ options.set(CoreOptions.BUCKET_KEY, "a");
+ createAppendOnlyTable(options);
+ SnapshotManager snapshotManager = table.snapshotManager();
+ StreamTableWrite write = table.newWrite(commitUser);
+ StreamTableCommit commit = table.newCommit(commitUser);
+
+ write.write(rowData(1, 10, 100L));
+ write.write(rowData(1, 11, 101L));
+ write.write(rowData(2, 10, 200L));
+ write.write(rowData(2, 11, 201L));
+ commit.commit(0, write.prepareCommit(true, 0));
+
+ StartingScanner.NextSnapshot earliestResult =
+ (StartingScanner.NextSnapshot)
+ new ContinuousCompactorStartingScanner(snapshotManager)
+ .scan(snapshotReader);
+ assertThat(earliestResult.nextSnapshotId()).isEqualTo(1L);
+
+ ContinuousCompactorStartingScanner scanner =
+ new ContinuousCompactorStartingScanner(
+ snapshotManager,
CoreOptions.CompactionInitialScanMode.LATEST);
+ StartingScanner.ScannedResult result =
+ (StartingScanner.ScannedResult) scanner.scan(snapshotReader);
+
+ Set<String> partitionBuckets = new HashSet<>();
+ for (Split split : result.splits()) {
+ DataSplit dataSplit = (DataSplit) split;
+ partitionBuckets.add(dataSplit.partition().getInt(0) + ":" +
dataSplit.bucket());
+ }
+ assertThat(partitionBuckets).hasSize(4);
+ assertThat(result.splits()).allMatch(split -> !((DataSplit)
split).dataFiles().isEmpty());
+
+ write.close();
+ commit.close();
+ }
+
+ @Test
+ public void testNoCompactSnapshotReadsLatestAsInitialBaseline() throws
Exception {
+ Options options = new Options();
+ options.set(CoreOptions.WRITE_ONLY, true);
+ createAppendOnlyTable(options);
+ SnapshotManager snapshotManager = table.snapshotManager();
+ StreamTableWrite write = table.newWrite(commitUser);
+ StreamTableCommit commit = table.newCommit(commitUser);
+
+ for (int i = 0; i < 5; i++) {
+ write.write(rowData(1, i, (long) i));
+ commit.commit(i, write.prepareCommit(true, i));
+ }
+
+ ContinuousCompactorStartingScanner scanner =
+ new ContinuousCompactorStartingScanner(
+ snapshotManager,
CoreOptions.CompactionInitialScanMode.LATEST);
+ StartingScanner.ScannedResult result =
+ (StartingScanner.ScannedResult) scanner.scan(snapshotReader);
+
+ assertThat(snapshotManager.earliestSnapshotId()).isEqualTo(1L);
+ assertThat(snapshotManager.latestSnapshotId()).isEqualTo(5L);
+
assertThat(snapshotManager.snapshot(5L).commitKind()).isEqualTo(Snapshot.CommitKind.APPEND);
+ assertThat(result.currentSnapshotId()).isEqualTo(5);
+ assertThat(result.plan().snapshotId()).isEqualTo(5);
+ assertThat(result.splits()).isNotEmpty();
+ assertThat(result.splits()).allMatch(split -> ((DataSplit)
split).snapshotId() == 5);
+ assertThat(result.splits()).allMatch(split -> !((DataSplit)
split).dataFiles().isEmpty());
+
+ write.close();
+ commit.close();
+ }
}