This is an automated email from the ASF dual-hosted git repository.

danny0405 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git


The following commit(s) were added to refs/heads/master by this push:
     new 20037d3c3f0e fix(flink): avoid blocking existing instant lookups on 
instant creation (#20031)
20037d3c3f0e is described below

commit 20037d3c3f0e58de450c61b3fa4edd57e1a24935
Author: Shuo Cheng <[email protected]>
AuthorDate: Thu Sep 24 10:10:35 2026 +0800

    fix(flink): avoid blocking existing instant lookups on instant creation 
(#20031)
---
 .../hudi/sink/StreamWriteOperatorCoordinator.java  |  6 ++----
 .../sink/TestStreamWriteOperatorCoordinator.java   | 25 ++++++++++++++++++++++
 2 files changed, 27 insertions(+), 4 deletions(-)

diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
index 6e1cd7b71fd3..11fccc32e87f 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteOperatorCoordinator.java
@@ -412,12 +412,10 @@ public class StreamWriteOperatorCoordinator
   }
 
   private CompletableFuture<CoordinationResponse> 
handleInstantRequest(Correspondent.InstantTimeRequest request) {
-    if (instantRequestExecutor.hasRunningTasks()) {
-      return 
CompletableFuture.completedFuture(CoordinationResponseSerDe.wrap(Correspondent.InstantTimeResponse.getInstance(null)));
-    }
     long checkpointId = request.getCheckpointId();
+    // Existing instants must remain available while another checkpoint's 
creation is blocked.
     Pair<String, EventBuffer> instantTimeAndEventBuffer = 
this.eventBuffers.getInstantAndEventBuffer(checkpointId);
-    if (instantTimeAndEventBuffer == null) {
+    if (instantTimeAndEventBuffer == null && 
!instantRequestExecutor.hasRunningTasks()) {
       instantRequestExecutor.execute(() -> {
         if (this.eventBuffers.getInstantAndEventBuffer(checkpointId) == null) {
           // Wait until previous instants are committed.
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
index d26e818d6c1b..bd661a9fdbae 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
@@ -858,6 +858,31 @@ public class TestStreamWriteOperatorCoordinator {
     assertTrue(inflights.containsInstant(instant));
   }
 
+  @Test
+  void testReadyInstantRemainsAvailableWhileAnotherCreationIsBlocked() throws 
Exception {
+    String existingInstant = requestInstantTime(1L);
+    MockOperatorCoordinatorContext ctx = (MockOperatorCoordinatorContext) 
coordinator.getContext();
+    CountDownLatch gate = new CountDownLatch(1);
+    NonThrownExecutor worker = new GatedInstantRequestExecutor(
+        Mockito.mock(Logger.class),
+        (errMsg, t) -> ctx.failJob(new HoodieException(errMsg, t)), gate);
+    coordinator.setInstantRequestExecutor(worker);
+    try {
+      Correspondent.InstantTimeResponse pending = 
CoordinationResponseSerDe.unwrap(
+          
coordinator.handleCoordinationRequest(Correspondent.InstantTimeRequest.getInstance(2L))
+              .get(1, TimeUnit.SECONDS));
+      assertNull(pending.getInstant());
+      assertTrue(worker.hasRunningTasks());
+      Correspondent.InstantTimeResponse ready = 
CoordinationResponseSerDe.unwrap(
+          
coordinator.handleCoordinationRequest(Correspondent.InstantTimeRequest.getInstance(1L))
+              .get(1, TimeUnit.SECONDS));
+      assertEquals(existingInstant, ready.getInstant(),
+          "An existing instant must remain available while another 
checkpoint's creation is blocked");
+    } finally {
+      gate.countDown();
+    }
+  }
+
   // -------------------------------------------------------------------------
   //  Utilities
   // -------------------------------------------------------------------------

Reply via email to