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 63d1190ea879 fix(flink): close all bulk insert handles and shut down
executor on failure (#20045)
63d1190ea879 is described below
commit 63d1190ea879e461930521144f3e85c6fc9b8123
Author: Shuo Cheng <[email protected]>
AuthorDate: Fri Sep 25 20:40:41 2026 +0800
fix(flink): close all bulk insert handles and shut down executor on failure
(#20045)
* fix(flink): close all bulk insert handles and shut down executor on
failure
---------
Co-authored-by: danny0405 <[email protected]>
---
.../hudi/sink/bulk/BulkInsertWriterHelper.java | 45 ++++---
.../hudi/sink/bulk/TestBulkInsertWriteHelper.java | 140 +++++++++++++++++++++
2 files changed, 162 insertions(+), 23 deletions(-)
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bulk/BulkInsertWriterHelper.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bulk/BulkInsertWriterHelper.java
index 4fa933d75db6..1fbb311fd736 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bulk/BulkInsertWriterHelper.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bulk/BulkInsertWriterHelper.java
@@ -51,11 +51,8 @@ import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
-import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
-import static org.apache.hudi.common.util.FutureUtils.allOf;
-
/**
* Helper class for bulk insert used by Flink.
*/
@@ -186,29 +183,31 @@ public class BulkInsertWriterHelper implements
AutoCloseable {
if (handles.isEmpty()) {
return;
}
- int handsSize = Math.min(handles.size(), 10);
- ExecutorService executorService = Executors.newFixedThreadPool(handsSize);
- allOf(handles.values().stream()
- .map(rowCreateHandle -> CompletableFuture.supplyAsync(() -> {
- try {
- log.info("Closing bulk insert file {}",
rowCreateHandle.getFileName());
- return rowCreateHandle.close();
- } catch (IOException e) {
- throw new HoodieIOException("IOE during rowCreateHandle.close()",
e);
- }
- }, executorService))
- .collect(Collectors.toList())
- ).whenComplete((result, throwable) -> {
- writeStatusList.addAll(result);
- }).join();
+ ExecutorService executorService =
Executors.newFixedThreadPool(Math.min(handles.size(), 10));
try {
+ List<CompletableFuture<WriteStatus>> futures = handles.values().stream()
+ .map(rowCreateHandle -> closeAsync(rowCreateHandle, executorService))
+ .collect(Collectors.toList());
+ // Use JDK allOf instead of FutureUtils.allOf, which cancels unfinished
futures on failure.
+ // Cancelling queued close tasks can leave handles unclosed and leak
resources; wait for every close attempt instead.
+ CompletableFuture.allOf(futures.toArray(new
CompletableFuture<?>[0])).join();
+ futures.forEach(future -> writeStatusList.add(future.join()));
+ handles.clear();
+ handle = null;
+ } finally {
executorService.shutdown();
- executorService.awaitTermination(10, TimeUnit.MINUTES);
- } catch (InterruptedException e) {
- throw new RuntimeException(e);
}
- handles.clear();
- handle = null;
+ }
+
+ private CompletableFuture<WriteStatus> closeAsync(HoodieRowDataCreateHandle
rowCreateHandle, ExecutorService executorService) {
+ return CompletableFuture.supplyAsync(() -> {
+ try {
+ log.info("Closing bulk insert file {}", rowCreateHandle.getFileName());
+ return rowCreateHandle.close();
+ } catch (IOException e) {
+ throw new HoodieIOException("IOE during rowCreateHandle.close()", e);
+ }
+ }, executorService);
}
private String getNextFileId() {
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/bulk/TestBulkInsertWriteHelper.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/bulk/TestBulkInsertWriteHelper.java
index da26b9d88b4c..14013df2fc31 100644
---
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/bulk/TestBulkInsertWriteHelper.java
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/bulk/TestBulkInsertWriteHelper.java
@@ -27,6 +27,7 @@ import org.apache.hudi.common.table.HoodieTableConfig;
import org.apache.hudi.common.util.ParquetUtils;
import org.apache.hudi.common.util.StringUtils;
import org.apache.hudi.configuration.FlinkOptions;
+import org.apache.hudi.io.storage.row.HoodieRowDataCreateHandle;
import org.apache.hudi.storage.StoragePath;
import org.apache.hudi.table.HoodieFlinkTable;
import org.apache.hudi.util.DataTypeUtils;
@@ -46,7 +47,9 @@ import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.CsvSource;
import org.junit.jupiter.params.provider.MethodSource;
+import org.mockito.MockedStatic;
import java.io.File;
import java.io.IOException;
@@ -56,6 +59,13 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.CompletionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.FutureTask;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
import java.util.stream.Stream;
@@ -64,6 +74,14 @@ import static org.hamcrest.MatcherAssert.assertThat;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
/**
* Test cases for {@link BulkInsertWriterHelper}.
@@ -200,6 +218,128 @@ public class TestBulkInsertWriteHelper {
assertThrows(IOException.class, () -> writerHelper.write(new
GenericRowData(0)));
}
+ @ParameterizedTest
+ @CsvSource({"false, false", "true, false", "false, true", "true, true"})
+ void testCloseFailureDoesNotSkipQueuedHandles(boolean runtimeFailure,
boolean multipleFailures) throws Exception {
+ BulkInsertWriterHelper helper = newWriterHelper();
+ List<HoodieRowDataCreateHandle> handles = new ArrayList<>();
+ for (int i = 0; i < 11; i++) {
+ HoodieRowDataCreateHandle handle = mock(HoodieRowDataCreateHandle.class);
+ when(handle.close()).thenReturn(new WriteStatus());
+ helper.handles.put("partition-" + i, handle);
+ handles.add(handle);
+ }
+ HoodieRowDataCreateHandle failedHandle =
helper.handles.values().iterator().next();
+ Exception failure = runtimeFailure ? new IllegalStateException("close
failed") : new IOException("close failed");
+ when(failedHandle.close()).thenThrow(failure);
+ List<Exception> failures = new ArrayList<>();
+ failures.add(failure);
+ if (multipleFailures) {
+ HoodieRowDataCreateHandle secondFailedHandle =
helper.handles.values().stream().skip(1).findFirst().get();
+ Exception secondFailure = runtimeFailure ? new IOException("second close
failed") : new IllegalStateException("second close failed");
+ when(secondFailedHandle.close()).thenThrow(secondFailure);
+ failures.add(secondFailure);
+ }
+
+ // Hold all close tasks in a queue until the helper is waiting for their
results.
+ // This makes cancellation of not-yet-started closes deterministic.
+ ExecutorService executor = mock(ExecutorService.class);
+ BlockingQueue<Runnable> tasks = new LinkedBlockingQueue<>();
+ doAnswer(invocation -> {
+ tasks.add(invocation.getArgument(0));
+ return null;
+ }).when(executor).execute(any(Runnable.class));
+ FutureTask<CompletionException> closeTask = new FutureTask<>(() -> {
+ try (MockedStatic<Executors> executors = mockStatic(Executors.class)) {
+ executors.when(() ->
Executors.newFixedThreadPool(10)).thenReturn(executor);
+ return assertThrows(CompletionException.class, helper::close);
+ }
+ });
+ Thread closingThread = new Thread(closeTask);
+ closingThread.setDaemon(true);
+ closingThread.start();
+ List<Runnable> submitted = new ArrayList<>();
+ try {
+ for (int i = 0; i < 11; i++) {
+ Runnable task = tasks.poll(10, TimeUnit.SECONDS);
+ assertTrue(task != null, "Every handle must be submitted for closing");
+ submitted.add(task);
+ }
+ long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
+ while (closingThread.getState() != Thread.State.WAITING &&
System.nanoTime() < deadline) {
+ Thread.yield();
+ }
+ assertEquals(Thread.State.WAITING, closingThread.getState());
+ submitted.get(0).run();
+ assertFalse(closeTask.isDone(), "A failed close must not complete the
helper before the remaining closes");
+ submitted.subList(1, submitted.size()).forEach(Runnable::run);
+ CompletionException thrown = closeTask.get(10, TimeUnit.SECONDS);
+ Throwable cause = thrown;
+ while (cause.getCause() != null) {
+ cause = cause.getCause();
+ }
+ // allOf preserves a close failure, but does not guarantee which one
when several handles fail.
+ assertTrue(failures.contains(cause));
+ assertEquals(0, thrown.getSuppressed().length);
+ verify(executor).shutdown();
+ for (HoodieRowDataCreateHandle handle : handles) {
+ verify(handle).close();
+ }
+ } finally {
+ // join() is uninterruptible: finish queued tasks even if an assertion
failed before they ran.
+ submitted.forEach(Runnable::run);
+ long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
+ while (closingThread.isAlive() && System.nanoTime() < deadline) {
+ Runnable task = tasks.poll(100, TimeUnit.MILLISECONDS);
+ if (task != null) {
+ task.run();
+ }
+ }
+ closingThread.join(1000);
+ assertFalse(closingThread.isAlive(), "The closing thread must terminate
after cleanup");
+ }
+ }
+
+ @Test
+ void testFailedCloseRetainsHandlesForCleanupRetry() throws Exception {
+ BulkInsertWriterHelper helper = newWriterHelper();
+ HoodieRowDataCreateHandle successfulHandle =
mock(HoodieRowDataCreateHandle.class);
+ HoodieRowDataCreateHandle failedHandle =
mock(HoodieRowDataCreateHandle.class);
+ when(successfulHandle.close()).thenReturn(new WriteStatus());
+ when(failedHandle.close()).thenThrow(new IOException("close failed"));
+ helper.handles.put("par1", successfulHandle);
+ helper.handles.put("par2", failedHandle);
+
+ assertThrows(CompletionException.class, helper::close);
+ assertEquals(2, helper.handles.size());
+ // A repeated status request must still fail, rather than returning an
incomplete set of statuses.
+ assertThrows(CompletionException.class, () -> helper.getWriteStatuses(1));
+ verify(successfulHandle, times(2)).close();
+ verify(failedHandle, times(2)).close();
+ }
+
+ @Test
+ void testSuccessfulCloseIsIdempotent() throws Exception {
+ BulkInsertWriterHelper helper = newWriterHelper();
+ HoodieRowDataCreateHandle handle = mock(HoodieRowDataCreateHandle.class);
+ WriteStatus status = new WriteStatus();
+ when(handle.close()).thenReturn(status);
+ helper.handles.put("par1", handle);
+
+ helper.close();
+ helper.close();
+
+ verify(handle).close();
+ assertTrue(helper.handles.isEmpty());
+ assertEquals(Arrays.asList(status), helper.getWriteStatuses(1));
+ }
+
+ private BulkInsertWriterHelper newWriterHelper() {
+ HoodieFlinkTable<?> table = FlinkTables.createTable(conf);
+ return new BulkInsertWriterHelper(conf, table, table.getConfig(),
WriteClientTestUtils.createNewInstantTime(),
+ 1, 1, 0, TestConfigurations.ROW_TYPE);
+ }
+
private void assertWriteStatus(List<WriteStatus> writeStatusList) {
String partitions = writeStatusList.stream()
.map(writeStatus ->
StringUtils.nullToEmpty(writeStatus.getStat().getPartitionPath()))