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)