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 fcc03c35bc5 [HUDI-9534] Blocking instant generation for Flink COW
table (#13464)
fcc03c35bc5 is described below
commit fcc03c35bc5bf71abf066922fe4ce773e3489bd2
Author: Shuo Cheng <[email protected]>
AuthorDate: Fri Jun 20 20:05:31 2025 +0800
[HUDI-9534] Blocking instant generation for Flink COW table (#13464)
* use blocking instant generation in Flink writer for COW table
* update default value for commit ack timeout and add test
---------
Co-authored-by: danny0405 <[email protected]>
---
.../apache/hudi/configuration/FlinkOptions.java | 2 +-
.../apache/hudi/configuration/OptionsResolver.java | 7 ++
.../hudi/sink/StreamWriteOperatorCoordinator.java | 38 +++++++--
.../org/apache/hudi/sink/event/Correspondent.java | 4 +-
.../org/apache/hudi/sink/utils/CommitGuard.java | 80 ++++++++++++++++++
...nseSeDe.java => CoordinationResponseSerDe.java} | 2 +-
.../org/apache/hudi/sink/utils/EventBuffers.java | 21 ++++-
.../org/apache/hudi/util/FlinkWriteClients.java | 8 +-
.../sink/TestStreamWriteOperatorCoordinator.java | 12 +--
.../org/apache/hudi/sink/TestWriteCopyOnWrite.java | 97 ++++++++++++++++++++--
.../apache/hudi/sink/utils/MockCorrespondent.java | 3 +-
...dent.java => MockCorrespondentWithTimeout.java} | 18 +++-
.../org/apache/hudi/sink/utils/TestWriteBase.java | 10 +++
13 files changed, 272 insertions(+), 30 deletions(-)
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/FlinkOptions.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/FlinkOptions.java
index 3a13a269b2f..d191bc9a825 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/FlinkOptions.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/FlinkOptions.java
@@ -671,7 +671,7 @@ public class FlinkOptions extends HoodieConfig {
public static final ConfigOption<Long> WRITE_COMMIT_ACK_TIMEOUT =
ConfigOptions
.key("write.commit.ack.timeout")
.longType()
- .defaultValue(-1L) // default at least once
+ .defaultValue(300_000L)
.withDescription("Timeout limit for a writer task after it finishes a
checkpoint and\n"
+ "waits for the instant commit success, only for internal use");
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/OptionsResolver.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/OptionsResolver.java
index 30b480c1f6e..e6b3a361b6d 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/OptionsResolver.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/configuration/OptionsResolver.java
@@ -480,6 +480,13 @@ public class OptionsResolver {
return
HoodieCleanConfig.FAILED_WRITES_CLEANER_POLICY.defaultValue().equalsIgnoreCase(HoodieFailedWritesCleaningPolicy.LAZY.name());
}
+ /**
+ * Returns whether the writers should use blocking instant time generation.
+ */
+ public static boolean isBlockingInstantGeneration(Configuration conf) {
+ return isCowTable(conf) && isUpsertOperation(conf);
+ }
+
// -------------------------------------------------------------------------
// Utilities
// -------------------------------------------------------------------------
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 0c57a2b1d82..4368ca3b365 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
@@ -38,7 +38,8 @@ import org.apache.hudi.hive.HiveSyncTool;
import org.apache.hudi.sink.common.AbstractStreamWriteFunction;
import org.apache.hudi.sink.event.Correspondent;
import org.apache.hudi.sink.event.WriteMetadataEvent;
-import org.apache.hudi.sink.utils.CoordinationResponseSeDe;
+import org.apache.hudi.sink.utils.CommitGuard;
+import org.apache.hudi.sink.utils.CoordinationResponseSerDe;
import org.apache.hudi.sink.utils.EventBuffers;
import org.apache.hudi.sink.utils.ExplicitClassloaderThreadFactory;
import org.apache.hudi.sink.utils.HiveSyncContext;
@@ -194,6 +195,11 @@ public class StreamWriteOperatorCoordinator
*/
private ClientIds clientIds;
+ /**
+ * The commit guard for blocking instant time generation.
+ */
+ private Option<CommitGuard> commitGuardOpt;
+
/**
* Constructs a StreamingSinkOperatorCoordinator.
*
@@ -214,8 +220,10 @@ public class StreamWriteOperatorCoordinator
// setup classloader for APIs that use reflection without taking
ClassLoader param
// reference:
https://stackoverflow.com/questions/1771679/difference-between-threads-context-class-loader-and-normal-classloader
Thread.currentThread().setContextClassLoader(getClass().getClassLoader());
+ this.tableState = TableState.create(conf);
+ initCommitGuard(this.conf);
// initialize event buffer
- this.eventBuffers = EventBuffers.getInstance();
+ this.eventBuffers = EventBuffers.getInstance(this.commitGuardOpt);
this.gateways = new SubtaskGateway[this.parallelism];
try {
// init table, create if not exists.
@@ -224,14 +232,13 @@ public class StreamWriteOperatorCoordinator
this.writeClient = FlinkWriteClients.createWriteClient(conf);
this.writeClient.tryUpgrade(instant, this.metaClient);
initMetadataTable(this.writeClient);
- this.tableState = TableState.create(conf);
// start the executor
this.executor = NonThrownExecutor.builder(LOG)
.threadFactory(getThreadFactory("meta-event-handle"))
.exceptionHook((errMsg, t) -> this.context.failJob(new
HoodieException(errMsg, t)))
.waitForTasksFinish(true).build();
this.instantRequestExecutor = NonThrownExecutor.builder(LOG)
- .threadFactory(getThreadFactory("instant-response"))
+ .threadFactory(getThreadFactory("instant-request"))
.exceptionHook((errMsg, t) -> this.context.failJob(new
HoodieException(errMsg, t)))
.build();
// start the executor if required
@@ -359,12 +366,14 @@ public class StreamWriteOperatorCoordinator
Pair<String, WriteMetadataEvent[]> instantTimeAndEventBuffer =
this.eventBuffers.getInstantAndEventBuffer(checkpointId);
final String instantTime;
if (instantTimeAndEventBuffer == null) {
+ // wait until previous instants are committed.
+ awaitAllInstantsToCompleteIfNecessary();
instantTime = startInstant();
this.eventBuffers.initNewEventBuffer(checkpointId, instantTime,
this.parallelism);
} else {
instantTime = instantTimeAndEventBuffer.getLeft();
}
-
response.complete(CoordinationResponseSeDe.wrap(Correspondent.InstantTimeResponse.getInstance(instantTime)));
+
response.complete(CoordinationResponseSerDe.wrap(Correspondent.InstantTimeResponse.getInstance(instantTime)));
}, "request instant time");
return response;
}
@@ -373,6 +382,12 @@ public class StreamWriteOperatorCoordinator
// Utilities
// -------------------------------------------------------------------------
+ private void awaitAllInstantsToCompleteIfNecessary() {
+ if (this.commitGuardOpt.isPresent() && this.eventBuffers.nonEmpty()) {
+
this.commitGuardOpt.get().blockFor(this.eventBuffers.getPendingInstants());
+ }
+ }
+
private ThreadFactory getThreadFactory(String threadName) {
return new ExplicitClassloaderThreadFactory(threadName,
context.getUserCodeClassloader());
}
@@ -426,6 +441,14 @@ public class StreamWriteOperatorCoordinator
this.clientIds.start();
}
+ private void initCommitGuard(Configuration conf) {
+ if (tableState.isBlockingInstantGeneration) {
+ this.commitGuardOpt =
Option.of(CommitGuard.create(conf.get(FlinkOptions.WRITE_COMMIT_ACK_TIMEOUT)));
+ } else {
+ this.commitGuardOpt = Option.empty();
+ }
+ }
+
private String startInstant() {
// refresh the meta client which is reused
metaClient.reloadActiveTimeline();
@@ -663,6 +686,10 @@ public class StreamWriteOperatorCoordinator
final boolean syncHive;
final boolean syncMetadata;
final boolean isDeltaTimeCompaction;
+ /**
+ * Whether the writer for the table applies blocking instant generation.
+ */
+ final boolean isBlockingInstantGeneration;
private TableState(Configuration conf) {
this.operationType =
WriteOperationType.fromValue(conf.getString(FlinkOptions.OPERATION));
@@ -674,6 +701,7 @@ public class StreamWriteOperatorCoordinator
this.syncHive = conf.getBoolean(FlinkOptions.HIVE_SYNC_ENABLED);
this.syncMetadata = conf.getBoolean(FlinkOptions.METADATA_ENABLED);
this.isDeltaTimeCompaction = OptionsResolver.isDeltaTimeCompaction(conf);
+ this.isBlockingInstantGeneration =
OptionsResolver.isBlockingInstantGeneration(conf);
}
public static TableState create(Configuration conf) {
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 d205e405534..e23e7e85127 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
@@ -20,7 +20,7 @@ package org.apache.hudi.sink.event;
import org.apache.hudi.common.util.VisibleForTesting;
import org.apache.hudi.exception.HoodieException;
-import org.apache.hudi.sink.utils.CoordinationResponseSeDe;
+import org.apache.hudi.sink.utils.CoordinationResponseSerDe;
import org.apache.flink.runtime.jobgraph.OperatorID;
import org.apache.flink.runtime.jobgraph.tasks.TaskOperatorEventGateway;
@@ -63,7 +63,7 @@ public class Correspondent {
*/
public String requestInstantTime(long checkpointId) {
try {
- InstantTimeResponse response =
CoordinationResponseSeDe.unwrap(this.gateway.sendRequestToCoordinator(this.operatorID,
+ InstantTimeResponse response =
CoordinationResponseSerDe.unwrap(this.gateway.sendRequestToCoordinator(this.operatorID,
new
SerializedValue<>(InstantTimeRequest.getInstance(checkpointId))).get());
return response.getInstant();
} catch (Exception e) {
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/CommitGuard.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/CommitGuard.java
new file mode 100644
index 00000000000..0be33a05dff
--- /dev/null
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/CommitGuard.java
@@ -0,0 +1,80 @@
+/*
+ * 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.apache.hudi.exception.HoodieException;
+
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.locks.Condition;
+import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantLock;
+
+/**
+ * The commit guard used for blocking instant time generation.
+ *
+ * <p>A new instant time can not be generated until all the pending instants
are committed.
+ */
+public class CommitGuard {
+ /**
+ * A lock used to block instant generation until being signaled after
pending instant is committed.
+ */
+ private final Lock lock;
+ private final Condition condition;
+ private final long commitAckTimeout;
+
+ private CommitGuard(long commitAckTimeout) {
+ this.lock = new ReentrantLock();
+ this.condition = lock.newCondition();
+ this.commitAckTimeout = commitAckTimeout;
+ }
+
+ public static CommitGuard create(long commitAckTimeout) {
+ return new CommitGuard(commitAckTimeout);
+ }
+
+ /**
+ * Wait until all the pending instants are committed.
+ *
+ * @param instants Current pending instants
+ */
+ public void blockFor(String instants) {
+ lock.lock();
+ try {
+ if (!condition.await(commitAckTimeout, TimeUnit.MILLISECONDS)) {
+ throw new HoodieException("Timeout(" + commitAckTimeout + "ms) while
waiting for instants [" + instants + "] to commit");
+ }
+ } catch (InterruptedException e) {
+ throw new HoodieException("Blocking for instants completion is
interrupted.", e);
+ } finally {
+ lock.unlock();
+ }
+ }
+
+ /**
+ * Signals all waited threads which are waiting on the condition.
+ */
+ public void unblock() {
+ lock.lock();
+ try {
+ condition.signalAll();
+ } finally {
+ lock.unlock();
+ }
+ }
+}
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/CoordinationResponseSeDe.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/CoordinationResponseSerDe.java
similarity index 99%
rename from
hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/CoordinationResponseSeDe.java
rename to
hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/CoordinationResponseSerDe.java
index 79b3723ae12..4d3cef43189 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/CoordinationResponseSeDe.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/CoordinationResponseSerDe.java
@@ -38,7 +38,7 @@ import java.util.List;
* Utilities for wrapping and unwrapping {@link CoordinationResponse}
* by {@link CollectCoordinationResponse}.
*/
-public class CoordinationResponseSeDe {
+public class CoordinationResponseSerDe {
private static final String MAGIC_VERSION = "__internal__";
private static final long MAGIC_OFFSET = 15213L;
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/EventBuffers.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/EventBuffers.java
index 5cf657d22f3..fa5595f9352 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/EventBuffers.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/EventBuffers.java
@@ -25,7 +25,9 @@ import org.apache.hudi.sink.event.WriteMetadataEvent;
import java.io.Serializable;
import java.util.Arrays;
import java.util.Map;
+import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.stream.Collectors;
import java.util.stream.Stream;
/**
@@ -36,13 +38,15 @@ public class EventBuffers implements Serializable {
// {checkpointId -> (instant, events)}
private final Map<Long, Pair<String, WriteMetadataEvent[]>> eventBuffers;
+ private final Option<CommitGuard> commitGuardOption;
- private EventBuffers(Map<Long, Pair<String, WriteMetadataEvent[]>>
eventBuffers) {
+ private EventBuffers(Map<Long, Pair<String, WriteMetadataEvent[]>>
eventBuffers, Option<CommitGuard> commitGuardOption) {
this.eventBuffers = eventBuffers;
+ this.commitGuardOption = commitGuardOption;
}
- public static EventBuffers getInstance() {
- return new EventBuffers(new ConcurrentHashMap<>());
+ public static EventBuffers getInstance(Option<CommitGuard>
commitGuardOption) {
+ return new EventBuffers(new ConcurrentHashMap<>(), commitGuardOption);
}
/**
@@ -116,5 +120,16 @@ public class EventBuffers implements Serializable {
public void reset(long checkpointId) {
this.eventBuffers.remove(checkpointId);
+ this.commitGuardOption.ifPresent(CommitGuard::unblock);
+ }
+
+ public boolean nonEmpty() {
+ return this.eventBuffers.values().stream()
+ .map(Pair::getValue)
+ .flatMap(Arrays::stream).anyMatch(Objects::nonNull);
+ }
+
+ public String getPendingInstants() {
+ return
this.eventBuffers.values().stream().map(Pair::getKey).collect(Collectors.joining(","));
}
}
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/FlinkWriteClients.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/FlinkWriteClients.java
index 70929c142d2..a1de831db31 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/FlinkWriteClients.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/FlinkWriteClients.java
@@ -244,9 +244,11 @@ public class FlinkWriteClients {
}
HoodieWriteConfig writeConfig = builder.build();
- // always LAZY for non-blocking instant time generation.
- writeConfig.setValue(HoodieCleanConfig.FAILED_WRITES_CLEANER_POLICY.key(),
- HoodieFailedWritesCleaningPolicy.LAZY.name());
+ if (!OptionsResolver.isBlockingInstantGeneration(conf)) {
+ // always LAZY for non-blocking instant time generation.
+
writeConfig.setValue(HoodieCleanConfig.FAILED_WRITES_CLEANER_POLICY.key(),
+ HoodieFailedWritesCleaningPolicy.LAZY.name());
+ }
if (loadFsViewStorageConfig &&
!conf.containsKey(FileSystemViewStorageConfig.REMOTE_HOST_NAME.key())) {
// do not use the builder to give a change for recovering the original
fs view storage config
FileSystemViewStorageConfig viewStorageConfig =
ViewStorageProperties.loadFromProperties(conf.getString(FlinkOptions.PATH),
conf);
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 3c2f385d14b..bf5c0553561 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
@@ -37,7 +37,7 @@ import org.apache.hudi.hadoop.fs.HadoopFSUtils;
import org.apache.hudi.metadata.HoodieTableMetadata;
import org.apache.hudi.sink.event.Correspondent;
import org.apache.hudi.sink.event.WriteMetadataEvent;
-import org.apache.hudi.sink.utils.CoordinationResponseSeDe;
+import org.apache.hudi.sink.utils.CoordinationResponseSerDe;
import org.apache.hudi.sink.utils.MockCoordinatorExecutor;
import org.apache.hudi.sink.utils.NonThrownExecutor;
import org.apache.hudi.storage.HoodieStorage;
@@ -433,16 +433,18 @@ public class TestStreamWriteOperatorCoordinator {
.createNewInstant(INSTANT_GENERATOR.createNewInstant(HoodieInstant.State.REQUESTED,
HoodieActiveTimeline.DELTA_COMMIT_ACTION, instant));
metadataTableMetaClient.getActiveTimeline().transitionRequestedToInflight(HoodieActiveTimeline.DELTA_COMMIT_ACTION,
instant);
metadataTableMetaClient.reloadActiveTimeline();
+ // reset the coordinator to mimic the job failover.
+ coordinator = createCoordinator(conf, 1);
- // write another commit with existing instant on the metadata timeline
+ // write another commit with new instant on the metadata timeline
instant = mockWriteWithMetadata(ckp);
metadataTableMetaClient.reloadActiveTimeline();
completedTimeline =
metadataTableMetaClient.getActiveTimeline().filterCompletedInstants();
assertThat("One instant need to sync to metadata table",
completedTimeline.countInstants(), is(metadataPartitions + 3));
- assertThat(completedTimeline.nthFromLastInstant(1).get().requestedTime(),
is(instant));
+ assertThat(completedTimeline.lastInstant().get().requestedTime(),
is(instant));
assertThat("The pending instant should be rolled back first",
- completedTimeline.lastInstant().get().getAction(),
is(HoodieTimeline.ROLLBACK_ACTION));
+ completedTimeline.nthFromLastInstant(1).get().getAction(),
is(HoodieTimeline.ROLLBACK_ACTION));
}
@Test
@@ -551,7 +553,7 @@ public class TestStreamWriteOperatorCoordinator {
private String requestInstantTime(StreamWriteOperatorCoordinator
coordinator, long checkpointId) {
try {
- Correspondent.InstantTimeResponse response =
CoordinationResponseSeDe.unwrap(coordinator.handleCoordinationRequest(Correspondent.InstantTimeRequest.getInstance(checkpointId)).get());
+ 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);
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestWriteCopyOnWrite.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestWriteCopyOnWrite.java
index 2eea0e357de..15fb8ab8c79 100644
---
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestWriteCopyOnWrite.java
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestWriteCopyOnWrite.java
@@ -29,6 +29,7 @@ import org.apache.hudi.config.HoodieCleanConfig;
import org.apache.hudi.config.HoodieWriteConfig;
import org.apache.hudi.configuration.FlinkOptions;
import org.apache.hudi.configuration.OptionsResolver;
+import org.apache.hudi.exception.HoodieException;
import org.apache.hudi.exception.HoodieWriteConflictException;
import org.apache.hudi.index.HoodieIndex;
import org.apache.hudi.sink.utils.TestWriteBase;
@@ -200,6 +201,82 @@ public class TestWriteCopyOnWrite extends TestWriteBase {
.end();
}
+ @Test
+ public void testNonBlockedInstantRequestAfterFailover() throws Exception {
+ conf.set(FlinkOptions.WRITE_BATCH_SIZE, BATCH_SIZE_MB);
+ conf.set(FlinkOptions.PRE_COMBINE, true);
+ conf.set(FlinkOptions.WRITE_COMMIT_ACK_TIMEOUT, 10_000L);
+ Map<String, String> expected = new HashMap<>();
+ expected.put("par1", "[id1,par1,id1,Danny,23,1,par1]");
+
+ preparePipeline()
+ // will eager flush
+ .consume(TestData.DATA_SET_INSERT_DUPLICATES)
+ .checkpoint(1)
+ .allDataFlushed()
+ .handleEvents(2)
+ .checkpointComplete(1)
+ .checkWrittenData(expected, 1)
+ // will eager flush
+ .consume(TestData.DATA_SET_INSERT_DUPLICATES)
+ .handleEvents(1)
+ // task failover, and send empty bootstrap event to coordinator
+ .subTaskFails(0, 1)
+ // handle the bootstrap event and reset buffer for subtask 0
+ .assertNextEvent()
+ // consume new data, will not be blocked
+ .consume(TestData.DATA_SET_INSERT)
+ .checkpoint(2)
+ .handleEvents(1)
+ .checkpointComplete(2)
+ .checkWrittenData(EXPECTED1, 4)
+ .end();
+ }
+
+ @Test
+ public void testBlockedInstantTimeRequest() throws Exception {
+ conf.set(FlinkOptions.WRITE_BATCH_SIZE, BATCH_SIZE_MB);
+ conf.set(FlinkOptions.PRE_COMBINE, true);
+ conf.set(FlinkOptions.WRITE_COMMIT_ACK_TIMEOUT, 10_000L);
+
+ Map<String, String> expected = new HashMap<>();
+ expected.put("par1", "[id1,par1,id1,Danny,23,1,par1]");
+
+ TestHarness testHarness = preparePipeline()
+ .consume(TestData.DATA_SET_INSERT_DUPLICATES)
+ .assertDataBuffer(1, 2)
+ .checkpoint(1)
+ .allDataFlushed()
+ .handleEvents(2);
+
+ Thread t1 = new Thread(() -> {
+ try {
+ Thread.sleep(3000);
+ testHarness.checkpointComplete(1);
+ testHarness.checkWrittenData(expected, 1);
+ } catch (Exception e) {
+ throw new HoodieException(e);
+ }
+ });
+ t1.start();
+
+ testHarness
+ // new records coming and flushing while cp1 is not completed yet,
+ // bucket assign function will upsert(U) new records to previous fg
stored in state.
+ // If async instant generation is used , HoodieMergedHandle will
either throw exception
+ // or get the wrong base file in the file group, since cp1 is not
committed yet.
+ .consume(TestData.DATA_SET_INSERT_DUPLICATES)
+ .consume(TestData.DATA_SET_INSERT)
+ .checkpoint(2)
+ .allDataFlushed();
+ t1.join();
+
+ testHarness.handleEvents(3)
+ .checkpointComplete(2)
+ .checkWrittenData(EXPECTED1, 4)
+ .end();
+ }
+
@Test
public void testInsert() throws Exception {
// open the function and ingest data
@@ -528,10 +605,11 @@ public class TestWriteCopyOnWrite extends TestWriteBase {
@Test
public void testWriteExactlyOnce() throws Exception {
// reset the config option
- conf.setLong(FlinkOptions.WRITE_COMMIT_ACK_TIMEOUT, 1L);
+ conf.setLong(FlinkOptions.WRITE_COMMIT_ACK_TIMEOUT, 1000L);
conf.set(FlinkOptions.WRITE_MEMORY_SEGMENT_PAGE_SIZE, 128);
conf.setDouble(FlinkOptions.WRITE_TASK_MAX_SIZE, 200.0006); // 630 bytes
buffer size
- preparePipeline(conf)
+ TestHarness pipeline = preparePipeline(conf)
+ .resetInstantTimeRequest(conf)
.consume(TestData.DATA_SET_INSERT)
.initialEventBuffer()
.checkpoint(1)
@@ -540,9 +618,18 @@ public class TestWriteCopyOnWrite extends TestWriteBase {
// requested instant with checkpoint id as 1
.consume(TestData.DATA_SET_INSERT)
.checkpoint(2)
- // requested instant with checkpoint id as 2
- .consume(TestData.DATA_SET_INSERT)
- .end();
+ .handleEvents(4);
+ // requested instant with checkpoint id as 2
+ if (OptionsResolver.isBlockingInstantGeneration(conf)) {
+ pipeline
+ .assertConsumeThrows(TestData.DATA_SET_INSERT,
+ "Timeout(1000ms) while waiting for instants")
+ .end();
+ } else {
+ pipeline
+ .consume(TestData.DATA_SET_INSERT)
+ .end();
+ }
}
// case1: txn2's time range is involved in txn1
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 aa2db0453c7..6669ed9f5ac 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
@@ -27,7 +27,6 @@ import org.apache.hudi.sink.event.Correspondent;
*/
public class MockCorrespondent extends Correspondent {
private final StreamWriteOperatorCoordinator coordinator;
-
public MockCorrespondent(StreamWriteOperatorCoordinator coordinator) {
this.coordinator = coordinator;
}
@@ -35,7 +34,7 @@ public class MockCorrespondent extends Correspondent {
@Override
public String requestInstantTime(long checkpointId) {
try {
- InstantTimeResponse response =
CoordinationResponseSeDe.unwrap(this.coordinator.handleCoordinationRequest(InstantTimeRequest.getInstance(checkpointId)).get());
+ 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);
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/MockCorrespondentWithTimeout.java
similarity index 61%
copy from
hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/MockCorrespondent.java
copy to
hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/MockCorrespondentWithTimeout.java
index aa2db0453c7..22b01a6432e 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/MockCorrespondentWithTimeout.java
@@ -18,24 +18,36 @@
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;
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.runtime.operators.coordination.CoordinationResponse;
+
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
+
/**
* A mock {@link Correspondent} that always return the latest instant.
+ *
+ * <p>A timeout is set up there to avoid the request hangs forever.
*/
-public class MockCorrespondent extends Correspondent {
+public class MockCorrespondentWithTimeout extends Correspondent {
private final StreamWriteOperatorCoordinator coordinator;
+ private final long commitAckTimeout;
- public MockCorrespondent(StreamWriteOperatorCoordinator coordinator) {
+ public MockCorrespondentWithTimeout(StreamWriteOperatorCoordinator
coordinator, Configuration conf) {
this.coordinator = coordinator;
+ this.commitAckTimeout = conf.get(FlinkOptions.WRITE_COMMIT_ACK_TIMEOUT);
}
@Override
public String requestInstantTime(long checkpointId) {
try {
- InstantTimeResponse response =
CoordinationResponseSeDe.unwrap(this.coordinator.handleCoordinationRequest(InstantTimeRequest.getInstance(checkpointId)).get());
+ 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);
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/TestWriteBase.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/TestWriteBase.java
index c65d7391e71..efa67a2ecf1 100644
---
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/TestWriteBase.java
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/utils/TestWriteBase.java
@@ -437,6 +437,16 @@ public class TestWriteBase {
return this;
}
+ /**
+ * Overrides the instant time request to avoid the write task hangs up.
+ */
+ public TestHarness resetInstantTimeRequest(Configuration conf) {
+ if (OptionsResolver.isBlockingInstantGeneration(conf)) {
+ this.pipeline.getWriteFunction().setCorrespondent(new
MockCorrespondentWithTimeout(this.pipeline.getCoordinator(), conf));
+ }
+ return this;
+ }
+
/**
* Mark the task with id {@code taskId} as failed.
*/