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());
+ }
+}