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 b193d0b24e [flink] Recover coordinator commit without intended 
failover (#8970)
b193d0b24e is described below

commit b193d0b24ec5dbfc9fcdbf077c4764c80ab2a7dc
Author: Biao Liu <[email protected]>
AuthorDate: Mon Aug 3 19:04:03 2026 +0800

    [flink] Recover coordinator commit without intended failover (#8970)
---
 .../paimon/flink/sink/RowAppendTableSink.java      |  18 +--
 .../CommittingWriteOperatorCoordinator.java        |  50 ++-------
 ...torCommittingRowDataStoreWriteOperatorTest.java |   3 +-
 .../CommittingWriteOperatorCoordinatorTest.java    | 121 +++++++--------------
 4 files changed, 56 insertions(+), 136 deletions(-)

diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/RowAppendTableSink.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/RowAppendTableSink.java
index b926dae4c2..443c9b6769 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/RowAppendTableSink.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/RowAppendTableSink.java
@@ -56,7 +56,6 @@ public class RowAppendTableSink extends 
AppendTableSink<InternalRow> {
                     // checkpointing on by default for the JM-side committer; 
bounded sources will
                     // be handled by end-input support in a follow-up PR
                     true,
-                    true,
                     createCommitterFactory());
         }
         return createNoStateRowWriteOperatorFactory(table, writeProvider, 
commitUser);
@@ -83,15 +82,9 @@ public class RowAppendTableSink extends 
AppendTableSink<InternalRow> {
                     StoreSinkWrite.Provider writeProvider,
                     String commitUser,
                     boolean streamingCheckpointEnabled,
-                    boolean failoverAfterRecovery,
                     Committer.Factory<Committable, ManifestCommittable> 
committerFactory) {
         return new CoordinatorCommittingFactory(
-                table,
-                writeProvider,
-                commitUser,
-                streamingCheckpointEnabled,
-                failoverAfterRecovery,
-                committerFactory);
+                table, writeProvider, commitUser, streamingCheckpointEnabled, 
committerFactory);
     }
 
     private static class CoordinatorCommittingFactory extends 
RowDataStoreWriteOperator.Factory
@@ -100,7 +93,6 @@ public class RowAppendTableSink extends 
AppendTableSink<InternalRow> {
         private static final long serialVersionUID = 1L;
 
         private final boolean streamingCheckpointEnabled;
-        private final boolean failoverAfterRecovery;
         private final Committer.Factory<Committable, ManifestCommittable> 
committerFactory;
 
         CoordinatorCommittingFactory(
@@ -108,11 +100,9 @@ public class RowAppendTableSink extends 
AppendTableSink<InternalRow> {
                 StoreSinkWrite.Provider storeSinkWriteProvider,
                 String initialCommitUser,
                 boolean streamingCheckpointEnabled,
-                boolean failoverAfterRecovery,
                 Committer.Factory<Committable, ManifestCommittable> 
committerFactory) {
             super(table, storeSinkWriteProvider, initialCommitUser);
             this.streamingCheckpointEnabled = streamingCheckpointEnabled;
-            this.failoverAfterRecovery = failoverAfterRecovery;
             this.committerFactory = committerFactory;
         }
 
@@ -120,11 +110,7 @@ public class RowAppendTableSink extends 
AppendTableSink<InternalRow> {
         public OperatorCoordinator.Provider getCoordinatorProvider(
                 String operatorName, OperatorID operatorID) {
             return new CommittingWriteOperatorCoordinator.Provider(
-                    operatorID,
-                    committerFactory,
-                    streamingCheckpointEnabled,
-                    initialCommitUser,
-                    failoverAfterRecovery);
+                    operatorID, committerFactory, streamingCheckpointEnabled, 
initialCommitUser);
         }
 
         @Override
diff --git 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/coordinator/CommittingWriteOperatorCoordinator.java
 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/coordinator/CommittingWriteOperatorCoordinator.java
index 20f3b65b88..3ef4b26d42 100644
--- 
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/coordinator/CommittingWriteOperatorCoordinator.java
+++ 
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/coordinator/CommittingWriteOperatorCoordinator.java
@@ -76,7 +76,6 @@ public class CommittingWriteOperatorCoordinator implements 
OperatorCoordinator {
     private final OperatorCoordinator.Context context;
     private final Committer.Factory<Committable, ManifestCommittable> 
committerFactory;
     private final boolean streamingCheckpointEnabled;
-    private final boolean failoverAfterRecovery;
     private final int parallelism;
 
     private final WriterCommittables[] subtaskCommittables;
@@ -101,13 +100,11 @@ public class CommittingWriteOperatorCoordinator 
implements OperatorCoordinator {
             OperatorCoordinator.Context context,
             Committer.Factory<Committable, ManifestCommittable> 
committerFactory,
             boolean streamingCheckpointEnabled,
-            String initialCommitUser,
-            boolean failoverAfterRecovery) {
+            String initialCommitUser) {
         this.context = context;
         this.committerFactory = committerFactory;
         this.streamingCheckpointEnabled = streamingCheckpointEnabled;
         this.commitUser = initialCommitUser;
-        this.failoverAfterRecovery = failoverAfterRecovery;
         this.parallelism = context.currentParallelism();
         this.subtaskCommittables = new WriterCommittables[parallelism];
         this.committablesSerializer =
@@ -354,31 +351,15 @@ public class CommittingWriteOperatorCoordinator 
implements OperatorCoordinator {
 
     // replaces CommittableStateManager because committables are not stored in 
the committer
     private void recover(long checkpointId) throws Exception {
-        if (failoverAfterRecovery) {
-            // recommit the restored committables and trigger a failover to 
reinitialize all writers
-            Map<Long, Long> watermarkPerCheckpoint =
-                    alignWatermarkPerCheckpoint(
-                            checkpointId, subtaskCommittables, 
watermarkAligner);
-            commitUpToCheckpoint(
-                    checkpointId,
-                    pollManifestCommittablesForCheckpoint(
-                            checkpointId, subtaskCommittables, 
watermarkPerCheckpoint, committer),
-                    watermarkPerCheckpoint,
-                    committables -> {
-                        int numCommitted = 
committer.filterAndCommit(committables, true, true);
-                        if (numCommitted > 0) {
-                            throw new RuntimeException(
-                                    "This exception is intentionally thrown 
after committing the "
-                                            + "restored checkpoints. By 
restarting the job we hope "
-                                            + "that writers can start writing 
based on these new commits.");
-                        }
-                    });
-        } else {
-            // just abandon the restoring committables
-            for (WriterCommittables subtaskCommit : subtaskCommittables) {
-                subtaskCommit.clearCommittablesBeforeCheckpoint(checkpointId, 
true);
-            }
-        }
+        // Mirror RestoreCommittableStateManager: re-commit restored 
committables and keep running.
+        Map<Long, Long> watermarkPerCheckpoint =
+                alignWatermarkPerCheckpoint(checkpointId, subtaskCommittables, 
watermarkAligner);
+        commitUpToCheckpoint(
+                checkpointId,
+                pollManifestCommittablesForCheckpoint(
+                        checkpointId, subtaskCommittables, 
watermarkPerCheckpoint, committer),
+                watermarkPerCheckpoint,
+                committables -> committer.filterAndCommit(committables, true, 
true));
     }
 
     @VisibleForTesting
@@ -619,29 +600,22 @@ public class CommittingWriteOperatorCoordinator 
implements OperatorCoordinator {
         private final Committer.Factory<Committable, ManifestCommittable> 
committerFactory;
         private final boolean streamingCheckpointEnabled;
         private final String initialCommitUser;
-        private final boolean failoverAfterRecovery;
 
         public Provider(
                 OperatorID operatorId,
                 Committer.Factory<Committable, ManifestCommittable> 
committerFactory,
                 boolean streamingCheckpointEnabled,
-                String initialCommitUser,
-                boolean failoverAfterRecovery) {
+                String initialCommitUser) {
             super(operatorId);
             this.committerFactory = committerFactory;
             this.streamingCheckpointEnabled = streamingCheckpointEnabled;
             this.initialCommitUser = initialCommitUser;
-            this.failoverAfterRecovery = failoverAfterRecovery;
         }
 
         @Override
         public OperatorCoordinator getCoordinator(OperatorCoordinator.Context 
context) {
             return new CommittingWriteOperatorCoordinator(
-                    context,
-                    committerFactory,
-                    streamingCheckpointEnabled,
-                    initialCommitUser,
-                    failoverAfterRecovery);
+                    context, committerFactory, streamingCheckpointEnabled, 
initialCommitUser);
         }
     }
 }
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/CoordinatorCommittingRowDataStoreWriteOperatorTest.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/CoordinatorCommittingRowDataStoreWriteOperatorTest.java
index bae18ea2c2..89911cc271 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/CoordinatorCommittingRowDataStoreWriteOperatorTest.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/CoordinatorCommittingRowDataStoreWriteOperatorTest.java
@@ -91,8 +91,7 @@ public class 
CoordinatorCommittingRowDataStoreWriteOperatorTest extends Committe
                                 new StoreCommitter(
                                         table, 
table.newCommit(context.commitUser()), context),
                         true,
-                        commitUser,
-                        false);
+                        commitUser);
         coordinator.start();
         coordinator.waitProcessAllActions();
 
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/coordinator/CommittingWriteOperatorCoordinatorTest.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/coordinator/CommittingWriteOperatorCoordinatorTest.java
index a43bba18d3..74f54da3aa 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/coordinator/CommittingWriteOperatorCoordinatorTest.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/coordinator/CommittingWriteOperatorCoordinatorTest.java
@@ -103,7 +103,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
     public void testCommitSingleSubtask() throws Exception {
         FileStoreTable table = createUnawareBucketTable();
         TestingContext context = new TestingContext(new OperatorID(), 1);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         coordinator.start();
         coordinator.waitProcessAllActions();
         assertThat(coordinator.getCurrentState())
@@ -124,7 +124,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
     public void testCommitFanInFromMultipleSubtasks() throws Exception {
         FileStoreTable table = createUnawareBucketTable();
         TestingContext context = new TestingContext(new OperatorID(), 2);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         coordinator.start();
         coordinator.waitProcessAllActions();
 
@@ -142,7 +142,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
     public void testWatermarkCommit() throws Exception {
         FileStoreTable table = createUnawareBucketTable();
         TestingContext context = new TestingContext(new OperatorID(), 1);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         coordinator.start();
         coordinator.waitProcessAllActions();
 
@@ -170,7 +170,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
         // barrier.
         FileStoreTable table = createUnawareBucketTable();
         TestingContext context = new TestingContext(new OperatorID(), 1);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         coordinator.start();
         coordinator.waitProcessAllActions();
 
@@ -195,7 +195,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
         // cannot advance a snapshot beyond what all writers had actually 
observed.
         FileStoreTable table = createUnawareBucketTable();
         TestingContext context = new TestingContext(new OperatorID(), 2);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         coordinator.start();
         coordinator.waitProcessAllActions();
 
@@ -214,7 +214,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
         TestingContext context = new TestingContext(new OperatorID(), 2);
 
         // first incarnation commits checkpoint 1 and captures the coordinator 
state
-        CommittingWriteOperatorCoordinator first = createCoordinator(table, 
context, false);
+        CommittingWriteOperatorCoordinator first = createCoordinator(table, 
context);
         first.start();
         first.waitProcessAllActions();
         first.handleEventFromOperator(0, 0, event(committable(table, 1, 1)));
@@ -228,7 +228,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
         assertResults(table, "1, 1", "2, 2");
 
         // second incarnation restores and stays RESTORING until both subtasks 
re-emit
-        CommittingWriteOperatorCoordinator second = createCoordinator(table, 
context, false);
+        CommittingWriteOperatorCoordinator second = createCoordinator(table, 
context);
         second.resetToCheckpoint(1, state);
         assertThat(second.getCurrentState())
                 .isEqualTo(CommittingWriteOperatorCoordinator.State.RESTORING);
@@ -248,61 +248,19 @@ public class CommittingWriteOperatorCoordinatorTest 
extends CommitterOperatorTes
         assertThat(second.getCurrentState())
                 .isEqualTo(CommittingWriteOperatorCoordinator.State.RUNNING);
 
-        // abandon path: restoring committables are dropped, not recommitted
+        // cp1 was already committed by the first incarnation; on restore the 
replayed cp1
+        // committables are filtered out by filterAndCommit (idempotent), so 
nothing is
+        // re-committed.
         assertResults(table, "1, 1", "2, 2");
         second.close();
     }
 
-    @Timeout(value = 30, unit = TimeUnit.SECONDS)
-    @Test
-    public void testSnapshotLostWhenFailed() throws Exception {
-        FileStoreTable table = createUnawareBucketTable();
-        TestingContext context = new TestingContext(new OperatorID(), 1);
-
-        // first incarnation: cp1 fully committed, cp2 snapshotted but never 
notified
-        CommittingWriteOperatorCoordinator first = createCoordinator(table, 
context, false);
-        first.start();
-        first.waitProcessAllActions();
-        first.handleEventFromOperator(0, 0, event(committable(table, 1, 1)));
-        first.notifyCheckpointComplete(1L);
-        first.waitProcessAllActions();
-        assertResults(table, "1, 1");
-
-        first.handleEventFromOperator(0, 0, event(committable(table, 2, 2)));
-        CompletableFuture<byte[]> cp2State = new CompletableFuture<>();
-        first.checkpointCoordinator(2L, cp2State);
-        first.waitProcessAllActions();
-        byte[] state = cp2State.get();
-        first.close();
-        // cp2 was never notified — only cp1 is in the table
-        assertResults(table, "1, 1");
-
-        // second incarnation: restore from cp2 state, replay the cp2 
restoring event. abandon
-        // mode drops it; the snapshot from cp1 stays untouched.
-        CommittingWriteOperatorCoordinator second = createCoordinator(table, 
context, false);
-        second.resetToCheckpoint(2L, state);
-        second.start();
-        second.waitProcessAllActions();
-        second.handleEventFromOperator(0, 0, restoreEvent(2L, 
committable(table, 2, 2)));
-        second.waitProcessAllActions();
-        assertThat(second.getCurrentState())
-                .isEqualTo(CommittingWriteOperatorCoordinator.State.RUNNING);
-        assertResults(table, "1, 1");
-
-        // a fresh checkpoint after recovery commits normally
-        second.handleEventFromOperator(0, 0, event(committable(table, 3, 3)));
-        second.notifyCheckpointComplete(3L);
-        second.waitProcessAllActions();
-        assertResults(table, "1, 1", "3, 3");
-        second.close();
-    }
-
     @Timeout(value = 30, unit = TimeUnit.SECONDS)
     @Test
     public void testRejectCheckpointWhileRestoring() throws Exception {
         FileStoreTable table = createUnawareBucketTable();
         TestingContext context = new TestingContext(new OperatorID(), 1);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         coordinator.resetToCheckpoint(2, emptyState());
         assertThat(coordinator.getCurrentState())
                 .isEqualTo(CommittingWriteOperatorCoordinator.State.RESTORING);
@@ -332,7 +290,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
     public void testCommittableEventInRestoringFailsJob() throws Exception {
         FileStoreTable table = createUnawareBucketTable();
         TestingContext context = new TestingContext(new OperatorID(), 1);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         coordinator.resetToCheckpoint(2, emptyState());
         coordinator.start();
         coordinator.waitProcessAllActions();
@@ -350,12 +308,12 @@ public class CommittingWriteOperatorCoordinatorTest 
extends CommitterOperatorTes
 
     @Timeout(value = 30, unit = TimeUnit.SECONDS)
     @Test
-    public void testFailIntentionallyAfterRestoring() throws Exception {
+    public void testRecommitOnRestoreWithoutFailover() throws Exception {
         FileStoreTable table = createUnawareBucketTable();
         TestingContext context = new TestingContext(new OperatorID(), 1);
 
         // capture coordinator state without committing checkpoint 1
-        CommittingWriteOperatorCoordinator first = createCoordinator(table, 
context, true);
+        CommittingWriteOperatorCoordinator first = createCoordinator(table, 
context);
         first.start();
         first.handleEventFromOperator(0, 0, event(committable(table, 1, 1)));
         CompletableFuture<byte[]> checkpoint = new CompletableFuture<>();
@@ -366,19 +324,26 @@ public class CommittingWriteOperatorCoordinatorTest 
extends CommitterOperatorTes
         // checkpoint 1 was never committed
         assertThat(table.latestSnapshot()).isNotPresent();
 
-        // restore with failoverAfterRecovery: the restored committables are 
recommitted and an
-        // intentional failure is raised to reinitialize all writers
-        CommittingWriteOperatorCoordinator second = createCoordinator(table, 
context, true);
+        // On restore the not-yet-committed checkpoint 1 is re-committed and 
the coordinator keeps
+        // running — no intentional failover. Unaware-append writers are 
stateless w.r.t. committed
+        // snapshots, so they need not restart after the restore-time commit.
+        CommittingWriteOperatorCoordinator second = createCoordinator(table, 
context);
         second.resetToCheckpoint(1, state);
         second.start();
         second.waitProcessAllActions();
         second.handleEventFromOperator(0, 0, restoreEvent(1L, 
committable(table, 1, 1)));
         second.waitProcessAllActions();
 
-        assertThat(failureCause).isInstanceOf(RuntimeException.class);
-        assertThat(failureCause).hasMessageContaining("intentionally thrown");
+        assertThat(failureCause).isNull();
+        assertThat(second.getCurrentState())
+                .isEqualTo(CommittingWriteOperatorCoordinator.State.RUNNING);
         assertResults(table, "1, 1");
-        failureCause = null;
+
+        // a fresh checkpoint after recovery commits normally, confirming the 
coordinator is live
+        second.handleEventFromOperator(0, 0, event(committable(table, 2, 2)));
+        second.notifyCheckpointComplete(2L);
+        second.waitProcessAllActions();
+        assertResults(table, "1, 1", "2, 2");
         second.close();
     }
 
@@ -389,7 +354,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
         // the last one. the coordinator must drain all pending checkpoints in 
a single commit.
         FileStoreTable table = createUnawareBucketTable();
         TestingContext context = new TestingContext(new OperatorID(), 1);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         coordinator.start();
         coordinator.waitProcessAllActions();
 
@@ -430,8 +395,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
                                                 commitContext),
                                         expected),
                         true,
-                        commitUser,
-                        false);
+                        commitUser);
         coordinator.start();
         coordinator.waitProcessAllActions();
 
@@ -451,7 +415,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
     public void testEmptyCommit() throws Exception {
         FileStoreTable table = createUnawareBucketTable();
         TestingContext context = new TestingContext(new OperatorID(), 1);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         coordinator.start();
         coordinator.waitProcessAllActions();
 
@@ -474,7 +438,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
                             
options.set(CoreOptions.COMMIT_FORCE_CREATE_SNAPSHOT, true);
                         });
         TestingContext context = new TestingContext(new OperatorID(), 1);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         coordinator.start();
         coordinator.waitProcessAllActions();
 
@@ -504,7 +468,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
                             
options.set(CoreOptions.COMMIT_FORCE_CREATE_SNAPSHOT, true);
                         });
         TestingContext context = new TestingContext(new OperatorID(), 1);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         coordinator.start();
         coordinator.waitProcessAllActions();
 
@@ -532,7 +496,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
     public void testAlignmentHonorsEmptyMinValueMarker() throws Exception {
         FileStoreTable table = createUnawareBucketTable();
         TestingContext context = new TestingContext(new OperatorID(), 2);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         coordinator.start();
         coordinator.waitProcessAllActions();
 
@@ -561,7 +525,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
     public void testAlignmentSkipsIdleSubtaskWhenSomeActive() throws Exception 
{
         FileStoreTable table = createUnawareBucketTable();
         TestingContext context = new TestingContext(new OperatorID(), 2);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         coordinator.start();
         coordinator.waitProcessAllActions();
 
@@ -806,7 +770,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
                         },
                         Collections.singletonList("a"));
         TestingContext context = new TestingContext(new OperatorID(), 2);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         // 1. start with non-restoring
         assertThat(coordinator.getCurrentState())
                 .isEqualTo(CommittingWriteOperatorCoordinator.State.CREATED);
@@ -901,7 +865,7 @@ public class CommittingWriteOperatorCoordinatorTest extends 
CommitterOperatorTes
                         },
                         Collections.singletonList("a"));
         TestingContext context = new TestingContext(new OperatorID(), 2);
-        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context, false);
+        CommittingWriteOperatorCoordinator coordinator = 
createCoordinator(table, context);
         // 1. start with non-restoring
         assertThat(coordinator.getCurrentState())
                 .isEqualTo(CommittingWriteOperatorCoordinator.State.CREATED);
@@ -1002,7 +966,7 @@ public class CommittingWriteOperatorCoordinatorTest 
extends CommitterOperatorTes
         TestingContext context = new TestingContext(new OperatorID(), 1);
 
         // 1. capture state from a coordinator without mark-done enabled
-        CommittingWriteOperatorCoordinator first = createCoordinator(table, 
context, false);
+        CommittingWriteOperatorCoordinator first = createCoordinator(table, 
context);
         first.start();
         first.waitProcessAllActions();
         first.handleEventFromOperator(0, 0, event(committable(table, 1, 1)));
@@ -1015,8 +979,7 @@ public class CommittingWriteOperatorCoordinatorTest 
extends CommitterOperatorTes
 
         // 2. restore with mark-done enabled — should initialize cleanly
         FileStoreTable markDoneTable = table.copy(markDoneOption);
-        CommittingWriteOperatorCoordinator second =
-                createCoordinator(markDoneTable, context, false);
+        CommittingWriteOperatorCoordinator second = 
createCoordinator(markDoneTable, context);
         second.resetToCheckpoint(1L, state);
         second.start();
         second.waitProcessAllActions();
@@ -1113,7 +1076,7 @@ public class CommittingWriteOperatorCoordinatorTest 
extends CommitterOperatorTes
     }
 
     private CommittingWriteOperatorCoordinator createCoordinator(
-            FileStoreTable table, TestingContext context, boolean 
failoverAfterRecovery) {
+            FileStoreTable table, TestingContext context) {
         return new CommittingWriteOperatorCoordinator(
                 context,
                 commitContext ->
@@ -1124,8 +1087,7 @@ public class CommittingWriteOperatorCoordinatorTest 
extends CommitterOperatorTes
                                         .newCommit(),
                                 commitContext),
                 true,
-                commitUser,
-                failoverAfterRecovery);
+                commitUser);
     }
 
     private CommittingWriteOperatorCoordinator 
createCoordinatorCapturingContext(
@@ -1144,8 +1106,7 @@ public class CommittingWriteOperatorCoordinatorTest 
extends CommitterOperatorTes
                             commitContext);
                 },
                 true,
-                commitUser,
-                false);
+                commitUser);
     }
 
     private Committable committable(FileStoreTable table, long checkpointId, 
int value)

Reply via email to