This is an automated email from the ASF dual-hosted git repository. voonhous pushed a commit to branch release-1.2.1 in repository https://gitbox.apache.org/repos/asf/hudi.git
commit d6ee4e285763d369b93cf8d6e4863ac1fbfc3609 Author: Shuo Cheng <[email protected]> AuthorDate: Tue Jun 2 08:47:25 2026 +0800 fix(flink): Trigger a failover after pending instants recommitted for both global and partitioned RLI (#18793) (cherry picked from commit ed9ea0ead908e3765b6c938da121e96ace13ba38) --- .../hudi/sink/utils/BulkInsertFunctionWrapper.java | 4 +--- .../hudi/sink/utils/InsertFunctionWrapper.java | 20 ++++++++++++++++---- .../hudi/sink/utils/StreamWriteFunctionWrapper.java | 6 ------ 3 files changed, 17 insertions(+), 13 deletions(-) diff --git a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/BulkInsertFunctionWrapper.java b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/BulkInsertFunctionWrapper.java index 26f79d870d66..159ead016fb5 100644 --- a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/BulkInsertFunctionWrapper.java +++ b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/BulkInsertFunctionWrapper.java @@ -167,9 +167,7 @@ public class BulkInsertFunctionWrapper<I> implements TestFunctionWrapper<I> { } public void coordinatorFails() throws Exception { - this.coordinator.close(); - this.coordinator.start(); - this.coordinator.setExecutor(new MockCoordinatorExecutor(coordinatorContext)); + // Do nothing since there is no state recovery for bulk insert. } public void restartCoordinator() throws Exception { diff --git a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/InsertFunctionWrapper.java b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/InsertFunctionWrapper.java index 8c0369a889c7..9a8262155973 100644 --- a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/InsertFunctionWrapper.java +++ b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/InsertFunctionWrapper.java @@ -49,6 +49,8 @@ import org.apache.flink.streaming.util.MockStreamTaskBuilder; import org.apache.flink.table.data.RowData; import org.apache.flink.table.types.logical.RowType; +import java.util.Map; +import java.util.TreeMap; import java.util.concurrent.CompletableFuture; /** @@ -70,6 +72,7 @@ public class InsertFunctionWrapper<I> implements TestFunctionWrapper<I> { private final boolean asyncClustering; private ClusteringFunctionWrapper clusteringFunctionWrapper; + private final TreeMap<Long, byte[]> coordinatorStateStore; /** * Append write function. @@ -97,6 +100,7 @@ public class InsertFunctionWrapper<I> implements TestFunctionWrapper<I> { this.coordinatorContext = new MockOperatorCoordinatorContext(new OperatorID(), 1); this.coordinator = new StreamWriteOperatorCoordinator(conf, this.coordinatorContext); this.stateInitializationContext = new MockStateInitializationContext(); + this.coordinatorStateStore = new TreeMap<>(); this.asyncClustering = OptionsResolver.needsAsyncClustering(conf); StreamConfig streamConfig = new StreamConfig(conf); @@ -142,8 +146,10 @@ public class InsertFunctionWrapper<I> implements TestFunctionWrapper<I> { } public void checkpointFunction(long checkpointId) throws Exception { + CompletableFuture<byte[]> completableFuture = new CompletableFuture<>(); // checkpoint the coordinator first - this.coordinator.checkpointCoordinator(checkpointId, new CompletableFuture<>()); + this.coordinator.checkpointCoordinator(checkpointId, completableFuture); + this.coordinatorStateStore.put(checkpointId, completableFuture.get()); writeFunction.snapshotState(new MockFunctionSnapshotContext(checkpointId)); stateInitializationContext.checkpointBegin(checkpointId); @@ -167,9 +173,15 @@ public class InsertFunctionWrapper<I> implements TestFunctionWrapper<I> { } public void coordinatorFails() throws Exception { - this.coordinator.close(); - this.coordinator.start(); - this.coordinator.setExecutor(new MockCoordinatorExecutor(coordinatorContext)); + resetCoordinatorToCheckpoint(); + } + + private void resetCoordinatorToCheckpoint() { + if (coordinatorStateStore.isEmpty()) { + return; + } + Map.Entry<Long, byte[]> latestState = this.coordinatorStateStore.lastEntry(); + this.coordinator.resetToCheckpoint(latestState.getKey(), latestState.getValue()); } public void restartCoordinator() throws Exception { diff --git a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/StreamWriteFunctionWrapper.java b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/StreamWriteFunctionWrapper.java index 7b1684c8e629..131c14602cde 100644 --- a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/StreamWriteFunctionWrapper.java +++ b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/StreamWriteFunctionWrapper.java @@ -368,13 +368,7 @@ public class StreamWriteFunctionWrapper<I> implements TestFunctionWrapper<I> { } public void coordinatorFails() throws Exception { - this.coordinator.close(); - if (isStreamingWriteIndexEnabled) { - this.coordinator.setExecutor(new MockCoordinatorExecutor(coordinatorContext)); - } resetCoordinatorToCheckpoint(); - this.coordinator.start(); - this.coordinator.setExecutor(new MockCoordinatorExecutor(coordinatorContext)); } public void restartCoordinator() throws Exception {
