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 65b10eaf95 [flink] Support unaligned checkpoints for coordinator 
commit (#9451)
65b10eaf95 is described below

commit 65b10eaf95d66be8c89e33158dadce732b6bccf0
Author: Jingsong Lee <[email protected]>
AuthorDate: Fri Aug 28 23:18:39 2026 +0800

    [flink] Support unaligned checkpoints for coordinator commit (#9451)
---
 .../org/apache/paimon/flink/sink/FlinkSink.java    | 13 +++--
 .../paimon/flink/CoordinatorCommitITCase.java      | 49 +++++++++++++++++++
 .../apache/paimon/flink/sink/FlinkSinkTest.java    | 56 ++++++++++++++++++++++
 3 files changed, 115 insertions(+), 3 deletions(-)

diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkSink.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkSink.java
index 1445061bfc..c28d280128 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkSink.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkSink.java
@@ -214,11 +214,12 @@ public abstract class FlinkSink<T> implements 
Serializable {
         CheckpointConfig checkpointConfig = env.getCheckpointConfig();
         boolean streamingCheckpointEnabled =
                 isStreaming(written) && 
checkpointConfig.isCheckpointingEnabled();
+        boolean coordinatorCommitEnabled = coordinatorCommitEnabled();
         if (streamingCheckpointEnabled) {
-            assertStreamingConfiguration(env);
+            assertStreamingConfiguration(env, coordinatorCommitEnabled);
         }
 
-        if (coordinatorCommitEnabled()) {
+        if (coordinatorCommitEnabled) {
             return doCoordinatorCommit(written, checkpointConfig, 
streamingCheckpointEnabled);
         }
         return doOperatorCommit(written, commitUser, 
streamingCheckpointEnabled);
@@ -309,8 +310,14 @@ public abstract class FlinkSink<T> implements Serializable 
{
     }
 
     public static void assertStreamingConfiguration(StreamExecutionEnvironment 
env) {
+        assertStreamingConfiguration(env, false);
+    }
+
+    private static void assertStreamingConfiguration(
+            StreamExecutionEnvironment env, boolean 
supportsUnalignedCheckpoints) {
         checkArgument(
-                !env.getCheckpointConfig().isUnalignedCheckpointsEnabled(),
+                supportsUnalignedCheckpoints
+                        || 
!env.getCheckpointConfig().isUnalignedCheckpointsEnabled(),
                 "Paimon sink currently does not support unaligned checkpoints. 
Please set "
                         + "execution.checkpointing.unaligned.enabled to 
false.");
         checkArgument(
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/CoordinatorCommitITCase.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/CoordinatorCommitITCase.java
index eaa541e099..61ee895919 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/CoordinatorCommitITCase.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/CoordinatorCommitITCase.java
@@ -31,6 +31,7 @@ import org.apache.paimon.table.FileStoreTable;
 import org.apache.paimon.utils.CloseableIterator;
 
 import org.apache.flink.api.common.JobID;
+import org.apache.flink.api.common.JobStatus;
 import org.apache.flink.api.common.eventtime.Watermark;
 import org.apache.flink.api.connector.source.Boundedness;
 import org.apache.flink.api.connector.source.ReaderOutput;
@@ -127,6 +128,21 @@ public class CoordinatorCommitITCase {
         assertThat(readRowCount(runningJob.table)).isGreaterThan(0L);
     }
 
+    @Timeout(value = 120, unit = TimeUnit.SECONDS)
+    @Test
+    public void 
testCoordinatorCommitWritesDataEvolutionTableWithUnalignedCheckpoints()
+            throws Exception {
+        RunningJob runningJob = startStreamingInsert(true, true, true);
+        
assertThat(runningJob.table.coreOptions().dataEvolutionEnabled()).isTrue();
+        
assertThat(runningJob.table.coreOptions().rowTrackingEnabled()).isTrue();
+        waitUntilWriterInputRecords(runningJob.jobId);
+        waitUntilCoordinatorCommitMetricsRegistered(runningJob.jobId);
+        waitUntilRowsCommitted(runningJob);
+        runningJob.cancel();
+
+        assertThat(readRowCount(runningJob.table)).isGreaterThan(0L);
+    }
+
     /**
      * Idle watermark parity: the snapshot watermark observed with 
coordinator-commit enabled must
      * match the one produced by the classic {@code CommitterOperator} path 
under the same input
@@ -236,11 +252,30 @@ public class CoordinatorCommitITCase {
     }
 
     private RunningJob startStreamingInsert(boolean coordinatorCommitEnabled) 
throws Exception {
+        return startStreamingInsert(coordinatorCommitEnabled, false, false);
+    }
+
+    private RunningJob startStreamingInsert(
+            boolean coordinatorCommitEnabled,
+            boolean dataEvolutionEnabled,
+            boolean unalignedCheckpointsEnabled)
+            throws Exception {
         String tableName = coordinatorCommitEnabled ? "T_COORDINATOR_COMMIT" : 
"T_DEFAULT_COMMIT";
+        if (dataEvolutionEnabled) {
+            tableName += "_DATA_EVOLUTION";
+        }
         TableEnvironment tEnv =
                 TableEnvironment.create(
                         
EnvironmentSettings.newInstance().inStreamingMode().build());
         
tEnv.getConfig().getConfiguration().setString("execution.checkpointing.interval",
 "200 ms");
+        tEnv.getConfig()
+                .getConfiguration()
+                .setString(
+                        "execution.checkpointing.unaligned.enabled",
+                        Boolean.toString(unalignedCheckpointsEnabled));
+        if (unalignedCheckpointsEnabled) {
+            
tEnv.getConfig().getConfiguration().setString("restart-strategy.type", "none");
+        }
 
         tEnv.executeSql(
                 "CREATE CATALOG mycat WITH ( 'type' = 'paimon', 'warehouse' = 
'"
@@ -251,12 +286,17 @@ public class CoordinatorCommitITCase {
                 coordinatorCommitEnabled
                         ? ", 'sink.coordinator-commit.enabled' = 'true', 
'write-only' = 'true'"
                         : "";
+        String dataEvolutionOptions =
+                dataEvolutionEnabled
+                        ? ", 'row-tracking.enabled' = 'true', 
'data-evolution.enabled' = 'true'"
+                        : "";
         tEnv.executeSql(
                 "CREATE TABLE "
                         + tableName
                         + " (id INT, data STRING) WITH ("
                         + "'bucket' = '-1'"
                         + coordinatorCommitOption
+                        + dataEvolutionOptions
                         + ")");
         tEnv.executeSql(
                 "CREATE TEMPORARY TABLE src (id INT, data STRING) WITH ("
@@ -402,6 +442,7 @@ public class CoordinatorCommitITCase {
     private void waitUntilRowsCommitted(RunningJob runningJob) throws 
Exception {
         long deadline = System.currentTimeMillis() + WAIT_TIMEOUT_MILLIS;
         while (System.currentTimeMillis() < deadline) {
+            runningJob.checkNotTerminated();
             if (readRowCount(runningJob.table) > 0) {
                 return;
             }
@@ -427,6 +468,14 @@ public class CoordinatorCommitITCase {
         private void cancel() throws Exception {
             client.cancel().get(30, TimeUnit.SECONDS);
         }
+
+        private void checkNotTerminated() throws Exception {
+            JobStatus status = client.getJobStatus().get(30, TimeUnit.SECONDS);
+            if (status == JobStatus.FAILED) {
+                client.getJobExecutionResult().get(30, TimeUnit.SECONDS);
+            }
+            assertThat(status.isTerminalState()).describedAs("job status: %s", 
status).isFalse();
+        }
     }
 
     /**
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/FlinkSinkTest.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/FlinkSinkTest.java
index 0df52e25c4..85f94ffcb5 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/FlinkSinkTest.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/FlinkSinkTest.java
@@ -30,12 +30,17 @@ import org.apache.paimon.types.DataType;
 import org.apache.paimon.types.DataTypes;
 import org.apache.paimon.types.RowType;
 
+import org.apache.flink.api.common.RuntimeExecutionMode;
+import org.apache.flink.streaming.api.CheckpointingMode;
+import org.apache.flink.streaming.api.datastream.DataStream;
 import org.apache.flink.streaming.api.environment.CheckpointConfig;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
 import org.junit.jupiter.api.Test;
 
 import java.util.Collections;
 import java.util.function.Consumer;
 
+import static org.assertj.core.api.Assertions.assertThatCode;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
 
 /** Tests for {@link FlinkSink}. */
@@ -125,6 +130,44 @@ public class FlinkSinkTest extends CommitterTestBase {
                 .isInstanceOf(IllegalArgumentException.class);
     }
 
+    @Test
+    public void testCoordinatorCommitAllowsUnalignedCheckpoints() throws 
Exception {
+        FileStoreTable table =
+                createUnawareBucketTable(
+                        options -> {
+                            options.set(CoreOptions.ROW_TRACKING_ENABLED, 
true);
+                            options.set(CoreOptions.DATA_EVOLUTION_ENABLED, 
true);
+                            options.set(
+                                    
FlinkConnectorOptions.SINK_COORDINATOR_COMMIT_ENABLED, true);
+                        });
+
+        assertThatCode(() -> buildCommitTopology(table, 
CheckpointingMode.EXACTLY_ONCE, true))
+                .doesNotThrowAnyException();
+    }
+
+    @Test
+    public void testOperatorCommitRejectsUnalignedCheckpoints() throws 
Exception {
+        FileStoreTable table = createUnawareBucketTable(options -> {});
+
+        assertThatThrownBy(() -> buildCommitTopology(table, 
CheckpointingMode.EXACTLY_ONCE, true))
+                .isInstanceOf(IllegalArgumentException.class)
+                .hasMessageContaining("does not support unaligned 
checkpoints");
+    }
+
+    @Test
+    public void testCoordinatorCommitRejectsAtLeastOnceCheckpoints() throws 
Exception {
+        FileStoreTable table =
+                createUnawareBucketTable(
+                        options ->
+                                options.set(
+                                        
FlinkConnectorOptions.SINK_COORDINATOR_COMMIT_ENABLED,
+                                        true));
+
+        assertThatThrownBy(() -> buildCommitTopology(table, 
CheckpointingMode.AT_LEAST_ONCE, false))
+                .isInstanceOf(IllegalArgumentException.class)
+                .hasMessageContaining("only supports EXACTLY_ONCE checkpoint 
mode");
+    }
+
     private FileStoreTable createUnawareBucketTable(Consumer<Options> 
setOptions) throws Exception {
         return createFileStoreTable(
                 options -> {
@@ -155,4 +198,17 @@ public class FlinkSinkTest extends CommitterTestBase {
         config.setMaxConcurrentCheckpoints(maxConcurrentCheckpoints);
         return config;
     }
+
+    private static void buildCommitTopology(
+            FileStoreTable table, CheckpointingMode checkpointingMode, boolean 
unaligned) {
+        StreamExecutionEnvironment env = 
StreamExecutionEnvironment.getExecutionEnvironment();
+        env.setRuntimeMode(RuntimeExecutionMode.STREAMING);
+        env.enableCheckpointing(10L);
+        env.getCheckpointConfig().setCheckpointingMode(checkpointingMode);
+        env.getCheckpointConfig().enableUnalignedCheckpoints(unaligned);
+        DataStream<Committable> written =
+                env.fromCollection(Collections.emptyList(), new 
CommittableTypeInfo());
+
+        new RowAppendTableSink(table, null, null).doCommit(written, 
"test-commit-user");
+    }
 }

Reply via email to