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
// -------------------------------------------------------------------------