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()
}
}