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 7430134280ae fix(flink): serve instant-time requests asynchronously to 
avoid coordination RPC timeout (#19960)
7430134280ae is described below

commit 7430134280ae5684daee64ece7660a77577bf65f
Author: fhan <[email protected]>
AuthorDate: Wed Sep 23 15:15:27 2026 +0800

    fix(flink): serve instant-time requests asynchronously to avoid 
coordination RPC timeout (#19960)
    
    * fix(flink): serve instant-time requests asynchronously to avoid 
coordination RPC timeout
    * refactor(flink): simplify asynchronous instant polling
    
    ---------
    
    Co-authored-by: fhan <[email protected]>
    Co-authored-by: danny0405 <[email protected]>
---
 .../hudi/sink/StreamWriteOperatorCoordinator.java  | 33 ++++-----
 .../hudi/sink/bulk/BulkInsertWriteFunction.java    |  3 +-
 .../sink/common/AbstractStreamWriteFunction.java   |  4 +-
 .../org/apache/hudi/sink/event/Correspondent.java  | 60 ++++++++++++++--
 .../apache/hudi/sink/utils/NonThrownExecutor.java  | 24 ++++++-
 .../sink/TestStreamWriteOperatorCoordinator.java   | 70 +++++++++++++++---
 .../common/TestAbstractStreamWriteFunction.java    | 15 ++--
 .../sink/event/TestCorrespondentEventModels.java   | 75 +++++++++++++++++++
 .../utils/BucketStreamWriteFunctionWrapper.java    |  1 +
 .../hudi/sink/utils/BulkInsertFunctionWrapper.java |  2 +
 .../hudi/sink/utils/InsertFunctionWrapper.java     |  2 +
 .../apache/hudi/sink/utils/MockCorrespondent.java  |  9 +--
 .../sink/utils/MockCorrespondentWithTimeout.java   | 12 +---
 .../sink/utils/StreamWriteFunctionWrapper.java     |  2 +
 .../hudi/sink/utils/TestNonThrownExecutor.java     | 84 ++++++++++++++++++++++
 15 files changed, 342 insertions(+), 54 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 e34acac60f98..6e1cd7b71fd3 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,22 +412,23 @@ public class StreamWriteOperatorCoordinator
   }
 
   private CompletableFuture<CoordinationResponse> 
handleInstantRequest(Correspondent.InstantTimeRequest request) {
-    CompletableFuture<CoordinationResponse> response = new 
CompletableFuture<>();
-    instantRequestExecutor.execute(() -> {
-      long checkpointId = request.getCheckpointId();
-      Pair<String, EventBuffer> instantTimeAndEventBuffer = 
this.eventBuffers.getInstantAndEventBuffer(checkpointId);
-      final String instantTime;
-      if (instantTimeAndEventBuffer == null) {
-        // wait until previous instants are committed.
-        eventBuffers.awaitAllInstantsToCompleteIfNecessary();
-        instantTime = startInstant();
-        this.eventBuffers.initNewEventBuffer(checkpointId, instantTime);
-      } else {
-        instantTime = instantTimeAndEventBuffer.getLeft();
-      }
-      
response.complete(CoordinationResponseSerDe.wrap(Correspondent.InstantTimeResponse.getInstance(instantTime)));
-    }, "request instant time");
-    return response;
+    if (instantRequestExecutor.hasRunningTasks()) {
+      return 
CompletableFuture.completedFuture(CoordinationResponseSerDe.wrap(Correspondent.InstantTimeResponse.getInstance(null)));
+    }
+    long checkpointId = request.getCheckpointId();
+    Pair<String, EventBuffer> instantTimeAndEventBuffer = 
this.eventBuffers.getInstantAndEventBuffer(checkpointId);
+    if (instantTimeAndEventBuffer == null) {
+      instantRequestExecutor.execute(() -> {
+        if (this.eventBuffers.getInstantAndEventBuffer(checkpointId) == null) {
+          // Wait until previous instants are committed.
+          eventBuffers.awaitAllInstantsToCompleteIfNecessary();
+          this.eventBuffers.initNewEventBuffer(checkpointId, startInstant());
+        }
+      }, "request instant time");
+      instantTimeAndEventBuffer = 
this.eventBuffers.getInstantAndEventBuffer(checkpointId);
+    }
+    String instantTime = instantTimeAndEventBuffer == null ? null : 
instantTimeAndEventBuffer.getLeft();
+    return 
CompletableFuture.completedFuture(CoordinationResponseSerDe.wrap(Correspondent.InstantTimeResponse.getInstance(instantTime)));
   }
 
   private CompletableFuture<CoordinationResponse> 
handleInFlightInstantsRequest(Correspondent.InflightInstantsRequest request) {
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bulk/BulkInsertWriteFunction.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bulk/BulkInsertWriteFunction.java
index 593b00333dcf..dbcbe5caf971 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bulk/BulkInsertWriteFunction.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bulk/BulkInsertWriteFunction.java
@@ -21,6 +21,7 @@ package org.apache.hudi.sink.bulk;
 import org.apache.hudi.client.HoodieFlinkWriteClient;
 import org.apache.hudi.client.WriteStatus;
 import org.apache.hudi.common.model.WriteOperationType;
+import org.apache.hudi.configuration.FlinkOptions;
 import org.apache.hudi.sink.StreamWriteOperatorCoordinator;
 import org.apache.hudi.sink.buffer.MemorySegmentPoolFactory;
 import org.apache.hudi.sink.common.AbstractWriteFunction;
@@ -169,6 +170,6 @@ public class BulkInsertWriteFunction<I>
    * Returns the instant to write.
    */
   private String instantToWrite() {
-    return this.correspondent.requestInstantTime(-1L);
+    return this.correspondent.requestInstantTime(-1L, 
this.config.get(FlinkOptions.WRITE_COMMIT_ACK_TIMEOUT));
   }
 }
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/common/AbstractStreamWriteFunction.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/common/AbstractStreamWriteFunction.java
index 567c312a156e..19b72a4020d8 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/common/AbstractStreamWriteFunction.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/common/AbstractStreamWriteFunction.java
@@ -23,6 +23,7 @@ import org.apache.hudi.client.WriteStatus;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.table.timeline.HoodieTimeline;
 import org.apache.hudi.common.util.ValidationUtils;
+import org.apache.hudi.configuration.FlinkOptions;
 import org.apache.hudi.sink.StreamWriteOperatorCoordinator;
 import org.apache.hudi.sink.buffer.MemorySegmentPoolFactory;
 import org.apache.hudi.sink.event.CommitAckEvent;
@@ -266,7 +267,8 @@ public abstract class AbstractStreamWriteFunction<I>
    * @return The instant time
    */
   protected String instantToWrite(boolean hasData) {
-    return 
Preconditions.checkNotNull(this.correspondent.requestInstantTime(this.checkpointId),
+    return Preconditions.checkNotNull(
+        this.correspondent.requestInstantTime(this.checkpointId, 
this.config.get(FlinkOptions.WRITE_COMMIT_ACK_TIMEOUT)),
         "No in-flight instant for checkpoint id: " + checkpointId);
   }
 
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/event/Correspondent.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/event/Correspondent.java
index 9bc8f8242ec5..dbd9a5b30aa6 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/event/Correspondent.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/event/Correspondent.java
@@ -34,6 +34,7 @@ import org.apache.flink.util.SerializedValue;
 import java.io.IOException;
 import java.util.HashMap;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
 
 /**
  * Correspondent between a write task with the coordinator.
@@ -42,6 +43,16 @@ import java.util.Map;
 @Getter
 public class Correspondent {
 
+  /**
+   * Initial backoff in milliseconds between two instant-time polls, allowing 
for instant creation latency.
+   */
+  private static final long POLL_BASE_MS = 200L;
+
+  /**
+   * Upper bound in milliseconds of the backoff between two instant-time polls.
+   */
+  private static final long POLL_CAP_MS = 1000L;
+
   private final OperatorID operatorID;
   private final TaskOperatorEventGateway gateway;
 
@@ -64,16 +75,50 @@ public class Correspondent {
   }
 
   /**
-   * Sends a request to the coordinator to fetch the instant time.
+   * Requests the instant time for the given checkpoint from the coordinator.
+   *
+   * <p>Polls with capped exponential backoff until the instant is non-null or 
the timeout expires.
+   * Request failures are propagated immediately.
+   *
+   * @param checkpointId The checkpoint id (or -1 for bulk insert)
+   * @param pollBudgetMs The overall budget to wait for an instant, in 
milliseconds
+   *
+   * @return the instant time to write with
    */
-  public String requestInstantTime(long checkpointId) {
+  public String requestInstantTime(long checkpointId, long pollBudgetMs) {
+    final long deadlineNanos = System.nanoTime() + 
TimeUnit.MILLISECONDS.toNanos(pollBudgetMs);
+    long backoffMs = POLL_BASE_MS;
     try {
-      InstantTimeResponse response = 
CoordinationResponseSerDe.unwrap(this.gateway.sendRequestToCoordinator(this.operatorID,
-          new 
SerializedValue<>(InstantTimeRequest.getInstance(checkpointId))).get());
-      return response.getInstant();
+      do {
+        String instant = fetchInstantTimeResponse(checkpointId).getInstant();
+        if (instant != null) {
+          return instant;
+        }
+        long remainingNanos = deadlineNanos - System.nanoTime();
+        if (remainingNanos <= 0) {
+          break;
+        }
+        TimeUnit.NANOSECONDS.sleep(Math.min(remainingNanos, 
TimeUnit.MILLISECONDS.toNanos(backoffMs)));
+        backoffMs = Math.min(backoffMs * 2, POLL_CAP_MS);
+      } while (System.nanoTime() < deadlineNanos);
+    } catch (InterruptedException e) {
+      Thread.currentThread().interrupt();
+      throw new HoodieException("Interrupted while requesting the instant time 
from the coordinator", e);
     } catch (Exception e) {
-      throw new HoodieException("Error requesting the instant time from the 
coordinator", e);
+      throw new HoodieException(
+          "Error requesting the instant time from the coordinator for 
checkpoint " + checkpointId, e);
     }
+    throw new HoodieException("Timeout waiting for the instant time from the 
coordinator for checkpoint " + checkpointId);
+  }
+
+  /**
+   * Sends a single instant-time request to the coordinator and returns its 
response.
+   *
+   * <p>Isolated so tests can stub the transport while reusing the poll loop 
in {@link #requestInstantTime}.
+   */
+  protected InstantTimeResponse fetchInstantTimeResponse(long checkpointId) 
throws Exception {
+    return 
CoordinationResponseSerDe.unwrap(this.gateway.sendRequestToCoordinator(this.operatorID,
+        new 
SerializedValue<>(InstantTimeRequest.getInstance(checkpointId))).get());
   }
 
   /**
@@ -121,6 +166,9 @@ public class Correspondent {
   @Getter
   public static class InstantTimeResponse implements CoordinationResponse {
 
+    /**
+     * The instant time, or null while the instant is still being created.
+     */
     private final String instant;
 
     public static InstantTimeResponse getInstance(String instant) {
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/NonThrownExecutor.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/NonThrownExecutor.java
index 4da4a7718a67..82d3193e874f 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/NonThrownExecutor.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/NonThrownExecutor.java
@@ -29,8 +29,10 @@ import java.util.Objects;
 import java.util.concurrent.ExecutionException;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
+import java.util.concurrent.RejectedExecutionException;
 import java.util.concurrent.ThreadFactory;
 import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
 import java.util.function.Supplier;
 
 /**
@@ -47,6 +49,8 @@ public class NonThrownExecutor implements AutoCloseable {
    */
   private final ExecutorService executor;
 
+  private final AtomicInteger pendingTasks = new AtomicInteger();
+
   /**
    * Exception hook for post-exception handling.
    */
@@ -88,7 +92,12 @@ public class NonThrownExecutor implements AutoCloseable {
       final ExceptionHook hook,
       final String actionName,
       final Object... actionParams) {
-    executor.execute(wrapAction(action, hook, actionName, actionParams));
+    try {
+      executor.execute(wrapAction(action, hook, actionName, actionParams));
+    } catch (RejectedExecutionException e) {
+      pendingTasks.decrementAndGet();
+      handleException(e, hook, getActionString(actionName, actionParams));
+    }
   }
 
   /**
@@ -97,6 +106,9 @@ public class NonThrownExecutor implements AutoCloseable {
   public void executeSync(ThrowingRunnable<Throwable> action, String 
actionName, Object... actionParams) {
     try {
       executor.submit(wrapAction(action, this.exceptionHook, actionName, 
actionParams)).get();
+    } catch (RejectedExecutionException e) {
+      pendingTasks.decrementAndGet();
+      handleException(e, this.exceptionHook, getActionString(actionName, 
actionParams));
     } catch (InterruptedException e) {
       handleException(e, this.exceptionHook, getActionString(actionName, 
actionParams));
     } catch (ExecutionException e) {
@@ -105,6 +117,13 @@ public class NonThrownExecutor implements AutoCloseable {
     }
   }
 
+  /**
+   * Returns whether any task is queued or running.
+   */
+  public boolean hasRunningTasks() {
+    return pendingTasks.get() > 0;
+  }
+
   @Override
   public void close() throws Exception {
     if (executor != null) {
@@ -125,6 +144,7 @@ public class NonThrownExecutor implements AutoCloseable {
       final String actionName,
       final Object... actionParams) {
 
+    pendingTasks.incrementAndGet();
     return () -> {
       final Supplier<String> actionString = getActionString(actionName, 
actionParams);
       try {
@@ -132,6 +152,8 @@ public class NonThrownExecutor implements AutoCloseable {
         logger.info("Executor executes action [{}] success!", 
actionString.get());
       } catch (Throwable t) {
         handleException(t, hook, actionString);
+      } finally {
+        pendingTasks.decrementAndGet();
       }
     };
   }
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 352229beff9c..d26e818d6c1b 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
@@ -48,6 +48,7 @@ import org.apache.hudi.sink.muttley.AthenaIngestionGateway;
 import org.apache.hudi.sink.utils.CoordinationResponseSerDe;
 import org.apache.hudi.sink.utils.EventBuffers;
 import org.apache.hudi.sink.utils.MockCoordinatorExecutor;
+import org.apache.hudi.sink.utils.MockCorrespondent;
 import org.apache.hudi.sink.utils.NonThrownExecutor;
 import org.apache.hudi.storage.HoodieStorage;
 import org.apache.hudi.storage.StoragePath;
@@ -64,6 +65,7 @@ import 
org.apache.flink.runtime.operators.coordination.MockOperatorCoordinatorCo
 import org.apache.flink.runtime.operators.coordination.OperatorCoordinator;
 import org.apache.flink.runtime.operators.coordination.OperatorEvent;
 import org.apache.flink.util.FileUtils;
+import org.apache.flink.util.function.ThrowingRunnable;
 import org.apache.hadoop.fs.FileSystem;
 import org.apache.hadoop.fs.Path;
 import org.junit.jupiter.api.AfterEach;
@@ -86,6 +88,7 @@ import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.TimeUnit;
 
 import static 
org.apache.hudi.common.testutils.HoodieTestUtils.INSTANT_GENERATOR;
@@ -382,8 +385,8 @@ public class TestStreamWriteOperatorCoordinator {
         
coordinator.getEventBuffer().getDataWriteEventBuffer()[0].getWriteStatuses().size(),
 is(1));
 
     long nextCkpId = 1;
-    
coordinator.handleCoordinationRequest(Correspondent.InstantTimeRequest.getInstance(nextCkpId));
-    OperatorEvent event4 = createOperatorEvent(0, nextCkpId, "002", "par1", 
false, false, 0.1);
+    String instant2 = requestInstantTime(nextCkpId);
+    OperatorEvent event4 = createOperatorEvent(0, nextCkpId, instant2, "par1", 
false, false, 0.1);
     coordinator.handleEventFromOperator(0, event4);
     assertThat("First instant is not committed yet, new event should not 
override the old event",
         
coordinator.getEventBuffer(0).getDataWriteEventBuffer()[0].getWriteStatuses().size(),
 is(1));
@@ -821,6 +824,40 @@ public class TestStreamWriteOperatorCoordinator {
     assertEquals(instant2, inflightInstants.get(2L));
   }
 
+  @Test
+  void testInstantRequestPollsWhileCreationBlockedThenSucceeds() throws 
Exception {
+    MockOperatorCoordinatorContext ctx = (MockOperatorCoordinatorContext) 
coordinator.getContext();
+    CountDownLatch lockHeld = new CountDownLatch(1);
+    NonThrownExecutor gatedWorker = Mockito.spy(new 
GatedInstantRequestExecutor(
+        Mockito.mock(Logger.class),
+        (errMsg, t) -> ctx.failJob(new HoodieException(errMsg, t)),
+        lockHeld));
+    coordinator.setInstantRequestExecutor(gatedWorker);
+
+    try {
+      // The first request starts creation; subsequent requests must not queue 
more work while it is blocked.
+      for (long checkpointId : new long[] {1L, 1L, 2L}) {
+        Correspondent.InstantTimeResponse pending = 
CoordinationResponseSerDe.unwrap(
+            
coordinator.handleCoordinationRequest(Correspondent.InstantTimeRequest.getInstance(checkpointId))
+                .get(1, TimeUnit.SECONDS));
+        assertNull(pending.getInstant());
+      }
+      assertTrue(gatedWorker.hasRunningTasks());
+      Mockito.verify(gatedWorker, Mockito.times(1)).execute(Mockito.any(), 
Mockito.eq("request instant time"));
+    } finally {
+      lockHeld.countDown();
+    }
+
+    String instant = requestInstantTime(1L);
+    assertNotNull(instant);
+    assertEquals(instant, requestInstantTime(1L), "Repeated requests must 
reuse the instant");
+    assertFalse(ctx.isJobFailed());
+    HoodieTimeline inflights = 
StreamerUtil.createMetaClient(TestConfigurations.getDefaultConf(tempFile.getAbsolutePath()))
+        .reloadActiveTimeline().filterInflights();
+    assertEquals(1, inflights.countInstants());
+    assertTrue(inflights.containsInstant(instant));
+  }
+
   // -------------------------------------------------------------------------
   //  Utilities
   // -------------------------------------------------------------------------
@@ -830,12 +867,8 @@ public class TestStreamWriteOperatorCoordinator {
   }
 
   private String requestInstantTime(StreamWriteOperatorCoordinator 
coordinator, long checkpointId) {
-    try {
-      Correspondent.InstantTimeResponse response = 
CoordinationResponseSerDe.unwrap(coordinator.handleCoordinationRequest(Correspondent.InstantTimeRequest.getInstance(checkpointId)).get());
-      return response.getInstant();
-    } catch (Exception e) {
-      throw new HoodieException("Error requesting the instant time from the 
coordinator", e);
-    }
+    return new MockCorrespondent(coordinator)
+        .requestInstantTime(checkpointId, TimeUnit.SECONDS.toMillis(10));
   }
 
   private void resetToMergeOnRead(Configuration conf) throws Exception {
@@ -863,6 +896,27 @@ public class TestStreamWriteOperatorCoordinator {
     return new StreamWriteOperatorCoordinator(conf, coordinatorContext);
   }
 
+  /**
+   * A real single-thread instant-request worker whose submitted creation task 
blocks on a latch before
+   * running, to deterministically delay instant creation (simulating a held 
table lock).
+   */
+  private static final class GatedInstantRequestExecutor extends 
NonThrownExecutor {
+    private final CountDownLatch gate;
+
+    private GatedInstantRequestExecutor(Logger logger, ExceptionHook 
exceptionHook, CountDownLatch gate) {
+      super(logger, null, exceptionHook, true);
+      this.gate = gate;
+    }
+
+    @Override
+    public void execute(ThrowingRunnable<Throwable> action, String actionName, 
Object... actionParams) {
+      super.execute(() -> {
+        gate.await();
+        action.run();
+      }, actionName, actionParams);
+    }
+  }
+
   private String mockWriteWithMetadata(long checkpointId) {
     String instant = requestInstantTime(checkpointId);
     OperatorEvent event = createOperatorEvent(0, checkpointId, instant, 
"par1", false, true, 0.1);
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/common/TestAbstractStreamWriteFunction.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/common/TestAbstractStreamWriteFunction.java
index b87ee0aa0a95..bfd9e06184fd 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/common/TestAbstractStreamWriteFunction.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/common/TestAbstractStreamWriteFunction.java
@@ -21,6 +21,7 @@ package org.apache.hudi.sink.common;
 import org.apache.hudi.client.HoodieFlinkWriteClient;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.table.timeline.HoodieTimeline;
+import org.apache.hudi.configuration.FlinkOptions;
 import org.apache.hudi.sink.event.Correspondent;
 import org.apache.hudi.sink.event.WriteMetadataEvent;
 import org.apache.hudi.sink.utils.MockOperatorStateStore;
@@ -84,7 +85,7 @@ class TestAbstractStreamWriteFunction {
     initialize(-1L, attempt);
 
     assertEquals("002", function.instantToWrite(true));
-    verify(correspondent).requestInstantTime(-1L);
+    verify(correspondent).requestInstantTime(-1L, instantRequestPollBudget());
     if (attempt == 0) {
       assertTrue(events.isEmpty());
     } else {
@@ -103,7 +104,7 @@ class TestAbstractStreamWriteFunction {
     initialize(42L, attempt);
 
     assertEquals("002", function.instantToWrite(true));
-    verify(correspondent).requestInstantTime(42L);
+    verify(correspondent).requestInstantTime(42L, instantRequestPollBudget());
     if (attempt == 0) {
       assertTrue(events.isEmpty());
     } else {
@@ -116,7 +117,7 @@ class TestAbstractStreamWriteFunction {
     assertEquals("002", snapshot.getInstantTime());
     assertTrue(snapshot.isBootstrap());
     function.instantToWrite(true);
-    verify(correspondent).requestInstantTime(43L);
+    verify(correspondent).requestInstantTime(43L, instantRequestPollBudget());
   }
 
   @Test
@@ -140,7 +141,7 @@ class TestAbstractStreamWriteFunction {
     assertEquals(41L, bootstrap.getCheckpointId());
     assertEquals("001", bootstrap.getInstantTime());
     function.instantToWrite(true);
-    verify(correspondent).requestInstantTime(42L);
+    verify(correspondent).requestInstantTime(42L, instantRequestPollBudget());
   }
 
   private void initialize(long checkpointId, int attempt) throws Exception {
@@ -149,7 +150,7 @@ class TestAbstractStreamWriteFunction {
     when(context.getOperatorStateStore()).thenReturn(stateStore);
     when(context.isRestored()).thenReturn(checkpointId >= 0);
     when(context.getRestoredCheckpointId()).thenReturn(checkpointId >= 0 ? 
OptionalLong.of(checkpointId) : OptionalLong.empty());
-    when(correspondent.requestInstantTime(anyLong())).thenReturn("002");
+    when(correspondent.requestInstantTime(anyLong(), 
anyLong())).thenReturn("002");
     function.setRuntimeContext(runtimeContext);
     function.setCorrespondent(correspondent);
     function.setOperatorEventGateway(events::add);
@@ -168,6 +169,10 @@ class TestAbstractStreamWriteFunction {
     }
   }
 
+  private long instantRequestPollBudget() {
+    return conf.get(FlinkOptions.WRITE_COMMIT_ACK_TIMEOUT);
+  }
+
   private void assertCleanupEvent(long checkpointId) {
     assertEquals(1, events.size());
     WriteMetadataEvent bootstrap = (WriteMetadataEvent) events.get(0);
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/event/TestCorrespondentEventModels.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/event/TestCorrespondentEventModels.java
index 8dc7102aad25..3539e0d8d5ca 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/event/TestCorrespondentEventModels.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/event/TestCorrespondentEventModels.java
@@ -18,15 +18,23 @@
 
 package org.apache.hudi.sink.event;
 
+import org.apache.hudi.exception.HoodieException;
+
 import org.apache.flink.runtime.jobgraph.OperatorID;
 import org.apache.flink.runtime.jobgraph.tasks.TaskOperatorEventGateway;
 import org.junit.jupiter.api.Test;
 
+import java.io.IOException;
 import java.util.HashMap;
+import java.util.concurrent.atomic.AtomicInteger;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.Mockito.mock;
 
 class TestCorrespondentEventModels {
@@ -41,6 +49,7 @@ class TestCorrespondentEventModels {
 
     assertEquals(9L, 
Correspondent.InstantTimeRequest.getInstance(9L).getCheckpointId());
     assertEquals("001", 
Correspondent.InstantTimeResponse.getInstance("001").getInstant());
+    
assertNull(Correspondent.InstantTimeResponse.getInstance(null).getInstant());
     assertNotNull(Correspondent.InflightInstantsRequest.getInstance());
 
     HashMap<Long, String> instants = new HashMap<>();
@@ -56,4 +65,70 @@ class TestCorrespondentEventModels {
     event.setCheckpointId(12L);
     assertEquals(12L, event.getCheckpointId());
   }
+
+  @Test
+  void testInstantRequestPollsUntilReady() {
+    AtomicInteger requestCount = new AtomicInteger();
+    Correspondent correspondent = new Correspondent() {
+      @Override
+      protected InstantTimeResponse fetchInstantTimeResponse(long 
checkpointId) {
+        return InstantTimeResponse.getInstance(requestCount.getAndIncrement() 
== 0 ? null : "001");
+      }
+    };
+
+    assertEquals("001", correspondent.requestInstantTime(9L, 10_000L));
+    assertEquals(2, requestCount.get());
+  }
+
+  @Test
+  void testInstantRequestTimeoutDoesNotRetry() {
+    AtomicInteger requestCount = new AtomicInteger();
+    Correspondent correspondent = new Correspondent() {
+      @Override
+      protected InstantTimeResponse fetchInstantTimeResponse(long 
checkpointId) {
+        requestCount.incrementAndGet();
+        return InstantTimeResponse.getInstance(null);
+      }
+    };
+
+    HoodieException error = assertThrows(HoodieException.class, () -> 
correspondent.requestInstantTime(9L, 0L));
+    assertEquals("Timeout waiting for the instant time from the coordinator 
for checkpoint 9", error.getMessage());
+    assertEquals(1, requestCount.get());
+  }
+
+  @Test
+  void testInstantPollingPreservesInterrupt() {
+    Correspondent correspondent = new Correspondent() {
+      @Override
+      protected InstantTimeResponse fetchInstantTimeResponse(long 
checkpointId) {
+        return InstantTimeResponse.getInstance(null);
+      }
+    };
+
+    Thread.currentThread().interrupt();
+    try {
+      HoodieException error = assertThrows(HoodieException.class, () -> 
correspondent.requestInstantTime(9L, 10_000L));
+      assertInstanceOf(InterruptedException.class, error.getCause());
+      assertTrue(Thread.currentThread().isInterrupted());
+    } finally {
+      Thread.interrupted();
+    }
+  }
+
+  @Test
+  void testInstantRequestFailureIsNotRetried() {
+    AtomicInteger requestCount = new AtomicInteger();
+    Correspondent correspondent = new Correspondent() {
+      @Override
+      protected InstantTimeResponse fetchInstantTimeResponse(long 
checkpointId) throws Exception {
+        requestCount.incrementAndGet();
+        throw new IOException("request failed");
+      }
+    };
+
+    HoodieException error = assertThrows(
+        HoodieException.class, () -> correspondent.requestInstantTime(9L, 
10_000L));
+    assertEquals("request failed", error.getCause().getMessage());
+    assertEquals(1, requestCount.get());
+  }
 }
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/BucketStreamWriteFunctionWrapper.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/BucketStreamWriteFunctionWrapper.java
index 80f6e7e33471..99dbd17c6352 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/BucketStreamWriteFunctionWrapper.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/BucketStreamWriteFunctionWrapper.java
@@ -126,6 +126,7 @@ public class BucketStreamWriteFunctionWrapper<I> implements 
TestFunctionWrapper<
   public void openFunction() throws Exception {
     this.coordinator.start();
     this.coordinator.setExecutor(new 
MockCoordinatorExecutor(coordinatorContext));
+    this.coordinator.setInstantRequestExecutor(new 
MockCoordinatorExecutor(coordinatorContext));
     toHoodieFunction = new RowDataToHoodieFunction<>(rowType, conf);
     toHoodieFunction.setRuntimeContext(runtimeContext);
     toHoodieFunction.open(conf);
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 1ad8ccdeec4f..5e45f3d2ee47 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
@@ -117,6 +117,7 @@ public class BulkInsertFunctionWrapper<I> implements 
TestFunctionWrapper<I> {
   public void openFunction() throws Exception {
     this.coordinator.start();
     this.coordinator.setExecutor(new 
MockCoordinatorExecutor(coordinatorContext));
+    this.coordinator.setInstantRequestExecutor(new 
MockCoordinatorExecutor(coordinatorContext));
     setupWriteFunction();
     setupMapFunction();
     if (needSortInput) {
@@ -184,6 +185,7 @@ public class BulkInsertFunctionWrapper<I> implements 
TestFunctionWrapper<I> {
     this.coordinator = new StreamWriteOperatorCoordinator(conf, 
this.coordinatorContext);
     this.coordinator.start();
     this.coordinator.setExecutor(new 
MockCoordinatorExecutor(coordinatorContext));
+    this.coordinator.setInstantRequestExecutor(new 
MockCoordinatorExecutor(coordinatorContext));
   }
 
   public void checkpointFails(long checkpointId) {
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 aabef24d5d04..c715a0c518e6 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
@@ -117,6 +117,7 @@ public class InsertFunctionWrapper<I> implements 
TestFunctionWrapper<I> {
   public void openFunction() throws Exception {
     this.coordinator.start();
     this.coordinator.setExecutor(new 
MockCoordinatorExecutor(coordinatorContext));
+    this.coordinator.setInstantRequestExecutor(new 
MockCoordinatorExecutor(coordinatorContext));
 
     setupWriteFunction();
 
@@ -191,6 +192,7 @@ public class InsertFunctionWrapper<I> implements 
TestFunctionWrapper<I> {
     this.coordinator = new StreamWriteOperatorCoordinator(conf, 
this.coordinatorContext);
     this.coordinator.start();
     this.coordinator.setExecutor(new 
MockCoordinatorExecutor(coordinatorContext));
+    this.coordinator.setInstantRequestExecutor(new 
MockCoordinatorExecutor(coordinatorContext));
   }
 
   public void checkpointFails(long checkpointId) {
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/MockCorrespondent.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/MockCorrespondent.java
index 6076048c3d34..f4852f5ff702 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/MockCorrespondent.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/MockCorrespondent.java
@@ -35,13 +35,8 @@ public class MockCorrespondent extends Correspondent {
   }
 
   @Override
-  public String requestInstantTime(long checkpointId) {
-    try {
-      InstantTimeResponse response = 
CoordinationResponseSerDe.unwrap(this.coordinator.handleCoordinationRequest(InstantTimeRequest.getInstance(checkpointId)).get());
-      return response.getInstant();
-    } catch (Exception e) {
-      throw new HoodieException("Error requesting the instant time from the 
coordinator", e);
-    }
+  protected InstantTimeResponse fetchInstantTimeResponse(long checkpointId) 
throws Exception {
+    return 
CoordinationResponseSerDe.unwrap(this.coordinator.handleCoordinationRequest(InstantTimeRequest.getInstance(checkpointId)).get());
   }
 
   @Override
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/MockCorrespondentWithTimeout.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/MockCorrespondentWithTimeout.java
index 22b01a6432e1..4dcada9d0a17 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/MockCorrespondentWithTimeout.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/MockCorrespondentWithTimeout.java
@@ -19,7 +19,6 @@
 package org.apache.hudi.sink.utils;
 
 import org.apache.hudi.configuration.FlinkOptions;
-import org.apache.hudi.exception.HoodieException;
 import org.apache.hudi.sink.StreamWriteOperatorCoordinator;
 import org.apache.hudi.sink.event.Correspondent;
 
@@ -44,13 +43,8 @@ public class MockCorrespondentWithTimeout extends 
Correspondent {
   }
 
   @Override
-  public String requestInstantTime(long checkpointId) {
-    try {
-      CompletableFuture<CoordinationResponse> future = 
this.coordinator.handleCoordinationRequest(InstantTimeRequest.getInstance(checkpointId));
-      InstantTimeResponse response = 
CoordinationResponseSerDe.unwrap(future.get(commitAckTimeout, 
TimeUnit.MILLISECONDS));
-      return response.getInstant();
-    } catch (Exception e) {
-      throw new HoodieException("Error requesting the instant time from the 
coordinator", e);
-    }
+  protected InstantTimeResponse fetchInstantTimeResponse(long checkpointId) 
throws Exception {
+    CompletableFuture<CoordinationResponse> future = 
this.coordinator.handleCoordinationRequest(InstantTimeRequest.getInstance(checkpointId));
+    return CoordinationResponseSerDe.unwrap(future.get(commitAckTimeout, 
TimeUnit.MILLISECONDS));
   }
 }
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 baab42555c1e..fbdb24214769 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
@@ -178,6 +178,7 @@ public class StreamWriteFunctionWrapper<I> implements 
TestFunctionWrapper<I> {
     resetCoordinatorToCheckpoint();
     this.coordinator.start();
     this.coordinator.setExecutor(new 
MockCoordinatorExecutor(coordinatorContext));
+    this.coordinator.setInstantRequestExecutor(new 
MockCoordinatorExecutor(coordinatorContext));
     toHoodieFunction = new RowDataToHoodieFunction<>(rowType, conf);
     toHoodieFunction.setRuntimeContext(runtimeContext);
     toHoodieFunction.open(conf);
@@ -377,6 +378,7 @@ public class StreamWriteFunctionWrapper<I> implements 
TestFunctionWrapper<I> {
     resetCoordinatorToCheckpoint();
     this.coordinator.start();
     this.coordinator.setExecutor(new 
MockCoordinatorExecutor(coordinatorContext));
+    this.coordinator.setInstantRequestExecutor(new 
MockCoordinatorExecutor(coordinatorContext));
     this.correspondent = new MockCorrespondent(coordinator);
   }
 
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/TestNonThrownExecutor.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/TestNonThrownExecutor.java
new file mode 100644
index 000000000000..8d268cfe85be
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/TestNonThrownExecutor.java
@@ -0,0 +1,84 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hudi.sink.utils;
+
+import org.junit.jupiter.api.Test;
+import org.slf4j.Logger;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.RejectedExecutionException;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+
+class TestNonThrownExecutor {
+
+  @Test
+  void testTracksTasksThroughCompletionFailureAndRejection() throws Exception {
+    CountDownLatch gate = new CountDownLatch(1);
+    AtomicInteger completed = new AtomicInteger();
+    AtomicInteger failures = new AtomicInteger();
+    AtomicReference<Throwable> failure = new AtomicReference<>();
+    NonThrownExecutor executor = NonThrownExecutor.builder(mock(Logger.class))
+        .exceptionHook((message, error) -> {
+          failures.incrementAndGet();
+          failure.set(error);
+        })
+        .waitForTasksFinish(true)
+        .build();
+    try {
+      assertFalse(executor.hasRunningTasks());
+      executor.execute(() -> {
+        gate.await();
+        throw new IllegalStateException("test failure");
+      }, "blocked task");
+      executor.execute(completed::incrementAndGet, "queued task");
+      assertTrue(executor.hasRunningTasks());
+    } finally {
+      gate.countDown();
+      executor.close();
+    }
+    assertEquals(1, failures.get());
+    assertEquals(1, completed.get());
+    assertFalse(executor.hasRunningTasks());
+    executor.execute(completed::incrementAndGet, "rejected task");
+    assertEquals(2, failures.get());
+    assertInstanceOf(RejectedExecutionException.class, failure.get());
+    assertFalse(executor.hasRunningTasks());
+    executor.executeSync(completed::incrementAndGet, "rejected synchronous 
task");
+    assertEquals(3, failures.get());
+    assertInstanceOf(RejectedExecutionException.class, failure.get());
+    assertFalse(executor.hasRunningTasks());
+
+    AtomicInteger customFailures = new AtomicInteger();
+    executor.execute(completed::incrementAndGet, (message, error) -> {
+      assertInstanceOf(RejectedExecutionException.class, error);
+      customFailures.incrementAndGet();
+    }, "rejected task with custom hook");
+    assertEquals(1, customFailures.get());
+    assertEquals(3, failures.get());
+    assertEquals(1, completed.get());
+    assertFalse(executor.hasRunningTasks());
+  }
+}

Reply via email to