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 23a62e4fe6bc fix(common): close native log writers independently on
failure (#20118)
23a62e4fe6bc is described below
commit 23a62e4fe6bc688425c0b5fb2da5c951c7038971
Author: Shuo Cheng <[email protected]>
AuthorDate: Tue Sep 29 17:59:02 2026 +0800
fix(common): close native log writers independently on failure (#20118)
---
.../hudi/io/cdc/HoodieNativeLogFormatWriter.java | 29 +++++-
.../io/cdc/TestHoodieNativeLogFormatWriter.java | 107 +++++++++++++++++++++
2 files changed, 131 insertions(+), 5 deletions(-)
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/cdc/HoodieNativeLogFormatWriter.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/cdc/HoodieNativeLogFormatWriter.java
index 274d4902f25f..f1ed4a9e821e 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/cdc/HoodieNativeLogFormatWriter.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/cdc/HoodieNativeLogFormatWriter.java
@@ -34,6 +34,7 @@ import org.apache.hudi.common.table.log.LogReaderUtils;
import org.apache.hudi.common.table.log.NativeLogFooterMetadata;
import org.apache.hudi.common.table.log.block.HoodieLogBlock;
import
org.apache.hudi.common.table.log.block.HoodieLogBlock.HeaderMetadataType;
+import org.apache.hudi.common.util.CloseableUtils;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.OrderingValues;
import org.apache.hudi.common.util.collection.ArrayComparable;
@@ -280,6 +281,20 @@ public class HoodieNativeLogFormatWriter extends
HoodieLogFormat.Writer {
}
private void closeFileWriters() throws IOException {
+ try {
+ closeDataFileWriter();
+ } catch (Throwable e) {
+ // A failure closing the data writer or reading its metadata must not
skip closing the delete writer.
+ CloseableUtils.closeSuppressing(this::closeDeleteFileWriter, e);
+ throw e;
+ } finally {
+ dataFileWriter = null;
+ dataRecordPositions.clear();
+ }
+ closeDeleteFileWriter();
+ }
+
+ private void closeDataFileWriter() throws IOException {
if (dataFileWriter != null) {
dataFileWriter.close();
if (writeConfig.isMetadataColumnStatsIndexEnabled()) {
@@ -293,14 +308,18 @@ public class HoodieNativeLogFormatWriter extends
HoodieLogFormat.Writer {
} else {
lastDataFileFormatMetadata = Option.empty();
}
- dataFileWriter = null;
}
- if (deleteFileWriter != null) {
- deleteFileWriter.close();
+ }
+
+ private void closeDeleteFileWriter() throws IOException {
+ try {
+ if (deleteFileWriter != null) {
+ deleteFileWriter.close();
+ }
+ } finally {
deleteFileWriter = null;
+ deleteRecordPositions.clear();
}
- dataRecordPositions.clear();
- deleteRecordPositions.clear();
}
private HoodieLogFile createNativeLogFile(int version, String logExtension)
throws IOException {
diff --git
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeLogFormatWriter.java
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeLogFormatWriter.java
index 182c1e2ee62e..d343ca4c077f 100644
---
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeLogFormatWriter.java
+++
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/cdc/TestHoodieNativeLogFormatWriter.java
@@ -39,6 +39,7 @@ import org.apache.hudi.common.util.OrderingValues;
import org.apache.hudi.config.HoodieWriteConfig;
import org.apache.hudi.core.io.storage.HoodieFileWriter;
import org.apache.hudi.core.io.storage.HoodieFileWriterFactory;
+import org.apache.hudi.exception.HoodieIOException;
import org.apache.hudi.storage.HoodieStorage;
import org.apache.hudi.storage.StoragePath;
import org.apache.hudi.storage.StoragePathInfo;
@@ -48,10 +49,12 @@ import org.apache.avro.generic.GenericRecord;
import org.apache.avro.generic.IndexedRecord;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.CsvSource;
import org.junit.jupiter.params.provider.ValueSource;
import org.mockito.ArgumentCaptor;
import org.mockito.MockedStatic;
+import java.io.IOException;
import java.lang.reflect.Field;
import java.util.ArrayList;
import java.util.Arrays;
@@ -61,10 +64,14 @@ import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
+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.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.mockStatic;
import static org.mockito.Mockito.never;
@@ -218,6 +225,106 @@ public class TestHoodieNativeLogFormatWriter {
assertFalse(metadata.isPresent());
}
+ @ParameterizedTest
+ @CsvSource({
+ "io, false, false", "io, false, true",
+ "io, true, false", "io, true, true",
+ "runtime, true, false", "runtime, true, true",
+ "metadata, false, false", "metadata, false, true",
+ "metadata, true, false", "metadata, true, true",
+ "error, true, false", "error, true, true",
+ "none, true, false", "none, true, true",
+ "shared, true, false", "shared, true, true"
+ })
+ public void testClosesBothWritersAndClearsStateOnFailure(
+ String failureMode, boolean failDelete, boolean flushAppend) throws
Exception {
+ HoodieSchema schema =
HoodieSchema.parse("{\"type\":\"record\",\"name\":\"test\",\"fields\":[]}");
+ HoodieWriteConfig config = mock(HoodieWriteConfig.class);
+ HoodieRecordMerger merger = mock(HoodieRecordMerger.class);
+ HoodieStorage storage = mock(HoodieStorage.class);
+ HoodieFileWriter dataWriter = mock(HoodieFileWriter.class);
+ HoodieFileWriter deleteWriter = mock(HoodieFileWriter.class);
+ HoodieFileWriter nextDataWriter = mock(HoodieFileWriter.class);
+ HoodieFileWriter nextDeleteWriter = mock(HoodieFileWriter.class);
+ when(config.getProps()).thenReturn(new TypedProperties());
+ when(config.getRecordMerger()).thenReturn(merger);
+ when(config.isMetadataColumnStatsIndexEnabled()).thenReturn(true);
+
when(merger.getRecordType()).thenReturn(HoodieRecord.HoodieRecordType.AVRO);
+ when(storage.getPathInfo(any(StoragePath.class))).thenAnswer(invocation ->
+ new StoragePathInfo(invocation.getArgument(0), 1L, false, (short) 1,
1L, 1L));
+
+ Throwable dataFailure = null;
+ if ("metadata".equals(failureMode)) {
+ dataFailure = new IllegalStateException("metadata failed");
+ when(dataWriter.getFileFormatMetadata()).thenThrow(dataFailure);
+ } else if (!"none".equals(failureMode)) {
+ dataFailure = "runtime".equals(failureMode) ? new
IllegalStateException("data close failed")
+ : "error".equals(failureMode) ? new AssertionError("data close
failed") : new IOException("data close failed");
+ doThrow(dataFailure).when(dataWriter).close();
+ }
+ Throwable deleteFailure = "shared".equals(failureMode) ? dataFailure : new
IOException("delete close failed");
+
+ try (MockedStatic<HoodieFileWriterFactory> writerFactory =
mockStatic(HoodieFileWriterFactory.class)) {
+ writerFactory.when(() -> HoodieFileWriterFactory.getFileWriter(
+ eq("100"), any(StoragePath.class), eq(storage), eq(config),
any(HoodieSchema.class),
+ any(TaskContextSupplier.class),
eq(HoodieRecord.HoodieRecordType.AVRO)))
+ .thenReturn(dataWriter, deleteWriter, nextDataWriter,
nextDeleteWriter);
+ HoodieNativeLogFormatWriter writer = new HoodieNativeLogFormatWriter(
+ 4096, storage, new StoragePath("/tmp/partition"), "file-1", "100",
1, "1-0-1", 1024L,
+ new LogFileCreationCallback() {
+ }, HoodieTableVersion.current(), config, HoodieFileFormat.PARQUET,
schema,
+ mock(TaskContextSupplier.class), new AvroRecordContext(), new
ArrayList<>(), Option.of("001"));
+ doAnswer(invocation -> {
+ assertTrue(writer.hasPendingWrites(), "The delete writer must remain
referenced until close is attempted");
+ if (failDelete) {
+ throw deleteFailure;
+ }
+ return null;
+ }).when(deleteWriter).close();
+ writer.appendRecord(recordWithPosition("key-1", 1L, schema), schema);
+ writer.appendDeleteRecord(recordWithPosition("key-2", 2L, schema),
schema);
+
+ Throwable expected = dataFailure != null ? dataFailure : deleteFailure;
+ Throwable thrown = assertThrows(Throwable.class, () -> {
+ if (flushAppend) {
+ writer.flushAppend(new HashMap<>());
+ } else {
+ writer.close();
+ }
+ });
+ if (!flushAppend && expected instanceof IOException) {
+ assertEquals(HoodieIOException.class, thrown.getClass());
+ thrown = thrown.getCause();
+ }
+ assertSame(expected, thrown);
+ if (dataFailure != null && failDelete && dataFailure != deleteFailure) {
+ assertEquals(1, thrown.getSuppressed().length);
+ assertSame(deleteFailure, thrown.getSuppressed()[0]);
+ } else {
+ assertEquals(0, thrown.getSuppressed().length);
+ }
+ assertFalse(writer.hasPendingWrites());
+ assertFalse(writer.getLastDataFileFormatMetadata().isPresent());
+ writer.close();
+ verify(dataWriter).close();
+ verify(deleteWriter).close();
+
+ // A subsequent append must not inherit positions from the failed close.
+ writer.appendRecord(recordWithPosition("key-3", 3L, schema), schema);
+ writer.appendDeleteRecord(recordWithPosition("key-4", 4L, schema),
schema);
+ Map<HeaderMetadataType, String> header = new HashMap<>();
+
header.put(HeaderMetadataType.BASE_FILE_INSTANT_TIME_OF_RECORD_POSITIONS,
"001");
+ writer.flushAppend(header);
+ for (HoodieFileWriter fileWriter : Arrays.asList(nextDataWriter,
nextDeleteWriter)) {
+ ArgumentCaptor<Map<String, String>> footerCaptor =
ArgumentCaptor.forClass(Map.class);
+ verify(fileWriter).addFooterMetadata(footerCaptor.capture());
+ Map<HeaderMetadataType, String> parsedHeader =
NativeLogFooterMetadata.fromFooterMetadata(footerCaptor.getValue());
+ assertEquals(Arrays.asList(fileWriter == nextDataWriter ? 3L : 4L),
+
LogReaderUtils.decodeRecordPositionsLongList(parsedHeader.get(HeaderMetadataType.RECORD_POSITIONS)));
+ }
+ }
+ }
+
private static Option<Object> writeDataLogAndGetFormatMetadata(
boolean columnStatsEnabled, Object formatMetadata, Throwable
metadataFailure) throws Exception {
String instantTime = "100";