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";

Reply via email to