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 513edcf8c5e0 fix(common): close LSM readers when initialization fails
(#20030)
513edcf8c5e0 is described below
commit 513edcf8c5e0a811699d2d9db2645252cd158199
Author: Shuo Cheng <[email protected]>
AuthorDate: Thu Sep 24 18:17:06 2026 +0800
fix(common): close LSM readers when initialization fails (#20030)
* fix(common): close all LSM readers even when one fails to close
---------
Co-authored-by: danny0405 <[email protected]>
---
.../table/read/lsm/LsmFileGroupRecordIterator.java | 67 ++++---
.../read/lsm/TestLsmFileGroupRecordIterator.java | 196 ++++++++++++++++++---
2 files changed, 219 insertions(+), 44 deletions(-)
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/LsmFileGroupRecordIterator.java
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/LsmFileGroupRecordIterator.java
index 4f3805405061..c9031b9af3b5 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/LsmFileGroupRecordIterator.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/LsmFileGroupRecordIterator.java
@@ -38,6 +38,7 @@ import org.apache.hudi.common.table.read.HoodieReadStats;
import org.apache.hudi.common.table.read.InputSplit;
import org.apache.hudi.common.table.read.ReaderParameters;
import org.apache.hudi.common.table.read.UpdateProcessor;
+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.common.util.VisibleForTesting;
@@ -153,28 +154,36 @@ public class LsmFileGroupRecordIterator<T> implements
ClosableIterator<BufferedR
*/
private List<SortedRunReader<T>> initializeReaders() throws IOException {
List<SortedRunReader<T>> sortedRunReaders = new ArrayList<>();
- int mergeOrder = 0;
- if (readBaseFile) {
- addReader(sortedRunReaders, mergeOrder++,
LsmFileIterators.createBaseFileIterator(
- readerContext, storage, inputSplit.getBaseFileOption().get(),
- inputSplit.getStart(), inputSplit.getLength(), orderingFieldNames,
false));
- }
+ try {
+ int mergeOrder = 0;
+ if (readBaseFile) {
+ addReader(sortedRunReaders, mergeOrder++,
LsmFileIterators.createBaseFileIterator(
+ readerContext, storage, inputSplit.getBaseFileOption().get(),
+ inputSplit.getStart(), inputSplit.getLength(), orderingFieldNames,
false));
+ }
- if (inputSplit.hasRecordIterator()) {
- addReader(sortedRunReaders, mergeOrder++,
createRecordIterator(inputSplit.getRecordIterator()));
- }
+ if (inputSplit.hasRecordIterator()) {
+ addReader(sortedRunReaders, mergeOrder++,
createRecordIterator(inputSplit.getRecordIterator()));
+ }
- List<LogReaderSpec> logReaderSpecs = new ArrayList<>();
- if (!inputSplit.hasRecordIterator()) {
- for (HoodieLogFile logFile : inputSplit.getLogFiles()) {
- logReaderSpecs.add(new LogReaderSpec(mergeOrder++, logFile));
+ List<LogReaderSpec> logReaderSpecs = new ArrayList<>();
+ if (!inputSplit.hasRecordIterator()) {
+ for (HoodieLogFile logFile : inputSplit.getLogFiles()) {
+ logReaderSpecs.add(new LogReaderSpec(mergeOrder++, logFile));
+ }
}
- }
- Set<Integer> directLogMergeOrders =
selectDirectLogMergeOrders(logReaderSpecs, readBaseFile);
- for (LogReaderSpec spec : logReaderSpecs) {
- ClosableIterator<BufferedRecord<T>> iterator =
LsmFileIterators.createLogFileIterator(
- readerContext, metaClient, storage, spec.logFile,
orderingFieldNames);
- addReader(sortedRunReaders, spec.mergeOrder,
maybeSpillIterator(directLogMergeOrders.contains(spec.mergeOrder), iterator));
+ Set<Integer> directLogMergeOrders =
selectDirectLogMergeOrders(logReaderSpecs, readBaseFile);
+ for (LogReaderSpec spec : logReaderSpecs) {
+ ClosableIterator<BufferedRecord<T>> iterator =
LsmFileIterators.createLogFileIterator(
+ readerContext, metaClient, storage, spec.logFile,
orderingFieldNames);
+ addReader(sortedRunReaders, spec.mergeOrder,
maybeSpillIterator(directLogMergeOrders.contains(spec.mergeOrder), iterator));
+ }
+ } catch (Throwable e) {
+ // Construction failed, so the caller cannot close the readers already
opened.
+ for (SortedRunReader<T> reader : sortedRunReaders) {
+ CloseableUtils.closeSuppressing(reader::close, e);
+ }
+ throw e;
}
return sortedRunReaders;
}
@@ -244,7 +253,14 @@ public class LsmFileGroupRecordIterator<T> implements
ClosableIterator<BufferedR
*/
private void addReader(List<SortedRunReader<T>> sortedRunReaders, int
mergeOrder, ClosableIterator<BufferedRecord<T>> iterator) {
SortedRunReader<T> sortedRunReader = new SortedRunReader<>(mergeOrder,
iterator);
- if (sortedRunReader.advance()) {
+ boolean hasRecord;
+ try {
+ hasRecord = sortedRunReader.advance();
+ } catch (Throwable failure) {
+ CloseableUtils.closeSuppressing(sortedRunReader::close, failure);
+ throw failure;
+ }
+ if (hasRecord) {
sortedRunReaders.add(sortedRunReader);
} else {
sortedRunReader.close();
@@ -486,7 +502,16 @@ public class LsmFileGroupRecordIterator<T> implements
ClosableIterator<BufferedR
}
private void close() {
- leaves.forEach(SortedRunReader::close);
+ Iterator<SortedRunReader<T>> iterator = leaves.iterator();
+ try {
+ while (iterator.hasNext()) {
+ iterator.next().close();
+ }
+ } catch (Throwable e) {
+ // Attempt to close every remaining reader before propagating the
first failure.
+ iterator.forEachRemaining(reader ->
CloseableUtils.closeSuppressing(reader::close, e));
+ throw e;
+ }
}
private static int nextPowerOfTwo(int value) {
diff --git
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestLsmFileGroupRecordIterator.java
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestLsmFileGroupRecordIterator.java
index 8090c4859601..0a1b8747a5f3 100644
---
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestLsmFileGroupRecordIterator.java
+++
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestLsmFileGroupRecordIterator.java
@@ -56,8 +56,12 @@ import org.apache.avro.generic.GenericData;
import org.apache.avro.generic.IndexedRecord;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.CsvSource;
+import org.junit.jupiter.params.provider.ValueSource;
import java.io.IOException;
+import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.Arrays;
@@ -71,6 +75,7 @@ import java.util.stream.Stream;
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;
@@ -78,7 +83,10 @@ import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -216,29 +224,8 @@ class TestLsmFileGroupRecordIterator {
@Test
void testSortedRunFixturesMergeAcrossBaseAndSpilledLogs() throws IOException
{
HoodieTableConfig tableConfig = new HoodieTableConfig();
- tableConfig.setValue(HoodieTableConfig.RECORDKEY_FIELDS, "id");
- tableConfig.setValue(HoodieTableConfig.ORDERING_FIELDS, "ts");
- tableConfig.setValue(HoodieTableConfig.RECORD_MERGE_MODE,
RecordMergeMode.EVENT_TIME_ORDERING.name());
- tableConfig.setValue(HoodieTableConfig.BASE_FILE_FORMAT,
HoodieFileFormat.PARQUET.name());
- tableConfig.setValue(HoodieTableConfig.META_FIELDS_MODE, "NONE");
- StorageConfiguration<?> storageConfiguration =
mock(StorageConfiguration.class);
- HoodieAvroReaderContext context = org.mockito.Mockito.spy(
- new HoodieAvroReaderContext(storageConfiguration, tableConfig,
Option.empty(), Option.empty()));
- TypedProperties props = new TypedProperties();
- props.setProperty(HoodieReaderConfig.LSM_SORT_MERGE_SPILL_THRESHOLD.key(),
"1");
- props.setProperty(HoodieMemoryConfig.SPILLABLE_MAP_BASE_PATH.key(),
tempDir.toString());
- context.initRecordMerger(props);
-
- FileGroupReaderSchemaHandler<IndexedRecord> schemaHandler =
mock(FileGroupReaderSchemaHandler.class);
- when(schemaHandler.getRequiredSchema()).thenReturn(tableSchema());
- when(schemaHandler.getRequestedSchema()).thenReturn(tableSchema());
- when(schemaHandler.getTableSchema()).thenReturn(tableSchema());
-
when(schemaHandler.getInternalSchema()).thenReturn(InternalSchema.getEmptyInternalSchema());
- when(schemaHandler.getDeleteContext()).thenReturn(new DeleteContext(props,
tableSchema()));
-
when(schemaHandler.getRequiredSchemaForFileAndRenamedColumns(any(StoragePath.class)))
- .thenReturn(Pair.of(tableSchema(), Collections.emptyMap()));
- context.setSchemaHandler(schemaHandler);
- context.setTablePath("/tmp/lsm-fixture");
+ TypedProperties props = readerProperties(1);
+ HoodieAvroReaderContext context = createReaderContext(tableConfig, props);
doAnswer(invocation -> {
StoragePathInfo pathInfo = invocation.getArgument(0);
@@ -303,6 +290,169 @@ class TestLsmFileGroupRecordIterator {
assertEquals(Arrays.asList("a:log2-a", "c:stale-c"), logOnly);
}
+ @ParameterizedTest
+ @ValueSource(booleans = {false, true})
+ void
testInitializationClosesReadersAndDeletesSpillsWhenOpeningLaterLogFails(boolean
closeFails) throws IOException {
+ HoodieTableConfig tableConfig = new HoodieTableConfig();
+ TypedProperties props = readerProperties(1);
+ HoodieAvroReaderContext context = createReaderContext(tableConfig, props);
+ ClosableIterator<IndexedRecord> baseIterator = spy(ClosableIterator.wrap(
+ Collections.singletonList(indexedRecord("a", "base", 1)).iterator()));
+ ClosableIterator<IndexedRecord> logIterator = spy(ClosableIterator.wrap(
+ Collections.singletonList(indexedRecord("a", "log", 2)).iterator()));
+ IOException openFailure = new IOException("Cannot open later log");
+ RuntimeException closeFailure = new RuntimeException("Cannot close base
reader");
+ if (closeFails) {
+ doThrow(closeFailure).when(baseIterator).close();
+ }
+ doAnswer(invocation -> {
+ String name = ((StoragePathInfo)
invocation.getArgument(0)).getPath().getName();
+ if (!name.contains(".log.")) {
+ return baseIterator;
+ }
+ if (name.contains("_001_")) {
+ return logIterator;
+ }
+ // Confirm the previous log actually spilled before failing the next
open.
+ assertEquals(1, spillFileCount());
+ throw openFailure;
+ }).when(context).getFileRecordIterator(any(StoragePathInfo.class),
anyLong(), anyLong(),
+ any(HoodieSchema.class), any(HoodieSchema.class),
any(HoodieStorage.class));
+
+ IOException thrown = assertThrows(IOException.class, () ->
createIterator(context, tableConfig, props));
+
+ assertSame(openFailure, thrown);
+ assertEquals(closeFails ? 1 : 0, thrown.getSuppressed().length);
+ if (closeFails) {
+ assertSame(closeFailure, thrown.getSuppressed()[0]);
+ }
+ verify(baseIterator).close();
+ verify(logIterator).close();
+ assertEquals(0, spillFileCount());
+ }
+
+ @ParameterizedTest
+ @CsvSource({"false, false", "false, true", "true, false", "true, true"})
+ void testInitializationClosesReadersWhenReadingLaterLogFails(boolean
failOnNext, boolean spill) throws IOException {
+ HoodieTableConfig tableConfig = new HoodieTableConfig();
+ TypedProperties props = readerProperties(spill ? 1 : 3);
+ HoodieAvroReaderContext context = createReaderContext(tableConfig, props);
+ ClosableIterator<IndexedRecord> baseIterator = spy(ClosableIterator.wrap(
+ Collections.singletonList(indexedRecord("a", "base", 1)).iterator()));
+ ClosableIterator<IndexedRecord> logIterator = spy(ClosableIterator.wrap(
+ Collections.singletonList(indexedRecord("a", "log", 2)).iterator()));
+ ClosableIterator<IndexedRecord> failingIterator =
mock(ClosableIterator.class);
+ RuntimeException readFailure = new RuntimeException("Cannot read later
log");
+ if (failOnNext) {
+ when(failingIterator.hasNext()).thenReturn(true);
+ when(failingIterator.next()).thenThrow(readFailure);
+ } else {
+ when(failingIterator.hasNext()).thenThrow(readFailure);
+ }
+ doReturn(baseIterator, logIterator,
failingIterator).when(context).getFileRecordIterator(
+ any(StoragePathInfo.class), anyLong(), anyLong(),
any(HoodieSchema.class), any(HoodieSchema.class), any(HoodieStorage.class));
+
+ RuntimeException thrown = assertThrows(RuntimeException.class, () ->
createIterator(context, tableConfig, props));
+
+ assertSame(readFailure, thrown);
+ verify(baseIterator).close();
+ verify(logIterator).close();
+ verify(failingIterator).close();
+ assertEquals(0, spillFileCount());
+ }
+
+ @ParameterizedTest
+ @CsvSource({"false, false", "false, true", "true, false", "true, true"})
+ void testCloseContinuesAfterReaderFailure(boolean baseFails, boolean error)
throws IOException {
+ HoodieTableConfig tableConfig = new HoodieTableConfig();
+ // Keep the base and first log direct, and spill the last log to exercise
file cleanup.
+ TypedProperties props = readerProperties(2);
+ HoodieAvroReaderContext context = createReaderContext(tableConfig, props);
+ ClosableIterator<IndexedRecord> baseIterator = spy(ClosableIterator.wrap(
+ Collections.singletonList(indexedRecord("a", "base", 1)).iterator()));
+ ClosableIterator<IndexedRecord> logIterator = spy(ClosableIterator.wrap(
+ Collections.singletonList(indexedRecord("a", "log", 2)).iterator()));
+ ClosableIterator<IndexedRecord> spilledSourceIterator =
spy(ClosableIterator.wrap(
+ Collections.singletonList(indexedRecord("a", "spilled",
3)).iterator()));
+ doReturn(baseIterator, logIterator,
spilledSourceIterator).when(context).getFileRecordIterator(
+ any(StoragePathInfo.class), anyLong(), anyLong(),
any(HoodieSchema.class), any(HoodieSchema.class), any(HoodieStorage.class));
+ Throwable firstFailure = error ? new AssertionError("reader close failed")
: new RuntimeException("reader close failed");
+ AssertionError secondFailure = new AssertionError("another reader close
failed");
+ doThrow(firstFailure).when(baseFails ? baseIterator : logIterator).close();
+ if (baseFails) {
+ doThrow(secondFailure).when(logIterator).close();
+ }
+ LsmFileGroupRecordIterator<IndexedRecord> iterator =
createIterator(context, tableConfig, props);
+ assertEquals(1, spillFileCount());
+
+ Throwable thrown = assertThrows(firstFailure.getClass(), iterator::close);
+
+ assertSame(firstFailure, thrown);
+ assertEquals(baseFails ? 1 : 0, thrown.getSuppressed().length);
+ if (baseFails) {
+ assertSame(secondFailure, thrown.getSuppressed()[0]);
+ }
+ verify(baseIterator).close();
+ verify(logIterator).close();
+ verify(spilledSourceIterator).close();
+ assertEquals(0, spillFileCount());
+ }
+
+ private TypedProperties readerProperties(int spillThreshold) {
+ TypedProperties props = new TypedProperties();
+ props.setProperty(HoodieReaderConfig.LSM_SORT_MERGE_SPILL_THRESHOLD.key(),
String.valueOf(spillThreshold));
+ props.setProperty(HoodieMemoryConfig.SPILLABLE_MAP_BASE_PATH.key(),
tempDir.toString());
+ return props;
+ }
+
+ private static HoodieAvroReaderContext createReaderContext(HoodieTableConfig
tableConfig, TypedProperties props) {
+ tableConfig.setValue(HoodieTableConfig.RECORDKEY_FIELDS, "id");
+ tableConfig.setValue(HoodieTableConfig.ORDERING_FIELDS, "ts");
+ tableConfig.setValue(HoodieTableConfig.RECORD_MERGE_MODE,
RecordMergeMode.EVENT_TIME_ORDERING.name());
+ tableConfig.setValue(HoodieTableConfig.BASE_FILE_FORMAT,
HoodieFileFormat.PARQUET.name());
+ tableConfig.setValue(HoodieTableConfig.META_FIELDS_MODE, "NONE");
+ StorageConfiguration<?> storageConfiguration =
mock(StorageConfiguration.class);
+ HoodieAvroReaderContext context = org.mockito.Mockito.spy(
+ new HoodieAvroReaderContext(storageConfiguration, tableConfig,
Option.empty(), Option.empty()));
+ context.initRecordMerger(props);
+
+ FileGroupReaderSchemaHandler<IndexedRecord> schemaHandler =
mock(FileGroupReaderSchemaHandler.class);
+ when(schemaHandler.getRequiredSchema()).thenReturn(tableSchema());
+ when(schemaHandler.getRequestedSchema()).thenReturn(tableSchema());
+ when(schemaHandler.getTableSchema()).thenReturn(tableSchema());
+
when(schemaHandler.getInternalSchema()).thenReturn(InternalSchema.getEmptyInternalSchema());
+ when(schemaHandler.getDeleteContext()).thenReturn(new DeleteContext(props,
tableSchema()));
+
when(schemaHandler.getRequiredSchemaForFileAndRenamedColumns(any(StoragePath.class)))
+ .thenReturn(Pair.of(tableSchema(), Collections.emptyMap()));
+ context.setSchemaHandler(schemaHandler);
+ context.setTablePath("/tmp/lsm-fixture");
+
+ return context;
+ }
+
+ private LsmFileGroupRecordIterator<IndexedRecord> createIterator(
+ HoodieAvroReaderContext context, HoodieTableConfig tableConfig,
TypedProperties props) throws IOException {
+ InputSplit split = InputSplit.builder()
+ .baseFileOption(Option.of(new
HoodieBaseFile(pathInfo("/tmp/file1_1-0-1_001.parquet", 100))))
+ .logFileStream(Stream.of(
+ new HoodieLogFile(pathInfo("/tmp/file1_1-0-1_001_1.log.parquet",
50)),
+ new HoodieLogFile(pathInfo("/tmp/file1_1-0-1_002_1.log.parquet",
60))))
+ .partitionPath("")
+ .start(0)
+ .length(Long.MAX_VALUE)
+ .build();
+ HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+ when(metaClient.getTableConfig()).thenReturn(tableConfig);
+ return new LsmFileGroupRecordIterator<>(context,
mock(HoodieStorage.class), split, Collections.singletonList("ts"),
+ metaClient, props,
ReaderParameters.builder().emitDeletes(false).build(), new HoodieReadStats(),
Option.empty());
+ }
+
+ private long spillFileCount() throws IOException {
+ try (Stream<Path> files = Files.list(tempDir)) {
+ return files.count();
+ }
+ }
+
private static LsmFileGroupRecordIterator.SortedRunReader<String>
sortedRunReader(int mergeOrder,
BufferedRecord<String>... records) {
LsmFileGroupRecordIterator.SortedRunReader<String> reader = new
LsmFileGroupRecordIterator.SortedRunReader<>(