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

Reply via email to