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 b83445fe6820 fix: close the file writers properly in fail cases 
(#18776)
b83445fe6820 is described below

commit b83445fe6820478b277d596af978c8ab549c3a95
Author: Danny Chan <[email protected]>
AuthorDate: Tue Sep 22 09:50:53 2026 +0800

    fix: close the file writers properly in fail cases (#18776)
    
    * fix: close the file writers properly in fail cases
---
 .../java/org/apache/hudi/io/BaseCreateHandle.java  |   8 ++
 .../hudi/io/FileGroupReaderBasedMergeHandle.java   |   3 +-
 .../org/apache/hudi/io/HoodieAppendHandle.java     |  29 ++++-
 .../org/apache/hudi/io/HoodieBinaryCopyHandle.java |  14 +++
 .../hudi/io/HoodieInlineLogAppendHandle.java       |   7 +-
 .../hudi/io/HoodieMergeHandleWithChangeLog.java    |   9 +-
 .../hudi/io/HoodieNativeLogAppendHandle.java       |  13 ++-
 .../apache/hudi/io/HoodieSortedMergeHandle.java    |  16 +--
 .../org/apache/hudi/io/HoodieWriteMergeHandle.java |  36 +++++--
 .../org/apache/hudi/io/TestHoodieAppendHandle.java |  45 +++++++-
 .../apache/hudi/io/TestHoodieBinaryCopyHandle.java |  74 +++++++++++++
 .../org/apache/hudi/io/TestHoodieCreateHandle.java |  52 ++++++---
 .../hudi/io/TestHoodieNativeLogAppendHandle.java   |  62 +++++++++++
 .../io/TestSortedAndChangeLogMergeHandles.java     | 107 +++++++++++++++++++
 .../HoodieRowDataParquetOutputStreamWriter.java    |   9 +-
 .../io/storage/HoodieSparkParquetStreamWriter.java |  27 +++--
 .../index/hfile/HFileBootstrapIndexWriter.java     |  58 +++++++---
 .../apache/hudi/common/util/CloseableUtils.java    |   9 +-
 .../hudi/common/util/TestCloseableUtils.java       |  46 ++++++--
 .../org/apache/hudi/common/util/ParquetUtils.java  |  13 +--
 .../io/storage/hadoop/HoodieAvroHFileWriter.java   |  17 ++-
 .../storage/hadoop/HoodieParquetStreamWriter.java  |  27 +++--
 .../parquet/io/HoodieParquetBinaryCopyBase.java    | 116 ++++++++++++--------
 .../parquet/io/HoodieParquetFileBinaryCopier.java  |  27 +++--
 .../io/hadoop/TestHoodieHFileReaderWriter.java     |  34 ++++++
 ...HoodieParquetBinaryCopyBaseSchemaEvolution.java |   2 +-
 .../io/TestHoodieParquetBinaryCopyLifecycle.java   | 118 +++++++++++++++++++++
 .../TestHoodieParquetFileBinaryCopierPrefetch.java |  21 ++++
 .../org/apache/spark/sql/hudi/SparkHelpers.scala   |  18 ++--
 29 files changed, 861 insertions(+), 156 deletions(-)

diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/BaseCreateHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/BaseCreateHandle.java
index 6eaba960606b..1ce8807e425b 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/BaseCreateHandle.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/BaseCreateHandle.java
@@ -30,6 +30,7 @@ import org.apache.hudi.common.model.IOType;
 import org.apache.hudi.common.model.MetaFieldsMode;
 import org.apache.hudi.common.model.MetadataValues;
 import org.apache.hudi.common.schema.HoodieSchema;
+import org.apache.hudi.common.util.CloseableUtils;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.StringUtils;
 import org.apache.hudi.config.HoodieWriteConfig;
@@ -128,6 +129,7 @@ public abstract class BaseCreateHandle<T, I, K, O> extends 
HoodieWriteHandle<T,
     } catch (Throwable t) {
       log.error("Error writing record {}", record, t);
       if (!config.getIgnoreWriteFailed()) {
+        closeFileWriterQuietly(t);
         throw new HoodieException(t.getMessage(), t);
       }
       writeStatus.markFailure(record, t, recordMetadata);
@@ -194,6 +196,12 @@ public abstract class BaseCreateHandle<T, I, K, O> extends 
HoodieWriteHandle<T,
     return record.prependMetaFields(schema, targetSchema, metadataValues, 
prop);
   }
 
+  private void closeFileWriterQuietly(Throwable failure) {
+    markClosed();
+    CloseableUtils.closeSuppressing(fileWriter, failure);
+    fileWriter = null;
+  }
+
   @Override
   public IOType getIOType() {
     return IOType.CREATE;
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
index 84fe098ff4cc..c7abeda7eeed 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
@@ -332,7 +332,8 @@ public class FileGroupReaderBasedMergeHandle<T, I, K, O> 
extends HoodieWriteMerg
         this.updatedRecordsWritten = readStats.getNumUpdates();
         this.recordsDeleted = readStats.getNumDeletes();
       }
-    } catch (IOException e) {
+    } catch (Exception e) {
+      closeFileWriterQuietly(e);
       throw new HoodieUpsertException("Failed to compact file group: " + 
fileId, e);
     }
   }
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
index 53bec47dca55..8bd2eff5e1ec 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieAppendHandle.java
@@ -40,6 +40,7 @@ import org.apache.hudi.common.table.log.block.HoodieLogBlock;
 import 
org.apache.hudi.common.table.log.block.HoodieLogBlock.HeaderMetadataType;
 import org.apache.hudi.common.table.timeline.HoodieInstantTimeGenerator;
 import org.apache.hudi.common.table.view.TableFileSystemView;
+import org.apache.hudi.common.util.CloseableUtils;
 import org.apache.hudi.common.util.HoodieRecordUtils;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.ReflectionUtils;
@@ -151,11 +152,16 @@ public abstract class HoodieAppendHandle<T, I, K, O> 
extends HoodieWriteHandle<T
    * Writes all records from the incoming iterator and flushes remaining 
pending data.
    */
   public void doAppend() {
-    while (recordItr.hasNext()) {
-      HoodieRecord<T> record = recordItr.next();
-      writeRecord(record);
+    try {
+      while (recordItr.hasNext()) {
+        HoodieRecord<T> record = recordItr.next();
+        writeRecord(record);
+      }
+      flushAppend();
+    } catch (RuntimeException e) {
+      closeLogWriterQuietly(e);
+      throw e;
     }
-    flushAppend();
   }
 
   /**
@@ -178,6 +184,7 @@ public abstract class HoodieAppendHandle<T, I, K, O> 
extends HoodieWriteHandle<T
       }
       flushAppend();
     } catch (Exception e) {
+      closeLogWriterQuietly(e);
       throw new HoodieUpsertException("Failed to compact blocks for fileId " + 
fileId, e);
     }
   }
@@ -214,8 +221,21 @@ public abstract class HoodieAppendHandle<T, I, K, O> 
extends HoodieWriteHandle<T
 
       return statuses;
     } catch (IOException e) {
+      closeLogWriterQuietly(e);
       throw new HoodieUpsertException("Failed to close " + 
getClass().getSimpleName(), e);
+    } catch (RuntimeException e) {
+      closeLogWriterQuietly(e);
+      throw e;
+    }
+  }
+
+  private void closeLogWriterQuietly(Throwable failure) {
+    markClosed();
+    CloseableUtils.closeSuppressing(this::closeLogWriter, failure);
+    if (recordItr instanceof Closeable) {
+      CloseableUtils.closeSuppressing((Closeable) recordItr, failure);
     }
+    recordItr = null;
   }
 
   @Override
@@ -382,6 +402,7 @@ public abstract class HoodieAppendHandle<T, I, K, O> 
extends HoodieWriteHandle<T
     } catch (Exception e) {
       log.error("Error writing record {}", hoodieRecord, e);
       if (!config.getIgnoreWriteFailed() || ExceptionUtil.isCausedBy(e, 
HoodieEarlyConflictDetectionException.class)) {
+        closeLogWriterQuietly(e);
         throw new HoodieException(e.getMessage(), e);
       }
       writeStatus.markFailure(hoodieRecord, e, recordMetadata);
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieBinaryCopyHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieBinaryCopyHandle.java
index e79faaec9291..07d0934b2fe0 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieBinaryCopyHandle.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieBinaryCopyHandle.java
@@ -23,6 +23,7 @@ import org.apache.hudi.common.engine.TaskContextSupplier;
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.HoodieWriteStat;
 import org.apache.hudi.common.model.IOType;
+import org.apache.hudi.common.util.CloseableUtils;
 import org.apache.hudi.common.util.HoodieTimer;
 import org.apache.hudi.common.util.ParquetUtils;
 import org.apache.hudi.config.HoodieWriteConfig;
@@ -120,7 +121,11 @@ public class HoodieBinaryCopyHandle<T, I, K, O> extends 
HoodieWriteHandle<T, I,
       log.info("Schema evolution enabled for binary copy: {}", 
schemaEvolutionEnabled);
       records = this.writer.binaryCopy(inputFiles, 
Collections.singletonList(path), writeScheMessageType, schemaEvolutionEnabled);
     } catch (IOException e) {
+      closeWriterQuietly(e);
       throw new HoodieIOException(e.getMessage(), e);
+    } catch (RuntimeException e) {
+      closeWriterQuietly(e);
+      throw e;
     } finally {
       this.recordsWritten = records;
       this.insertRecordsWritten = records;
@@ -128,10 +133,19 @@ public class HoodieBinaryCopyHandle<T, I, K, O> extends 
HoodieWriteHandle<T, I,
     log.info("Finish rewriting {}. Using {} mills", this.path, 
timer.endTimer());
   }
 
+  private void closeWriterQuietly(Throwable failure) {
+    markClosed();
+    CloseableUtils.closeSuppressing(writer::close, failure);
+  }
+
   @Override
   public List<WriteStatus> close() {
     log.info("Closing the file {} as we are done with all the records {}", 
writeStatus.getFileId(), recordsWritten);
     try {
+      if (isClosed()) {
+        return Collections.singletonList(writeStatus);
+      }
+      markClosed();
       this.writer.close();
 
       HoodieWriteStat stat = writeStatus.getStat();
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieInlineLogAppendHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieInlineLogAppendHandle.java
index 59709e74b1b0..e0a2d184425d 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieInlineLogAppendHandle.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieInlineLogAppendHandle.java
@@ -252,8 +252,11 @@ public class HoodieInlineLogAppendHandle<T, I, K, O> 
extends HoodieAppendHandle<
 
   @Override
   protected void closeLogWriter() throws IOException {
-    if (writer != null) {
-      writer.close();
+    try {
+      if (writer != null) {
+        writer.close();
+      }
+    } finally {
       writer = null;
     }
   }
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleWithChangeLog.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleWithChangeLog.java
index e635294dbba0..3dd86c8699fb 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleWithChangeLog.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieMergeHandleWithChangeLog.java
@@ -27,6 +27,7 @@ import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.HoodieWriteStat;
 import org.apache.hudi.common.schema.HoodieSchema;
 import org.apache.hudi.common.table.cdc.HoodieCDCUtils;
+import org.apache.hudi.common.util.CloseableUtils;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.io.cdc.HoodieCDCLogWriter;
@@ -112,7 +113,13 @@ public class HoodieMergeHandleWithChangeLog<T, I, K, O> 
extends HoodieWriteMerge
 
   @Override
   public List<WriteStatus> close() {
-    List<WriteStatus> writeStatuses = super.close();
+    List<WriteStatus> writeStatuses;
+    try {
+      writeStatuses = super.close();
+    } catch (RuntimeException e) {
+      CloseableUtils.closeSuppressing(cdcLogger, e);
+      throw e;
+    }
     // if there are cdc data written, set the CDC-related information.
 
     if (cdcLogger == null || recordsWritten == 0L || (recordsWritten == 
insertRecordsWritten)) {
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogAppendHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogAppendHandle.java
index 658a24d610e8..aae1a6b476dc 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogAppendHandle.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogAppendHandle.java
@@ -30,6 +30,7 @@ import org.apache.hudi.common.util.HoodieRecordUtils;
 import org.apache.hudi.common.util.Lazy;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.ParquetUtils;
+import org.apache.hudi.common.util.VisibleForTesting;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.exception.HoodieAppendException;
 import org.apache.hudi.exception.HoodieException;
@@ -62,6 +63,11 @@ public class HoodieNativeLogAppendHandle<T, I, K, O> extends 
HoodieAppendHandle<
 
   private HoodieNativeLogFormatWriter writer;
 
+  @VisibleForTesting
+  HoodieNativeLogFormatWriter getWriter() {
+    return writer;
+  }
+
   public HoodieNativeLogAppendHandle(HoodieWriteConfig config, String 
instantTime, HoodieTable<T, I, K, O> hoodieTable,
                                      String partitionPath, String fileId, 
Iterator<HoodieRecord<T>> recordItr,
                                      TaskContextSupplier taskContextSupplier) {
@@ -163,8 +169,11 @@ public class HoodieNativeLogAppendHandle<T, I, K, O> 
extends HoodieAppendHandle<
 
   @Override
   protected void closeLogWriter() {
-    if (writer != null) {
-      writer.close();
+    try {
+      if (writer != null) {
+        writer.close();
+      }
+    } finally {
       writer = null;
     }
   }
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieSortedMergeHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieSortedMergeHandle.java
index 458567b1fd91..dabb71ebdd7a 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieSortedMergeHandle.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieSortedMergeHandle.java
@@ -18,7 +18,6 @@
 
 package org.apache.hudi.io;
 
-import org.apache.hudi.client.WriteStatus;
 import org.apache.hudi.common.engine.TaskContextSupplier;
 import org.apache.hudi.common.model.HoodieBaseFile;
 import org.apache.hudi.common.model.HoodieRecord;
@@ -33,7 +32,6 @@ import org.apache.hudi.table.HoodieTable;
 import javax.annotation.concurrent.NotThreadSafe;
 
 import java.io.IOException;
-import java.util.List;
 import java.util.Map;
 import java.util.PriorityQueue;
 import java.util.Queue;
@@ -103,10 +101,10 @@ public class HoodieSortedMergeHandle<T, I, K, O> extends 
HoodieWriteMergeHandle<
   }
 
   @Override
-  public List<WriteStatus> close() {
+  protected void writeIncomingRecords() throws IOException {
     // write out any pending records (this can happen when inserts are turned 
into updates)
-    while (!newRecordKeysSorted.isEmpty()) {
-      try {
+    try {
+      while (!newRecordKeysSorted.isEmpty()) {
         String key = newRecordKeysSorted.poll();
         HoodieRecord<T> hoodieRecord = keyToNewRecords.get(key);
         if (!writtenRecordKeys.contains(hoodieRecord.getRecordKey())) {
@@ -118,13 +116,9 @@ public class HoodieSortedMergeHandle<T, I, K, O> extends 
HoodieWriteMergeHandle<
           insertRecordsWritten++;
           writtenRecordKeys.add(hoodieRecord.getRecordKey());
         }
-      } catch (IOException e) {
-        throw new HoodieUpsertException("Failed to close UpdateHandle", e);
       }
+    } finally {
+      newRecordKeysSorted.clear();
     }
-
-    newRecordKeysSorted.clear();
-
-    return super.close();
   }
 }
diff --git 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieWriteMergeHandle.java
 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieWriteMergeHandle.java
index b29746ec583e..4a426b2496d2 100644
--- 
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieWriteMergeHandle.java
+++ 
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieWriteMergeHandle.java
@@ -34,6 +34,7 @@ import org.apache.hudi.common.schema.HoodieSchema;
 import org.apache.hudi.common.serialization.DefaultSerializer;
 import org.apache.hudi.common.table.read.BufferedRecord;
 import org.apache.hudi.common.table.read.BufferedRecords;
+import org.apache.hudi.common.util.CloseableUtils;
 import org.apache.hudi.common.util.ConfigUtils;
 import org.apache.hudi.common.util.DefaultSizeEstimator;
 import org.apache.hudi.common.util.HoodieRecordSizeEstimator;
@@ -473,17 +474,10 @@ public class HoodieWriteMergeHandle<T, I, K, O> extends 
HoodieAbstractMergeHandl
       }
 
       markClosed();
-      writeIncomingRecords();
-
-      if (keyToNewRecords instanceof Closeable) {
-        ((Closeable) keyToNewRecords).close();
+      try (Closeable records = keyToNewRecords instanceof Closeable ? 
(Closeable) keyToNewRecords : null) {
+        writeIncomingRecords();
       }
-
-      keyToNewRecords = null;
-      writtenRecordKeys = null;
-
-      fileWriter.close();
-      fileWriter = null;
+      closeFileWriter();
 
       long fileSizeInBytes = storage.getPathInfo(newFilePath).getLength();
       HoodieWriteStat stat = writeStatus.getStat();
@@ -506,10 +500,32 @@ public class HoodieWriteMergeHandle<T, I, K, O> extends 
HoodieAbstractMergeHandl
 
       return Collections.singletonList(writeStatus);
     } catch (IOException e) {
+      closeFileWriterQuietly(e);
       throw new HoodieUpsertException("Failed to close UpdateHandle", e);
+    } catch (RuntimeException e) {
+      closeFileWriterQuietly(e);
+      throw e;
+    } finally {
+      keyToNewRecords = null;
+      writtenRecordKeys = null;
     }
   }
 
+  private void closeFileWriter() throws IOException {
+    try {
+      if (fileWriter != null) {
+        fileWriter.close();
+      }
+    } finally {
+      fileWriter = null;
+    }
+  }
+
+  protected void closeFileWriterQuietly(Throwable failure) {
+    CloseableUtils.closeSuppressing(fileWriter, failure);
+    fileWriter = null;
+  }
+
   public void performMergeDataValidationCheck(WriteStatus writeStatus) {
     if (!config.isMergeDataValidationCheckEnabled() || baseFileToMerge == 
null) {
       return;
diff --git 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieAppendHandle.java
 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieAppendHandle.java
index b1f194796c62..85b8aa18c23b 100644
--- 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieAppendHandle.java
+++ 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieAppendHandle.java
@@ -37,15 +37,27 @@ import org.junit.jupiter.api.extension.ExtendWith;
 import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.Arguments;
 import org.junit.jupiter.params.provider.MethodSource;
+import org.junit.jupiter.params.provider.ValueSource;
 import org.mockito.Mock;
 import org.mockito.junit.jupiter.MockitoExtension;
 
 import java.io.IOException;
+import java.util.Collections;
 import java.util.stream.Stream;
 
 import static 
org.apache.hudi.common.testutils.HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA;
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+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.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.spy;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
 @ExtendWith(MockitoExtension.class)
@@ -87,6 +99,37 @@ public class TestHoodieAppendHandle extends 
HoodieCommonTestHarness {
     when(mockHoodieTable.getMetaClient()).thenReturn(metaClient);
   }
 
+  @ParameterizedTest
+  @ValueSource(booleans = {false, true})
+  void testFailedFlushClosesWriterAndPreventsAnotherFlush(boolean closeFails) 
throws IOException {
+    writeConfig = HoodieWriteConfig.newBuilder()
+        .withProps(writeConfig.getProps())
+        .withWriteTableVersion(HoodieTableVersion.SIX.versionCode())
+        .build();
+    when(mockHoodieTable.getStorage()).thenReturn(metaClient.getStorage());
+    HoodieInlineLogAppendHandle<Object, Object, Object, Object> handle = spy(
+        new HoodieInlineLogAppendHandle<>(writeConfig, TEST_INSTANT_TIME, 
mockHoodieTable,
+            TEST_PARTITION_PATH, TEST_FILE_ID, taskContextSupplier));
+    HoodieLogFormat.Writer writer = mock(HoodieLogFormat.Writer.class);
+    handle.writer = writer;
+    handle.recordItr = Collections.emptyIterator();
+    RuntimeException failure = new IllegalStateException("flush failed");
+    RuntimeException closeFailure = new IllegalStateException("close failed");
+    doThrow(failure).when(handle).flushAppend();
+    if (closeFails) {
+      doThrow(closeFailure).when(writer).close();
+    }
+
+    assertSame(failure, assertThrows(IllegalStateException.class, 
handle::doAppend));
+    assertTrue(handle.isClosed());
+    assertNull(handle.writer);
+    assertNull(handle.recordItr);
+    assertArrayEquals(closeFails ? new Throwable[] {closeFailure} : new 
Throwable[0], failure.getSuppressed());
+    assertDoesNotThrow(handle::close);
+    verify(handle, times(1)).flushAppend();
+    verify(writer, times(1)).close();
+  }
+
   private static Stream<Arguments> versionsSixAndAbove() {
     return Stream.of(
         Arguments.of(HoodieTableVersion.SIX),
@@ -114,7 +157,7 @@ public class TestHoodieAppendHandle extends 
HoodieCommonTestHarness {
       when(mockedFSView.getLatestBaseFile(TEST_PARTITION_PATH, 
TEST_FILE_ID)).thenReturn(Option.empty());
     }
 
-    HoodieAppendHandle<Object, Object, Object, Object> appendHandle =
+    HoodieInlineLogAppendHandle<Object, Object, Object, Object> appendHandle =
         new HoodieInlineLogAppendHandle<>(writeConfig, TEST_INSTANT_TIME, 
mockHoodieTable, TEST_PARTITION_PATH, TEST_FILE_ID, taskContextSupplier);
 
     FileSlice mockFileSlice = mock(FileSlice.class);
diff --git 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieBinaryCopyHandle.java
 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieBinaryCopyHandle.java
new file mode 100644
index 000000000000..5d0c59e8f6c2
--- /dev/null
+++ 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieBinaryCopyHandle.java
@@ -0,0 +1,74 @@
+/*
+ * 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.io;
+
+import org.apache.hudi.common.engine.LocalTaskContextSupplier;
+import org.apache.hudi.common.testutils.HoodieCommonTestHarness;
+import org.apache.hudi.config.HoodieClusteringConfig;
+import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.exception.HoodieIOException;
+import org.apache.hudi.parquet.io.HoodieParquetFileBinaryCopier;
+import org.apache.hudi.storage.StoragePath;
+import org.apache.hudi.table.HoodieTable;
+import org.apache.hudi.table.TestBaseHoodieTable;
+
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedConstruction;
+
+import java.io.IOException;
+import java.util.Collections;
+
+import static 
org.apache.hudi.common.testutils.HoodieTestDataGenerator.TRIP_EXAMPLE_SCHEMA;
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyBoolean;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.verify;
+
+class TestHoodieBinaryCopyHandle extends HoodieCommonTestHarness {
+  @Test
+  void testWriteFailureClosesCopierAndPreservesOriginalException() throws 
Exception {
+    initPath();
+    initMetaClient();
+    HoodieWriteConfig config = 
HoodieWriteConfig.newBuilder().withPath(basePath)
+        .withSchema(TRIP_EXAMPLE_SCHEMA).build();
+    
config.setValue(HoodieClusteringConfig.FILE_STITCHING_BINARY_COPY_SCHEMA_EVOLUTION_ENABLE,
 "true");
+    HoodieTable table = new TestBaseHoodieTable(config, getEngineContext(), 
metaClient);
+    IOException failure = new IOException("copy failed");
+    IOException closeFailure = new IOException("close failed");
+    try (MockedConstruction<HoodieParquetFileBinaryCopier> copiers = 
mockConstruction(
+        HoodieParquetFileBinaryCopier.class, (copier, context) -> {
+          doThrow(failure).when(copier).binaryCopy(any(), any(), any(), 
anyBoolean());
+          doThrow(closeFailure).when(copier).close();
+        })) {
+      HoodieBinaryCopyHandle handle = new HoodieBinaryCopyHandle(config, 
"100", "partition", "file-1",
+          table, new LocalTaskContextSupplier(), Collections.singletonList(new 
StoragePath(basePath, "input.parquet")));
+      assertSame(failure, assertThrows(HoodieIOException.class, 
handle::write).getCause());
+      assertArrayEquals(new Throwable[] {closeFailure}, 
failure.getSuppressed());
+      assertDoesNotThrow(handle::close);
+      verify(copiers.constructed().get(0)).close();
+    } finally {
+      cleanMetaClient();
+    }
+  }
+}
diff --git 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieCreateHandle.java
 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieCreateHandle.java
index 240cb6622ca1..5c34a9e1abdb 100644
--- 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieCreateHandle.java
+++ 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieCreateHandle.java
@@ -35,6 +35,7 @@ import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.ParquetUtils;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.core.io.storage.HoodieFileWriter;
+import org.apache.hudi.exception.HoodieException;
 import org.apache.hudi.exception.HoodieInsertException;
 import org.apache.hudi.io.storage.TestFileWriter;
 import org.apache.hudi.storage.HoodieStorage;
@@ -342,21 +343,6 @@ public class TestHoodieCreateHandle extends 
HoodieCommonTestHarness {
 
   @Test
   void testMarkerFileCreatedWhenFileWriterWriteFails() throws Exception {
-    // Create a custom HoodieCreateHandle that uses TestFileWriter that 
succeeds on initialization but fails on write
-    class CreateHandleWithFileWriterWriteFailure extends 
HoodieCreateHandle<Object, Object, Object, Object> {
-      public CreateHandleWithFileWriterWriteFailure(HoodieWriteConfig config, 
String instantTime,
-                                                HoodieTable<Object, Object, 
Object, Object> hoodieTable,
-                                                String partitionPath, String 
fileId,
-                                                TaskContextSupplier 
taskContextSupplier) {
-        super(config, instantTime, hoodieTable, partitionPath, fileId, 
taskContextSupplier);
-      }
-
-      @Override
-      protected HoodieFileWriter initializeFileWriter() throws IOException {
-        return new TestFileWriter(path, hoodieTable.getStorage(), false, true);
-      }
-    }
-
     CreateHandleWithFileWriterWriteFailure createHandle = new 
CreateHandleWithFileWriterWriteFailure(
         writeConfig, TEST_INSTANT_TIME, hoodieTable, TEST_PARTITION_PATH, 
TEST_FILE_ID, taskContextSupplier);
     
@@ -392,6 +378,42 @@ public class TestHoodieCreateHandle extends 
HoodieCommonTestHarness {
     assertDoesNotThrow(createHandle::close);
   }
 
+  @Test
+  void testFileWriterClosedWhenDoWriteFails() throws Exception {
+    HoodieWriteConfig failOnWriteConfig = HoodieWriteConfig.newBuilder()
+        .withProps(writeConfig.getProps())
+        .withWriteIgnoreFailed(false)
+        .build();
+    HoodieTable failOnWriteTable = new TestBaseHoodieTable(failOnWriteConfig, 
getEngineContext(), metaClient);
+    CreateHandleWithFileWriterWriteFailure createHandle = new 
CreateHandleWithFileWriterWriteFailure(
+        failOnWriteConfig, TEST_INSTANT_TIME, failOnWriteTable, 
TEST_PARTITION_PATH, TEST_FILE_ID, taskContextSupplier);
+    HoodieRecord testRecord = dataGen.generateInserts(TEST_INSTANT_TIME, 
1).get(0);
+
+    HoodieFileWriter fileWriter = createHandle.fileWriter;
+    HoodieException exception = assertThrows(HoodieException.class, () ->
+        createHandle.doWrite(testRecord, TEST_SCHEMA, new TypedProperties()));
+
+    assertEquals("Simulated file writer write failure", 
exception.getMessage());
+    assertNull(createHandle.fileWriter);
+    assertTrue(createHandle.isClosed());
+    assertFalse(fileWriter.canWrite());
+    assertDoesNotThrow(createHandle::close);
+  }
+
+  private static class CreateHandleWithFileWriterWriteFailure extends 
HoodieCreateHandle<Object, Object, Object, Object> {
+    CreateHandleWithFileWriterWriteFailure(HoodieWriteConfig config, String 
instantTime,
+                                           HoodieTable<Object, Object, Object, 
Object> hoodieTable,
+                                           String partitionPath, String fileId,
+                                           TaskContextSupplier 
taskContextSupplier) {
+      super(config, instantTime, hoodieTable, partitionPath, fileId, 
taskContextSupplier);
+    }
+
+    @Override
+    protected HoodieFileWriter initializeFileWriter() throws IOException {
+      return new TestFileWriter(path, hoodieTable.getStorage(), false, true);
+    }
+  }
+
   /**
    * Validates that the written records match the original records and that 
all write statistics are correct.
    */
diff --git 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieNativeLogAppendHandle.java
 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieNativeLogAppendHandle.java
index e0ff62eb0ea4..cf68cf4367e8 100644
--- 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieNativeLogAppendHandle.java
+++ 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestHoodieNativeLogAppendHandle.java
@@ -35,12 +35,16 @@ import org.apache.hudi.common.table.log.AppendResult;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.exception.HoodieAppendException;
+import org.apache.hudi.exception.HoodieEarlyConflictDetectionException;
+import org.apache.hudi.exception.HoodieException;
 import org.apache.hudi.io.cdc.HoodieNativeLogFormatWriter;
 import org.apache.hudi.storage.HoodieStorage;
 import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.table.HoodieTable;
 
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
 import org.mockito.InOrder;
 import org.mockito.MockedConstruction;
 
@@ -50,7 +54,11 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
 
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+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;
@@ -214,6 +222,60 @@ public class TestHoodieNativeLogAppendHandle {
     }
   }
 
+  @ParameterizedTest
+  @ValueSource(booleans = {false, true})
+  public void testWriteFailureClosesNativeWriter(boolean useRecordMap) throws 
Exception {
+    HoodieWriteConfig config = config();
+    // Conflict detection must fail the handle even when individual write 
failures may be ignored.
+    config.setValue(HoodieWriteConfig.IGNORE_FAILED, "true");
+    HoodieTable table = table(config);
+    HoodieEarlyConflictDetectionException failure = new 
HoodieEarlyConflictDetectionException("conflict");
+    RuntimeException closeFailure = new IllegalStateException("close failed");
+    try (MockedConstruction<HoodieNativeLogFormatWriter> writers = 
mockConstruction(
+        HoodieNativeLogFormatWriter.class, (writer, context) -> {
+          when(writer.canWriteDataFile()).thenReturn(true);
+          doThrow(failure).when(writer).appendRecord(any(), any());
+          doThrow(closeFailure).when(writer).close();
+        })) {
+      TestableNativeLogAppendHandle handle = new 
TestableNativeLogAppendHandle(config, table);
+      handle.createWriter();
+      handle.doInit = false;
+      HoodieRecord record = mock(HoodieRecord.class);
+      when(record.getPartitionPath()).thenReturn("partition");
+      when(record.getMetadata()).thenReturn(Option.empty());
+      when(record.prependMetaFields(any(), any(), any(), 
any())).thenReturn(record);
+
+      assertThrows(HoodieException.class, () -> {
+        if (useRecordMap) {
+          handle.write(Collections.singletonMap("key", record));
+        } else {
+          handle.writeRecord(record);
+        }
+      });
+      assertTrue(handle.isClosed());
+      assertNull(handle.getWriter());
+      assertDoesNotThrow(handle::close);
+      assertArrayEquals(new Throwable[] {closeFailure}, 
failure.getSuppressed());
+      verify(writers.constructed().get(0)).close();
+    }
+  }
+
+  @Test
+  public void testCloseFailureClearsNativeWriter() throws Exception {
+    HoodieWriteConfig config = config();
+    RuntimeException failure = new IllegalStateException("close failed");
+    try (MockedConstruction<HoodieNativeLogFormatWriter> writers = 
mockConstruction(
+        HoodieNativeLogFormatWriter.class, (writer, context) -> 
doThrow(failure).when(writer).close())) {
+      TestableNativeLogAppendHandle handle = new 
TestableNativeLogAppendHandle(config, table(config));
+      handle.createWriter();
+      assertSame(failure, assertThrows(IllegalStateException.class, 
handle::close));
+      assertTrue(handle.isClosed());
+      assertNull(handle.getWriter());
+      assertDoesNotThrow(handle::close);
+      verify(writers.constructed().get(0)).close();
+    }
+  }
+
   private static HoodieWriteConfig config() {
     return HoodieWriteConfig.newBuilder()
         .withPath("/tmp")
diff --git 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestSortedAndChangeLogMergeHandles.java
 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestSortedAndChangeLogMergeHandles.java
index d5731ab37ba6..54322e4a392e 100644
--- 
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestSortedAndChangeLogMergeHandles.java
+++ 
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/io/TestSortedAndChangeLogMergeHandles.java
@@ -20,6 +20,7 @@ package org.apache.hudi.io;
 
 import org.apache.hudi.client.WriteStatus;
 import org.apache.hudi.common.config.RecordMergeMode;
+import org.apache.hudi.common.config.TypedProperties;
 import org.apache.hudi.common.engine.HoodieReaderContext;
 import org.apache.hudi.common.engine.LocalTaskContextSupplier;
 import org.apache.hudi.common.engine.ReaderContextFactory;
@@ -33,7 +34,10 @@ import org.apache.hudi.common.schema.HoodieSchema;
 import org.apache.hudi.common.table.HoodieTableConfig;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.table.HoodieTableVersion;
+import org.apache.hudi.common.table.read.HoodieRecordReader;
 import org.apache.hudi.common.util.Option;
+import org.apache.hudi.common.util.collection.ClosableIterator;
+import org.apache.hudi.common.util.collection.ExternalSpillableMap;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.core.io.storage.HoodieFileWriter;
 import org.apache.hudi.core.io.storage.HoodieFileWriterFactory;
@@ -61,15 +65,22 @@ import java.util.List;
 import java.util.Map;
 import java.util.Properties;
 
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
 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.anyBoolean;
 import static org.mockito.ArgumentMatchers.anyString;
+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.spy;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
@@ -240,6 +251,101 @@ public class TestSortedAndChangeLogMergeHandles {
     }
   }
 
+  @Test
+  public void testWriterCloseFailureDoesNotCloseWriterTwice() throws Exception 
{
+    HoodieWriteConfig config = config();
+    TestContext context = new TestContext(config);
+    IOException failure = new IOException("writer close failed");
+    try (MockedStatic<WriteMarkersFactory> markers = 
mockStatic(WriteMarkersFactory.class);
+         MockedStatic<HoodieFileWriterFactory> writers = 
mockStatic(HoodieFileWriterFactory.class)) {
+      context.stubWriters(markers, writers);
+      TestableWriteMergeHandle handle = new TestableWriteMergeHandle(config, 
context.table, new HashMap<>());
+      doThrow(failure).when(context.fileWriter).close();
+      assertSame(failure, assertThrows(HoodieUpsertException.class, 
handle::close).getCause());
+      assertNull(handle.fileWriter);
+      assertDoesNotThrow(handle::close);
+      verify(context.fileWriter).close();
+    }
+  }
+
+  @Test
+  public void testPendingRecordFailureClosesSpillableMap() throws Exception {
+    HoodieWriteConfig config = config();
+    TestContext context = new TestContext(config);
+    IOException failure = new IOException("pending records failed");
+    try (MockedStatic<WriteMarkersFactory> markers = 
mockStatic(WriteMarkersFactory.class);
+         MockedStatic<HoodieFileWriterFactory> writers = 
mockStatic(HoodieFileWriterFactory.class)) {
+      context.stubWriters(markers, writers);
+      TestableWriteMergeHandle handle = spy(new 
TestableWriteMergeHandle(config, context.table, new HashMap<>()));
+      ExternalSpillableMap map = mock(ExternalSpillableMap.class);
+      handle.keyToNewRecords = map;
+      doThrow(failure).when(handle).writeIncomingRecords();
+      assertSame(failure, assertThrows(HoodieUpsertException.class, 
handle::close).getCause());
+      assertNull(handle.keyToNewRecords);
+      assertDoesNotThrow(handle::close);
+      verify(map).close();
+      verify(context.fileWriter).close();
+    }
+  }
+
+  @Test
+  public void testFileGroupMergeWriteFailureClosesWriter() throws Exception {
+    HoodieWriteConfig config = 
HoodieWriteConfig.newBuilder().withProps(config().getProps())
+        .withWriteIgnoreFailed(false).build();
+    TestContext context = new TestContext(config);
+    IOException failure = new IOException("record write failed");
+    IOException closeFailure = new IOException("writer close failed");
+    try (MockedStatic<WriteMarkersFactory> markers = 
mockStatic(WriteMarkersFactory.class);
+         MockedStatic<HoodieFileWriterFactory> writers = 
mockStatic(HoodieFileWriterFactory.class)) {
+      context.stubWriters(markers, writers);
+      FileGroupReaderBasedMergeHandle handle = spy(new 
FileGroupReaderBasedMergeHandle(
+          config, "100", context.table, 
MergeContext.create(Collections.emptyIterator()), "partition", "file-1",
+          new LocalTaskContextSupplier(), null, Option.empty()));
+      HoodieRecordReader reader = mock(HoodieRecordReader.class);
+      ClosableIterator iterator = mock(ClosableIterator.class);
+      when(reader.getClosableHoodieRecordIterator()).thenReturn(iterator);
+      when(iterator.hasNext()).thenReturn(true);
+      HoodieRecord inputRecord = record("key");
+      when(iterator.next()).thenReturn(inputRecord);
+      doReturn(reader).when(handle).getFileGroupReader(anyBoolean(), any(), 
any(), any(), any());
+      doThrow(failure).when(handle).writeToFile(any(), any(), any(), any(), 
anyBoolean());
+      doThrow(closeFailure).when(context.fileWriter).close();
+
+      HoodieUpsertException actual = assertThrows(HoodieUpsertException.class, 
handle::doMerge);
+      assertSame(failure, actual.getCause().getCause());
+      assertArrayEquals(new Throwable[] {closeFailure}, 
actual.getCause().getSuppressed());
+      assertNull(handle.fileWriter);
+      verify(context.fileWriter).close();
+      verify(iterator).close();
+      verify(reader).close();
+    }
+  }
+
+  @Test
+  public void testMergeCloseFailureClosesCDCWriter() throws Exception {
+    HoodieWriteConfig config = config();
+    TestContext context = new TestContext(config);
+    HoodieCDCLogWriter<IndexedRecord> cdcWriter = 
mock(HoodieCDCLogWriter.class);
+    IOException failure = new IOException("base writer close failed");
+    RuntimeException cdcFailure = new IllegalStateException("CDC close 
failed");
+    try (MockedStatic<WriteMarkersFactory> markers = 
mockStatic(WriteMarkersFactory.class);
+         MockedStatic<HoodieFileWriterFactory> writers = 
mockStatic(HoodieFileWriterFactory.class);
+         MockedStatic<HoodieCDCLogWriterFactory> cdcWriters = 
mockStatic(HoodieCDCLogWriterFactory.class)) {
+      context.stubWriters(markers, writers);
+      cdcWriters.when(() -> HoodieCDCLogWriterFactory.createAvroCDCLogWriter(
+          anyString(), any(), any(), anyString(), any(), any(), anyString(), 
anyString(), any(), any(), any()))
+          .thenReturn(cdcWriter);
+      TestableChangeLogMergeHandle handle = new 
TestableChangeLogMergeHandle(config, context.table, new HashMap<>());
+      doThrow(failure).when(context.fileWriter).close();
+      doThrow(cdcFailure).when(cdcWriter).close();
+      HoodieUpsertException actual = assertThrows(HoodieUpsertException.class, 
handle::close);
+      assertSame(failure, actual.getCause());
+      assertArrayEquals(new Throwable[] {cdcFailure}, actual.getSuppressed());
+      verify(context.fileWriter).close();
+      verify(cdcWriter).close();
+    }
+  }
+
   private static HoodieWriteConfig config() {
     return HoodieWriteConfig.newBuilder()
         .withPath("/tmp")
@@ -280,6 +386,7 @@ public class TestSortedAndChangeLogMergeHandles {
       when(metaClient.getBasePath()).thenReturn(new StoragePath("/tmp"));
       when(metaClient.getTableConfig()).thenReturn(tableConfig);
       when(metaClient.getIndexMetadata()).thenReturn(Option.empty());
+      when(tableConfig.getProps()).thenReturn(new TypedProperties());
       when(tableConfig.getTableVersion()).thenReturn(HoodieTableVersion.TEN);
       when(tableConfig.getMetaFieldsMode()).thenReturn(MetaFieldsMode.ALL);
       
when(tableConfig.getRecordMergeMode()).thenReturn(RecordMergeMode.COMMIT_TIME_ORDERING);
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataParquetOutputStreamWriter.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataParquetOutputStreamWriter.java
index ca65291e2c69..dbe4e19677c9 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataParquetOutputStreamWriter.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataParquetOutputStreamWriter.java
@@ -32,6 +32,8 @@ import org.apache.parquet.hadoop.api.WriteSupport;
 
 import java.io.IOException;
 
+import static org.apache.hudi.io.util.FileIOUtils.closeQuietly;
+
 /**
  * An implementation for {@link HoodieRowDataFileWriter} which is used to 
write hoodie records into an {@link FSDataOutputStream}.
  *
@@ -67,7 +69,12 @@ public class HoodieRowDataParquetOutputStreamWriter 
implements HoodieRowDataFile
     
parquetWriterbuilder.withDictionaryEncoding(parquetConfig.isDictionaryEnabled());
     
parquetWriterbuilder.withWriterVersion(ParquetWriter.DEFAULT_WRITER_VERSION);
     
parquetWriterbuilder.withConf(parquetConfig.getStorageConf().unwrapAs(Configuration.class));
-    this.writer = parquetWriterbuilder.build();
+    try {
+      this.writer = parquetWriterbuilder.build();
+    } catch (IOException | RuntimeException e) {
+      closeQuietly(outputStream);
+      throw e;
+    }
   }
 
   @Override
diff --git 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/HoodieSparkParquetStreamWriter.java
 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/HoodieSparkParquetStreamWriter.java
index ab3e4277e286..fbad6e654a0b 100644
--- 
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/HoodieSparkParquetStreamWriter.java
+++ 
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/io/storage/HoodieSparkParquetStreamWriter.java
@@ -35,6 +35,8 @@ import org.apache.spark.unsafe.types.UTF8String;
 
 import java.io.IOException;
 
+import static org.apache.hudi.io.util.FileIOUtils.closeQuietly;
+
 public class HoodieSparkParquetStreamWriter implements HoodieSparkFileWriter, 
AutoCloseable {
   private final ParquetWriter<InternalRow> writer;
   private final HoodieRowParquetWriteSupport writeSupport;
@@ -42,16 +44,21 @@ public class HoodieSparkParquetStreamWriter implements 
HoodieSparkFileWriter, Au
   public HoodieSparkParquetStreamWriter(FSDataOutputStream outputStream,
       HoodieRowParquetConfig parquetConfig) throws IOException {
     this.writeSupport = parquetConfig.getWriteSupport();
-    this.writer = new Builder<>(new 
OutputStreamBackedOutputFile(outputStream), writeSupport)
-        .withWriteMode(ParquetFileWriter.Mode.CREATE)
-        .withCompressionCodec(parquetConfig.getCompressionCodecName())
-        .withRowGroupSize(parquetConfig.getBlockSize())
-        .withPageSize(parquetConfig.getPageSize())
-        .withDictionaryPageSize(parquetConfig.getPageSize())
-        .withDictionaryEncoding(parquetConfig.isDictionaryEnabled())
-        .withWriterVersion(ParquetWriter.DEFAULT_WRITER_VERSION)
-        .withConf(parquetConfig.getHadoopConf())
-        .build();
+    try {
+      this.writer = new Builder<>(new 
OutputStreamBackedOutputFile(outputStream), writeSupport)
+          .withWriteMode(ParquetFileWriter.Mode.CREATE)
+          .withCompressionCodec(parquetConfig.getCompressionCodecName())
+          .withRowGroupSize(parquetConfig.getBlockSize())
+          .withPageSize(parquetConfig.getPageSize())
+          .withDictionaryPageSize(parquetConfig.getPageSize())
+          .withDictionaryEncoding(parquetConfig.isDictionaryEnabled())
+          .withWriterVersion(ParquetWriter.DEFAULT_WRITER_VERSION)
+          .withConf(parquetConfig.getHadoopConf())
+          .build();
+    } catch (IOException | RuntimeException e) {
+      closeQuietly(outputStream);
+      throw e;
+    }
   }
 
   @Override
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/bootstrap/index/hfile/HFileBootstrapIndexWriter.java
 
b/hudi-common/src/main/java/org/apache/hudi/common/bootstrap/index/hfile/HFileBootstrapIndexWriter.java
index 8fd232a87402..d1b5a1704019 100644
--- 
a/hudi-common/src/main/java/org/apache/hudi/common/bootstrap/index/hfile/HFileBootstrapIndexWriter.java
+++ 
b/hudi-common/src/main/java/org/apache/hudi/common/bootstrap/index/hfile/HFileBootstrapIndexWriter.java
@@ -26,6 +26,7 @@ import org.apache.hudi.common.bootstrap.index.BootstrapIndex;
 import org.apache.hudi.common.model.BootstrapFileMapping;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.table.timeline.TimelineMetadataUtils;
+import org.apache.hudi.common.util.CloseableUtils;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.collection.Pair;
 import org.apache.hudi.exception.HoodieException;
@@ -51,6 +52,7 @@ import static 
org.apache.hudi.common.bootstrap.index.hfile.HFileBootstrapIndex.f
 import static 
org.apache.hudi.common.bootstrap.index.hfile.HFileBootstrapIndex.getFileGroupKey;
 import static 
org.apache.hudi.common.bootstrap.index.hfile.HFileBootstrapIndex.getPartitionKey;
 import static 
org.apache.hudi.common.bootstrap.index.hfile.HFileBootstrapIndex.partitionIndexPath;
+import static org.apache.hudi.io.util.FileIOUtils.closeQuietly;
 
 @Slf4j
 public class HFileBootstrapIndexWriter extends BootstrapIndex.IndexWriter {
@@ -174,14 +176,16 @@ public class HFileBootstrapIndexWriter extends 
BootstrapIndex.IndexWriter {
    * Close Writer Handles.
    */
   public void close() {
-    try {
-      if (!closed) {
-        indexByPartitionWriter.close();
-        indexByFileIdWriter.close();
-        closed = true;
-      }
-    } catch (IOException ioe) {
-      throw new HoodieIOException(ioe.getMessage(), ioe);
+    if (closed) {
+      return;
+    }
+    Exception failure = closeHFileWriter(indexByPartitionWriter, null);
+    failure = closeHFileWriter(indexByFileIdWriter, failure);
+    indexByPartitionWriter = null;
+    indexByFileIdWriter = null;
+    closed = true;
+    if (failure != null) {
+      throw new HoodieException(failure.getMessage(), failure);
     }
   }
 
@@ -189,13 +193,43 @@ public class HFileBootstrapIndexWriter extends 
BootstrapIndex.IndexWriter {
   public void begin() {
     try {
       HFileContext context = HFileContext.builder().build();
-      OutputStream outputStreamForPartitionWriter = 
metaClient.getStorage().create(indexByPartitionPath);
-      this.indexByPartitionWriter = new HFileWriterImpl(context, 
outputStreamForPartitionWriter);
-      OutputStream outputStreamForFileIdWriter = 
metaClient.getStorage().create(indexByFileIdPath);
-      this.indexByFileIdWriter = new HFileWriterImpl(context, 
outputStreamForFileIdWriter);
+      this.indexByPartitionWriter = createHFileWriter(indexByPartitionPath, 
context);
+      this.indexByFileIdWriter = createHFileWriter(indexByFileIdPath, context);
     } catch (IOException ioe) {
+      CloseableUtils.closeSuppressing(this::close, ioe);
       throw new HoodieIOException(ioe.getMessage(), ioe);
+    } catch (RuntimeException re) {
+      CloseableUtils.closeSuppressing(this::close, re);
+      throw re;
+    }
+  }
+
+  private HFileWriter createHFileWriter(StoragePath path, HFileContext 
context) throws IOException {
+    OutputStream outputStream = metaClient.getStorage().create(path);
+    HFileWriter writer = null;
+    try {
+      writer = new HFileWriterImpl(context, outputStream);
+      return writer;
+    } finally {
+      if (writer == null) {
+        closeQuietly(outputStream);
+      }
+    }
+  }
+
+  private Exception closeHFileWriter(HFileWriter writer, Exception failure) {
+    if (writer == null) {
+      return failure;
+    }
+    try {
+      writer.close();
+    } catch (IOException | RuntimeException e) {
+      if (failure == null) {
+        return e;
+      }
+      failure.addSuppressed(e);
     }
+    return failure;
   }
 
   @Override
diff --git 
a/hudi-common/src/main/java/org/apache/hudi/common/util/CloseableUtils.java 
b/hudi-common/src/main/java/org/apache/hudi/common/util/CloseableUtils.java
index 50b6e2e6ee80..44c8480ed970 100644
--- a/hudi-common/src/main/java/org/apache/hudi/common/util/CloseableUtils.java
+++ b/hudi-common/src/main/java/org/apache/hudi/common/util/CloseableUtils.java
@@ -24,12 +24,17 @@ public final class CloseableUtils {
   private CloseableUtils() {
   }
 
-  /** Closes {@code closeable}, attaching any failure to {@code primary} as a 
suppressed exception. */
+  /** Closes a non-null resource, attaching any distinct close failure to the 
original failure. */
   public static void closeSuppressing(AutoCloseable closeable, Throwable 
primary) {
+    if (closeable == null) {
+      return;
+    }
     try {
       closeable.close();
     } catch (Throwable closeError) {
-      primary.addSuppressed(closeError);
+      if (closeError != primary) {
+        primary.addSuppressed(closeError);
+      }
     }
   }
 }
diff --git 
a/hudi-common/src/test/java/org/apache/hudi/common/util/TestCloseableUtils.java 
b/hudi-common/src/test/java/org/apache/hudi/common/util/TestCloseableUtils.java
index f18339047905..57a6997a9e83 100644
--- 
a/hudi-common/src/test/java/org/apache/hudi/common/util/TestCloseableUtils.java
+++ 
b/hudi-common/src/test/java/org/apache/hudi/common/util/TestCloseableUtils.java
@@ -19,22 +19,54 @@
 package org.apache.hudi.common.util;
 
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.MethodSource;
 
 import java.io.IOException;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.stream.Stream;
 
 import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 
 class TestCloseableUtils {
 
+  static Stream<Throwable> closeFailures() {
+    return Stream.of(new IOException("IO failure"), new 
IllegalStateException("runtime failure"),
+        new Exception("checked failure"), new AssertionError("error"));
+  }
+
   @Test
-  void testCloseSuppressing() {
-    IOException primary = new IOException("primary");
-    IOException closeError = new IOException("close");
+  void testNullAndSuccessfulClose() {
+    IOException failure = new IOException("write failed");
+    assertDoesNotThrow(() -> CloseableUtils.closeSuppressing(null, failure));
+    AtomicBoolean closed = new AtomicBoolean();
+    CloseableUtils.closeSuppressing(() -> closed.set(true), failure);
+    assertTrue(closed.get());
+    assertArrayEquals(new Throwable[0], failure.getSuppressed());
+  }
 
-    CloseableUtils.closeSuppressing(() -> {
-      throw closeError;
-    }, primary);
+  @ParameterizedTest
+  @MethodSource("closeFailures")
+  void testPreservesOriginalFailure(Throwable closeFailure) {
+    IOException failure = new IOException("write failed");
+    AutoCloseable closeable = () -> {
+      if (closeFailure instanceof Error) {
+        throw (Error) closeFailure;
+      }
+      throw (Exception) closeFailure;
+    };
+    assertDoesNotThrow(() -> CloseableUtils.closeSuppressing(closeable, 
failure));
+    assertArrayEquals(new Throwable[] {closeFailure}, failure.getSuppressed());
+  }
 
-    assertArrayEquals(new Throwable[] {closeError}, primary.getSuppressed());
+  @Test
+  void testDoesNotSuppressFailureOnItself() {
+    IOException failure = new IOException("write failed");
+    assertDoesNotThrow(() -> CloseableUtils.closeSuppressing(() -> {
+      throw failure;
+    }, failure));
+    assertArrayEquals(new Throwable[0], failure.getSuppressed());
   }
 }
diff --git 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/ParquetUtils.java
 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/ParquetUtils.java
index ee391bb826ce..31b0e0868fe6 100644
--- 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/ParquetUtils.java
+++ 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/common/util/ParquetUtils.java
@@ -532,13 +532,14 @@ public class ParquetUtils extends FileFormatUtils {
 
     HoodieFileWriter parquetWriter = HoodieFileWriterFactory.getFileWriter(
         HoodieFileFormat.PARQUET, outputStream, storage, config, writerSchema, 
recordType);
-    while (recordItr.hasNext()) {
-      HoodieRecord record = recordItr.next();
-      String recordKey = record.getRecordKey(readerSchema, keyFieldName);
-      parquetWriter.write(recordKey, record, writerSchema);
+    try (HoodieFileWriter writerToClose = parquetWriter) {
+      while (recordItr.hasNext()) {
+        HoodieRecord record = recordItr.next();
+        String recordKey = record.getRecordKey(readerSchema, keyFieldName);
+        parquetWriter.write(recordKey, record, writerSchema);
+      }
+      outputStream.flush();
     }
-    outputStream.flush();
-    parquetWriter.close();
     return Pair.of(outputStream, parquetWriter.getFileFormatMetadata());
   }
 
diff --git 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroHFileWriter.java
 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroHFileWriter.java
index 901a4766ad8a..62c77401067e 100644
--- 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroHFileWriter.java
+++ 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieAvroHFileWriter.java
@@ -24,6 +24,7 @@ import org.apache.hudi.common.engine.TaskContextSupplier;
 import org.apache.hudi.common.model.HoodieKey;
 import org.apache.hudi.common.schema.HoodieSchema;
 import org.apache.hudi.common.schema.HoodieSchemaField;
+import org.apache.hudi.common.util.CloseableUtils;
 import org.apache.hudi.common.util.HFileUtils;
 import org.apache.hudi.common.util.HoodieStorageUtils;
 import org.apache.hudi.common.util.Option;
@@ -106,11 +107,17 @@ public class HoodieAvroHFileWriter
         .build();
     StorageConfiguration<Configuration> storageConf = new 
HadoopStorageConfiguration(conf);
     StoragePath filePath = new StoragePath(this.file.toUri());
-    OutputStream outputStream =  HoodieStorageUtils.getStorage(filePath, 
storageConf).create(filePath);
-    this.writer = new HFileWriterImpl(context, outputStream);
-    this.prevRecordKey = "";
-    writer.appendFileInfo(
-        HoodieAvroHFileReaderImplBase.SCHEMA_KEY, 
getUTF8Bytes(schema.toString()));
+    OutputStream outputStream = HoodieStorageUtils.getStorage(filePath, 
storageConf).create(filePath);
+    try {
+      this.writer = new HFileWriterImpl(context, outputStream);
+      this.prevRecordKey = "";
+      writer.appendFileInfo(
+          HoodieAvroHFileReaderImplBase.SCHEMA_KEY, 
getUTF8Bytes(schema.toString()));
+    } catch (RuntimeException e) {
+      CloseableUtils.closeSuppressing(writer != null ? writer : outputStream, 
e);
+      writer = null;
+      throw e;
+    }
   }
 
   @Override
diff --git 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieParquetStreamWriter.java
 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieParquetStreamWriter.java
index eb10677a1ba7..c517e9e02d33 100644
--- 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieParquetStreamWriter.java
+++ 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/io/storage/hadoop/HoodieParquetStreamWriter.java
@@ -36,6 +36,8 @@ import org.apache.parquet.io.OutputFile;
 
 import java.io.IOException;
 
+import static org.apache.hudi.io.util.FileIOUtils.closeQuietly;
+
 /**
  * Hudi log block writer for parquet format.
  * <p>
@@ -49,16 +51,21 @@ public class HoodieParquetStreamWriter implements 
HoodieAvroFileWriter, AutoClos
   public HoodieParquetStreamWriter(FSDataOutputStream outputStream,
                                    HoodieParquetConfig<HoodieAvroWriteSupport> 
parquetConfig) throws IOException {
     this.writeSupport = parquetConfig.getWriteSupport();
-    this.writer = new Builder<IndexedRecord>(new 
OutputStreamBackedOutputFile(outputStream), writeSupport)
-        .withWriteMode(ParquetFileWriter.Mode.CREATE)
-        .withCompressionCodec(parquetConfig.getCompressionCodecName())
-        .withRowGroupSize(parquetConfig.getBlockSize())
-        .withPageSize(parquetConfig.getPageSize())
-        .withDictionaryPageSize(parquetConfig.getPageSize())
-        .withDictionaryEncoding(parquetConfig.isDictionaryEnabled())
-        .withWriterVersion(ParquetWriter.DEFAULT_WRITER_VERSION)
-        .withConf(parquetConfig.getStorageConf().unwrapAs(Configuration.class))
-        .build();
+    try {
+      this.writer = new Builder<IndexedRecord>(new 
OutputStreamBackedOutputFile(outputStream), writeSupport)
+          .withWriteMode(ParquetFileWriter.Mode.CREATE)
+          .withCompressionCodec(parquetConfig.getCompressionCodecName())
+          .withRowGroupSize(parquetConfig.getBlockSize())
+          .withPageSize(parquetConfig.getPageSize())
+          .withDictionaryPageSize(parquetConfig.getPageSize())
+          .withDictionaryEncoding(parquetConfig.isDictionaryEnabled())
+          .withWriterVersion(ParquetWriter.DEFAULT_WRITER_VERSION)
+          
.withConf(parquetConfig.getStorageConf().unwrapAs(Configuration.class))
+          .build();
+    } catch (IOException | RuntimeException e) {
+      closeQuietly(outputStream);
+      throw e;
+    }
   }
 
   @Override
diff --git 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetBinaryCopyBase.java
 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetBinaryCopyBase.java
index d12268569ffb..e01e2f6cd6f8 100644
--- 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetBinaryCopyBase.java
+++ 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetBinaryCopyBase.java
@@ -20,6 +20,7 @@ package org.apache.hudi.parquet.io;
 
 import org.apache.hudi.common.model.HoodieRecord;
 import org.apache.hudi.common.model.MetaFieldsMode;
+import org.apache.hudi.common.util.CloseableUtils;
 import org.apache.hudi.common.util.VisibleForTesting;
 import org.apache.hudi.exception.HoodieException;
 
@@ -159,6 +160,7 @@ public abstract class HoodieParquetBinaryCopyBase 
implements Closeable {
       writer.start();
       log.info("init writer ");
     } catch (Exception e) {
+      closeParquetFileWriterQuietly(e);
       log.error("failed to init parquet writer", e);
       throw new HoodieException(e);
     }
@@ -166,17 +168,47 @@ public abstract class HoodieParquetBinaryCopyBase 
implements Closeable {
 
   @Override
   public void close() throws IOException {
-    Map<String, String> extraMetaData = finalizeMetadata();
-    extraMetaData = extraMetaData == null ? new HashMap<>() : extraMetaData;
-    extraMetaData.remove("parquet.avro.schema");
-    extraMetaData.remove("org.apache.spark.sql.parquet.row.metadata");
-    writer.end(extraMetaData);
-    // Release the buffer
-    reusableBlockBuffer = null;
+    if (writer == null) {
+      return;
+    }
+    try {
+      Map<String, String> extraMetaData = finalizeMetadata();
+      extraMetaData = extraMetaData == null ? new HashMap<>() : extraMetaData;
+      extraMetaData.remove("parquet.avro.schema");
+      extraMetaData.remove("org.apache.spark.sql.parquet.row.metadata");
+      writer.end(extraMetaData);
+    } catch (IOException | RuntimeException e) {
+      closeParquetFileWriterQuietly(e);
+      throw e;
+    } finally {
+      writer = null;
+      // Release the buffer
+      reusableBlockBuffer = null;
+    }
   }
 
   protected abstract Map<String, String> finalizeMetadata();
 
+  private void closeParquetFileWriterQuietly(Throwable failure) {
+    // Parquet 1.12.x/1.13.x have no close() API; a failed start()/end() can 
leave the stream open.
+    // Newer versions implement AutoCloseable, allowing explicit cleanup here.
+    ParquetFileWriter parquetFileWriter = writer;
+    writer = null;
+    if (parquetFileWriter instanceof AutoCloseable) {
+      CloseableUtils.closeSuppressing((AutoCloseable) parquetFileWriter, 
failure);
+    }
+  }
+
+  @VisibleForTesting
+  void setWriter(ParquetFileWriter writer) {
+    this.writer = writer;
+  }
+
+  @VisibleForTesting
+  ParquetFileWriter getWriter() {
+    return writer;
+  }
+
   public void 
processBlocksFromReader(CompressionConverter.TransParquetFileReader reader, 
BlockMetaData block, String originalCreatedBy) throws IOException {
 
     totalRecordsWritten += block.getRowCount();
@@ -492,22 +524,20 @@ public abstract class HoodieParquetBinaryCopyBase 
implements Closeable {
         new ColumnChunkPageWriteStore(compressor, newSchema, 
props.getAllocator(), props.getColumnIndexTruncateLength(), 
props.getPageWriteChecksumEnabled(), null, numBlocksRewritten);
     ColumnWriteStore cStore = props.newColumnWriteStore(newSchema, cPageStore);
     ColumnWriter cWriter = cStore.getColumnWriter(descriptor);
+    try (Closeable columnWriter = cWriter::close; Closeable columnStore = 
cStore::close) {
+      // For masked column, we assume it's present (DL = max) and not repeated 
(RL = 0)
+      // This is valid for _hoodie_file_name which is a top-level field.
+      int rlvl = 0;
+      int dlvl = descriptor.getMaxDefinitionLevel();
 
-    // For masked column, we assume it's present (DL = max) and not repeated 
(RL = 0)
-    // This is valid for _hoodie_file_name which is a top-level field.
-    int rlvl = 0;
-    int dlvl = descriptor.getMaxDefinitionLevel();
+      for (int i = 0; i < totalChunkValues; i++) {
+        cWriter.write(maskValue, rlvl, dlvl);
+        cStore.endRecord();
+      }
 
-    for (int i = 0; i < totalChunkValues; i++) {
-      cWriter.write(maskValue, rlvl, dlvl);
-      cStore.endRecord();
+      cStore.flush();
+      cPageStore.flushToFileWriter(writer);
     }
-
-    cStore.flush();
-    cPageStore.flushToFileWriter(writer);
-
-    cStore.close();
-    cWriter.close();
   }
 
   private void addNullColumn(ColumnDescriptor descriptor, long 
totalChunkValues, EncodingStats encodingStats, ParquetFileWriter writer, 
MessageType schema, CompressionCodecName newCodecName)
@@ -524,32 +554,32 @@ public abstract class HoodieParquetBinaryCopyBase 
implements Closeable {
         new ColumnChunkPageWriteStore(compressor, newSchema, 
props.getAllocator(), props.getColumnIndexTruncateLength(), 
props.getPageWriteChecksumEnabled(), null, numBlocksRewritten);
     ColumnWriteStore cStore = props.newColumnWriteStore(newSchema, cPageStore);
     ColumnWriter cWriter = cStore.getColumnWriter(descriptor);
-    int dMax = descriptor.getMaxDefinitionLevel();
-
-    for (int i = 0; i < totalChunkValues; i++) {
-      int rlvl = 0;
-      int dlvl = 0;
-      if (dlvl == dMax) {
-        // since we checked ether optional or repeated, dlvl should be > 0
-        if (dlvl == 0) {
-          throw new IOException("definition level is detected to be 0 for 
column " + Arrays.stream(descriptor.getPath()).collect(Collectors.joining(".")) 
+ " to be nullified");
-        }
-        // we just write one null for the whole list at the top level,
-        // instead of nullify the elements in the list one by one
-        if (rlvl == 0) {
-          cWriter.writeNull(rlvl, dlvl - 1);
+    try (Closeable columnWriter = cWriter::close; Closeable columnStore = 
cStore::close) {
+      int dMax = descriptor.getMaxDefinitionLevel();
+
+      for (int i = 0; i < totalChunkValues; i++) {
+        int rlvl = 0;
+        int dlvl = 0;
+        if (dlvl == dMax) {
+          // since we checked ether optional or repeated, dlvl should be > 0
+          if (dlvl == 0) {
+            throw new IOException("definition level is detected to be 0 for 
column " + Arrays.stream(descriptor.getPath()).collect(Collectors.joining(".")) 
+ " to be nullified");
+          }
+          // we just write one null for the whole list at the top level,
+          // instead of nullify the elements in the list one by one
+          if (rlvl == 0) {
+            cWriter.writeNull(rlvl, dlvl - 1);
+          }
+        } else {
+          // A zero repetition level starts a new record; a zero definition 
level marks the field absent.
+          cWriter.writeNull(rlvl, dlvl);
         }
-      } else {
-        cWriter.writeNull(rlvl, dlvl); // 因为repeatition 
level没有重复所以后面都是以0在第一层,definition level是字段path的第0层
+        cStore.endRecord();
       }
-      cStore.endRecord();
-    }
-
-    cStore.flush();
-    cPageStore.flushToFileWriter(writer);
 
-    cStore.close();
-    cWriter.close();
+      cStore.flush();
+      cPageStore.flushToFileWriter(writer);
+    }
   }
 
   private List<ColumnDescriptor> missedColumns(MessageType requiredSchema, 
MessageType fileSchema) {
diff --git 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetFileBinaryCopier.java
 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetFileBinaryCopier.java
index 3589ec3fd0de..84eb474bfa96 100644
--- 
a/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetFileBinaryCopier.java
+++ 
b/hudi-hadoop-common/src/main/java/org/apache/hudi/parquet/io/HoodieParquetFileBinaryCopier.java
@@ -18,6 +18,7 @@
 
 package org.apache.hudi.parquet.io;
 
+import org.apache.hudi.common.util.VisibleForTesting;
 import org.apache.hudi.core.io.storage.HoodieFileBinaryCopier;
 import org.apache.hudi.core.io.storage.HoodieFileMetadataMerger;
 import org.apache.hudi.storage.StoragePath;
@@ -202,6 +203,16 @@ public class HoodieParquetFileBinaryCopier extends 
HoodieParquetBinaryCopyBase i
     this.prefetchExecutor = Executors.newSingleThreadExecutor();
   }
 
+  @VisibleForTesting
+  ExecutorService getPrefetchExecutor() {
+    return prefetchExecutor;
+  }
+
+  @VisibleForTesting
+  void setReader(CompressionConverter.TransParquetFileReader reader) {
+    this.reader = reader;
+  }
+
   @Override
   protected Map<String, String> finalizeMetadata() {
     return this.extraMetaData;
@@ -248,13 +259,17 @@ public class HoodieParquetFileBinaryCopier extends 
HoodieParquetBinaryCopyBase i
 
   @Override
   public void close() throws IOException {
-    super.close();
-    if (prefetchExecutor != null) {
-      prefetchExecutor.shutdownNow();
+    try (CompressionConverter.TransParquetFileReader input = reader) {
+      super.close();
+    } finally {
+      reader = null;
+      if (prefetchExecutor != null) {
+        prefetchExecutor.shutdownNow();
+      }
+      // Release buffers
+      currentBuffer = null;
+      nextBuffer = null;
     }
-    // Release buffers
-    currentBuffer = null;
-    nextBuffer = null;
   }
 
   // Queue input files to be processed
diff --git 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/io/hadoop/TestHoodieHFileReaderWriter.java
 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/io/hadoop/TestHoodieHFileReaderWriter.java
index 6f61cb5b24b0..d9d9605743e6 100644
--- 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/io/hadoop/TestHoodieHFileReaderWriter.java
+++ 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/io/hadoop/TestHoodieHFileReaderWriter.java
@@ -44,6 +44,7 @@ import 
org.apache.hudi.core.io.storage.HoodieNativeAvroHFileReader;
 import org.apache.hudi.exception.MetadataNotFoundException;
 import org.apache.hudi.hadoop.fs.HadoopFSUtils;
 import org.apache.hudi.io.hfile.HFileReader;
+import org.apache.hudi.io.hfile.HFileWriterImpl;
 import org.apache.hudi.io.hfile.UTF8StringKey;
 import org.apache.hudi.io.storage.hadoop.HoodieAvroHFileWriter;
 import org.apache.hudi.io.util.FileIOUtils;
@@ -62,8 +63,10 @@ import org.junit.jupiter.params.provider.Arguments;
 import org.junit.jupiter.params.provider.CsvSource;
 import org.junit.jupiter.params.provider.MethodSource;
 import org.junit.jupiter.params.provider.ValueSource;
+import org.mockito.MockedConstruction;
 
 import java.io.IOException;
+import java.io.OutputStream;
 import java.nio.ByteBuffer;
 import java.util.ArrayList;
 import java.util.Arrays;
@@ -98,12 +101,19 @@ import static 
org.apache.hudi.io.hfile.TestHFileReader.BOOTSTRAP_INDEX_HFILE_SUF
 import static 
org.apache.hudi.io.hfile.TestHFileReader.COMPLEX_SCHEMA_HFILE_SUFFIX;
 import static 
org.apache.hudi.io.hfile.TestHFileReader.SIMPLE_SCHEMA_HFILE_SUFFIX;
 import static org.apache.hudi.io.hfile.TestHFileReader.readHFileFromResources;
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertNull;
+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.doThrow;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockConstruction;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
@@ -204,6 +214,30 @@ public class TestHoodieHFileReaderWriter extends 
TestHoodieReaderWriterBase {
     }
   }
 
+  @ParameterizedTest
+  @ValueSource(booleans = {false, true})
+  public void testConstructorClosesWriterWhenFileInfoFails(boolean closeFails) 
throws Exception {
+    HoodieSchema schema = 
getSchemaFromResource(TestHoodieOrcReaderWriter.class, 
"/exampleSchemaWithMetaFields.avsc");
+    RuntimeException failure = new IllegalStateException("file info failed");
+    IOException closeFailure = new IOException("close failed");
+    try (MockedConstruction<HFileWriterImpl> construction = 
mockConstruction(HFileWriterImpl.class, (writer, context) -> {
+      doThrow(failure).when(writer).appendFileInfo(eq(SCHEMA_KEY), 
any(byte[].class));
+      OutputStream outputStream = (OutputStream) context.arguments().get(1);
+      doAnswer(invocation -> {
+        outputStream.close();
+        if (closeFails) {
+          throw closeFailure;
+        }
+        return null;
+      }).when(writer).close();
+    })) {
+      assertSame(failure, assertThrows(IllegalStateException.class, () -> 
createWriter(schema, false)));
+      assertEquals(1, construction.constructed().size());
+      verify(construction.constructed().get(0)).close();
+      assertArrayEquals(closeFails ? new Throwable[] {closeFailure} : new 
Throwable[0], failure.getSuppressed());
+    }
+  }
+
   @Test
   public void testReadFooterMetadata() throws Exception {
     HoodieSchema schema = 
getSchemaFromResource(TestHoodieOrcReaderWriter.class, 
"/exampleSchemaWithMetaFields.avsc");
diff --git 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestHoodieParquetBinaryCopyBaseSchemaEvolution.java
 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestHoodieParquetBinaryCopyBaseSchemaEvolution.java
index b556889af5d4..dab1d2411a8a 100644
--- 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestHoodieParquetBinaryCopyBaseSchemaEvolution.java
+++ 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestHoodieParquetBinaryCopyBaseSchemaEvolution.java
@@ -365,4 +365,4 @@ public class TestHoodieParquetBinaryCopyBaseSchemaEvolution 
{
       return new HashMap<>();
     }
   }
-}
\ No newline at end of file
+}
diff --git 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestHoodieParquetBinaryCopyLifecycle.java
 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestHoodieParquetBinaryCopyLifecycle.java
new file mode 100644
index 000000000000..6dc2b9353594
--- /dev/null
+++ 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestHoodieParquetBinaryCopyLifecycle.java
@@ -0,0 +1,118 @@
+/*
+ * 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.parquet.io;
+
+import org.apache.hudi.core.io.storage.HoodieFileMetadataMerger;
+
+import org.apache.hadoop.conf.Configuration;
+import org.apache.parquet.hadoop.ParquetFileWriter;
+import org.apache.parquet.hadoop.metadata.CompressionCodecName;
+import org.apache.parquet.hadoop.util.CompressionConverter;
+import org.junit.jupiter.api.Test;
+
+import java.io.IOException;
+import java.util.Collections;
+import java.util.Map;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertNull;
+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.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+
+class TestHoodieParquetBinaryCopyLifecycle {
+  @Test
+  void testCloseClearsWriterWhenEndFails() throws Exception {
+    TestCopy copy = new TestCopy();
+    ParquetFileWriter writer = mock(ParquetFileWriter.class);
+    copy.setWriter(writer);
+    IOException failure = new IOException("end failed");
+    doThrow(failure).when(writer).end(any());
+    assertSame(failure, assertThrows(IOException.class, copy::close));
+    assertNull(copy.getWriter());
+    assertDoesNotThrow(copy::close);
+    verify(writer).end(any());
+  }
+
+  @Test
+  void testClosePreservesMetadataFailure() throws Exception {
+    TestCopy copy = new TestCopy();
+    ParquetFileWriter writer = mock(ParquetFileWriter.class);
+    copy.setWriter(writer);
+    copy.metadataFailure = new IllegalStateException("metadata failed");
+    IOException closeFailure = new IOException("close failed");
+    if (writer instanceof AutoCloseable) {
+      doThrow(closeFailure).when((AutoCloseable) writer).close();
+    }
+    assertSame(copy.metadataFailure, assertThrows(IllegalStateException.class, 
copy::close));
+    assertArrayEquals(writer instanceof AutoCloseable ? new Throwable[] 
{closeFailure} : new Throwable[0],
+        copy.metadataFailure.getSuppressed());
+    assertNull(copy.getWriter());
+  }
+
+  @Test
+  void testFailedFinalizationClosesReaderAndExecutor() throws Exception {
+    HoodieParquetFileBinaryCopier copy = new HoodieParquetFileBinaryCopier(
+        new Configuration(), CompressionCodecName.UNCOMPRESSED, new 
HoodieFileMetadataMerger());
+    ExecutorService executor = copy.getPrefetchExecutor();
+    executor.submit(() -> { }).get();
+    CompressionConverter.TransParquetFileReader reader = 
mock(CompressionConverter.TransParquetFileReader.class);
+    IOException closeFailure = new IOException("reader close failed");
+    doThrow(closeFailure).when(reader).close();
+    copy.setReader(reader);
+    ParquetFileWriter writer = mock(ParquetFileWriter.class);
+    IOException failure = new IOException("end failed");
+    doThrow(failure).when(writer).end(any());
+    copy.setWriter(writer);
+    try {
+      assertSame(failure, assertThrows(IOException.class, copy::close));
+      assertArrayEquals(new Throwable[] {closeFailure}, 
failure.getSuppressed());
+      assertTrue(executor.isShutdown());
+      assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+      assertDoesNotThrow(copy::close);
+      verify(reader).close();
+      verify(writer).end(any());
+    } finally {
+      executor.shutdownNow();
+    }
+  }
+
+  private static class TestCopy extends HoodieParquetBinaryCopyBase {
+    private RuntimeException metadataFailure;
+
+    private TestCopy() {
+      super(new Configuration());
+    }
+
+    @Override
+    protected Map<String, String> finalizeMetadata() {
+      if (metadataFailure != null) {
+        throw metadataFailure;
+      }
+      return Collections.emptyMap();
+    }
+  }
+}
diff --git 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestHoodieParquetFileBinaryCopierPrefetch.java
 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestHoodieParquetFileBinaryCopierPrefetch.java
index be65c8e77bfb..c3ce9b315a24 100644
--- 
a/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestHoodieParquetFileBinaryCopierPrefetch.java
+++ 
b/hudi-hadoop-common/src/test/java/org/apache/hudi/parquet/io/TestHoodieParquetFileBinaryCopierPrefetch.java
@@ -38,12 +38,15 @@ import java.io.IOException;
 import java.nio.file.Files;
 import java.nio.file.Paths;
 import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.Collections;
 import java.util.List;
+import java.util.concurrent.ExecutorService;
 import java.util.stream.Collectors;
 
 import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 
 public class TestHoodieParquetFileBinaryCopierPrefetch {
@@ -74,6 +77,24 @@ public class TestHoodieParquetFileBinaryCopierPrefetch {
     }
   }
 
+  @Test
+  public void testCloseAfterPrefetchFailureReleasesResources() throws 
Exception {
+    TestHoodieParquetFileBinaryCopier.TestFile first = new 
TestHoodieParquetFileBinaryCopier.TestFileBuilder(conf, schema)
+        .withNumRecord(5).build();
+    inputFiles.add(first);
+    ExecutorService executor;
+    try (HoodieParquetFileBinaryCopier copier = new 
HoodieParquetFileBinaryCopier(
+        conf, CompressionCodecName.UNCOMPRESSED, new 
HoodieFileMetadataMerger())) {
+      executor = copier.getPrefetchExecutor();
+      assertThrows(IOException.class, () -> copier.binaryCopy(
+          Arrays.asList(new StoragePath(first.getFileName()), new 
StoragePath(first.getFileName() + ".missing")),
+          Collections.singletonList(new StoragePath(outputFile)), schema, 
false));
+      assertEquals(5, copier.totalRecordsWritten);
+      assertTrue(Files.exists(Paths.get(outputFile)));
+    }
+    assertTrue(executor.isShutdown());
+  }
+
   @Test
   public void testPrefetchingMultipleFiles() throws IOException {
     // Create multiple small files
diff --git 
a/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/SparkHelpers.scala
 
b/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/SparkHelpers.scala
index 4b4c39bc14ed..fee135de700f 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/SparkHelpers.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/main/scala/org/apache/spark/sql/hudi/SparkHelpers.scala
@@ -24,7 +24,7 @@ import org.apache.hudi.common.config.{HoodieParquetConfig, 
HoodieStorageConfig}
 import 
org.apache.hudi.common.config.HoodieStorageConfig.{BLOOM_FILTER_DYNAMIC_MAX_ENTRIES,
 BLOOM_FILTER_FPP_VALUE, BLOOM_FILTER_NUM_ENTRIES_VALUE, BLOOM_FILTER_TYPE}
 import org.apache.hudi.common.model.{HoodieFileFormat, HoodieRecord}
 import org.apache.hudi.common.schema.HoodieSchema
-import org.apache.hudi.common.util.Option
+import org.apache.hudi.common.util.{CloseableUtils, Option}
 import org.apache.hudi.core.io.storage.HoodieIOFactory
 import org.apache.hudi.io.storage.hadoop.HoodieAvroParquetWriter
 import org.apache.hudi.storage.{HoodieStorage, StorageConfiguration, 
StoragePath}
@@ -72,14 +72,20 @@ object SparkHelpers {
     conf.unwrap().setClassLoader(Thread.currentThread.getContextClassLoader)
 
     val writer = new HoodieAvroParquetWriter(destinationFile, parquetConfig, 
instantTime, new SparkTaskContextSupplier(), true)
-    for (rec <- sourceRecords) {
-      val key: String = 
rec.get(HoodieRecord.RECORD_KEY_METADATA_FIELD).toString
-      if (!keysToSkip.contains(key)) {
+    try {
+      for (rec <- sourceRecords) {
+        val key: String = 
rec.get(HoodieRecord.RECORD_KEY_METADATA_FIELD).toString
+        if (!keysToSkip.contains(key)) {
 
-        writer.writeAvro(key, rec)
+          writer.writeAvro(key, rec)
+        }
       }
+    } catch {
+      case t: Throwable =>
+        CloseableUtils.closeSuppressing(writer, t)
+        throw t
     }
-    writer.close
+    writer.close()
   }
 }
 

Reply via email to