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 6572fe07dd1074b09d55fb6ebe621b8174f15b41
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 {

Reply via email to