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 d78a7b5baa42 fix(flink): avoid reusing split reader functions across 
fetchers (#19315)
d78a7b5baa42 is described below

commit d78a7b5baa42eb5cc4d493ed2101e0af35a441e2
Author: Shuo Cheng <[email protected]>
AuthorDate: Tue Jul 21 09:20:16 2026 +0800

    fix(flink): avoid reusing split reader functions across fetchers (#19315)
---
 .../java/org/apache/hudi/source/HoodieSource.java  | 12 ++--
 .../hudi/source/reader/HoodieSourceReader.java     |  5 +-
 .../source/reader/HoodieSourceSplitReader.java     |  7 +-
 .../org/apache/hudi/table/HoodieTableSource.java   | 27 +++-----
 .../org/apache/hudi/source/TestHoodieSource.java   | 18 ++---
 .../source/reader/TestHoodieSourceSplitReader.java | 80 ++++++++++++++--------
 6 files changed, 87 insertions(+), 62 deletions(-)

diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/HoodieSource.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/HoodieSource.java
index a42dbc1e62b5..7126dfd290ef 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/HoodieSource.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/HoodieSource.java
@@ -18,6 +18,7 @@
 
 package org.apache.hudi.source;
 
+import org.apache.hudi.common.function.SerializableSupplier;
 import org.apache.hudi.common.model.HoodieTableType;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.util.Option;
@@ -73,7 +74,7 @@ import java.util.stream.Collectors;
 @Slf4j
 public class HoodieSource<T> extends FileIndexReader implements Source<T, 
HoodieSourceSplit, HoodieSplitEnumeratorState> {
   private final HoodieScanContext scanContext;
-  private final SplitReaderFunction<T> readerFunction;
+  private final SerializableSupplier<SplitReaderFunction<T>> 
readerFunctionSupplier;
   private final SerializableComparator<HoodieSourceSplit> splitComparator;
   private final HoodieTableMetaClient metaClient;
   private final HoodieRecordEmitter<T> recordEmitter;
@@ -81,18 +82,18 @@ public class HoodieSource<T> extends FileIndexReader 
implements Source<T, Hoodie
 
   public HoodieSource(
       HoodieScanContext scanContext,
-      SplitReaderFunction<T> readerFunction,
+      SerializableSupplier<SplitReaderFunction<T>> readerFunctionSupplier,
       SerializableComparator<HoodieSourceSplit> splitComparator,
       HoodieTableMetaClient metaClient,
       HoodieRecordEmitter<T> recordEmitter) {
     ValidationUtils.checkArgument(scanContext != null, "scanContext can't be 
null.");
-    ValidationUtils.checkArgument(readerFunction != null, "readerFunction 
can't be null.");
+    ValidationUtils.checkArgument(readerFunctionSupplier != null, 
"readerFunctionSupplier can't be null.");
     ValidationUtils.checkArgument(splitComparator != null, "splitComparator 
can't be null.");
     ValidationUtils.checkArgument(metaClient != null, "metaClient can't be 
null.");
     ValidationUtils.checkArgument(recordEmitter != null, "recordEmitter can't 
be null.");
 
     this.scanContext = scanContext;
-    this.readerFunction = readerFunction;
+    this.readerFunctionSupplier = readerFunctionSupplier;
     this.splitComparator = splitComparator;
     this.metaClient = metaClient;
     this.recordEmitter = recordEmitter;
@@ -127,7 +128,8 @@ public class HoodieSource<T> extends FileIndexReader 
implements Source<T, Hoodie
 
   @Override
   public SourceReader<T, HoodieSourceSplit> createReader(SourceReaderContext 
readerContext) throws Exception {
-    return new HoodieSourceReader<T>(tableName, recordEmitter, scanContext, 
readerContext, readerFunction, splitComparator);
+    return new HoodieSourceReader<T>(
+        tableName, recordEmitter, scanContext, readerContext, 
readerFunctionSupplier, splitComparator);
   }
 
   private SplitEnumerator<HoodieSourceSplit, HoodieSplitEnumeratorState> 
createEnumerator(
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/reader/HoodieSourceReader.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/reader/HoodieSourceReader.java
index 7bc3c87be0a8..294039597fb7 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/reader/HoodieSourceReader.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/reader/HoodieSourceReader.java
@@ -18,6 +18,7 @@
 
 package org.apache.hudi.source.reader;
 
+import org.apache.hudi.common.function.SerializableSupplier;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.source.HoodieScanContext;
 import org.apache.hudi.source.reader.function.SplitReaderFunction;
@@ -45,9 +46,9 @@ public class HoodieSourceReader<T> extends
           RecordEmitter<HoodieRecordWithPosition<T>, T, HoodieSourceSplit> 
recordEmitter,
           HoodieScanContext scanContext,
           SourceReaderContext context,
-          SplitReaderFunction<T> readerFunction,
+          SerializableSupplier<SplitReaderFunction<T>> readerFunctionSupplier,
           SerializableComparator<HoodieSourceSplit> splitComparator) {
-    super(() -> new HoodieSourceSplitReader<>(tableName, context, 
readerFunction, splitComparator,
+    super(() -> new HoodieSourceSplitReader<>(tableName, context, 
readerFunctionSupplier, splitComparator,
             scanContext.getLimit() == RecordLimiter.NO_LIMIT ? Option.empty() 
: Option.of(new RecordLimiter(scanContext.getLimit()))),
         recordEmitter, scanContext.getConf(), context);
   }
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/reader/HoodieSourceSplitReader.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/reader/HoodieSourceSplitReader.java
index e53c5f82eed9..c12c4309df83 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/reader/HoodieSourceSplitReader.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/source/reader/HoodieSourceSplitReader.java
@@ -18,6 +18,7 @@
 
 package org.apache.hudi.source.reader;
 
+import org.apache.hudi.common.function.SerializableSupplier;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.metrics.FlinkStreamReadMetrics;
 import org.apache.hudi.source.reader.function.SplitReaderFunction;
@@ -74,12 +75,14 @@ public class HoodieSourceSplitReader<T> implements 
SplitReader<HoodieRecordWithP
   public HoodieSourceSplitReader(
       String tableName,
       SourceReaderContext context,
-      SplitReaderFunction<T> readerFunction,
+      SerializableSupplier<SplitReaderFunction<T>> readerFunctionSupplier,
       SerializableComparator<HoodieSourceSplit> splitComparator,
       Option<RecordLimiter> recordLimiter) {
     this.splitComparator = splitComparator;
     this.splits = new ArrayDeque<>();
-    this.readerFunction = readerFunction;
+    // Flink can start a new fetcher before the previous idle fetcher finishes 
closing. Supply a
+    // fresh stateful cursor for each split reader so the old fetcher's close 
cannot affect the new one.
+    this.readerFunction = readerFunctionSupplier.get();
     this.recordLimiter = recordLimiter;
     this.readerMetrics = new FlinkStreamReadMetrics(context.metricGroup(), 
tableName);
     this.readerMetrics.registerMetrics();
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/HoodieTableSource.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/HoodieTableSource.java
index c3ba129eaa27..a32014a3c67f 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/HoodieTableSource.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/HoodieTableSource.java
@@ -20,6 +20,7 @@ package org.apache.hudi.table;
 
 import org.apache.hudi.adapter.DataStreamScanProviderAdapter;
 import org.apache.hudi.adapter.InputFormatSourceFunctionAdapter;
+import org.apache.hudi.common.function.SerializableSupplier;
 import org.apache.hudi.common.model.FileSlice;
 import org.apache.hudi.common.model.HoodieCommitMetadata;
 import org.apache.hudi.common.model.HoodieTableType;
@@ -304,7 +305,7 @@ public class HoodieTableSource extends FileIndexReader 
implements
 
     HoodieScanContext context = createHoodieScanContext(rowType);
     final HoodieTableType tableType = 
HoodieTableType.valueOf(this.conf.get(FlinkOptions.TABLE_TYPE));
-    final SplitReaderFunction<RowData> splitReaderFunction;
+    final SerializableSupplier<SplitReaderFunction<RowData>> 
splitReaderFunctionSupplier;
     final MergeOnReadTableState<HoodieSourceSplit> hoodieTableState = new 
MergeOnReadTableState<>(
             rowType,
             requiredRowType,
@@ -314,24 +315,16 @@ public class HoodieTableSource extends FileIndexReader 
implements
     boolean emitDelete = tableType == HoodieTableType.MERGE_ON_READ && 
context.isStreaming();
     if (conf.get(FlinkOptions.CDC_ENABLED)) {
       List<DataType> fieldTypes = rowDataType.getChildren();
-      splitReaderFunction = new HoodieCdcSplitReaderFunction(
-          conf,
-          hoodieTableState,
-          internalSchemaManager,
-          fieldTypes,
-          predicates,
-          emitDelete);
+      splitReaderFunctionSupplier = () -> new HoodieCdcSplitReaderFunction(
+          conf, hoodieTableState, internalSchemaManager, fieldTypes, 
predicates, emitDelete);
     } else {
-      splitReaderFunction = new HoodieSplitReaderFunction(
-          conf,
-          tableSchema,
-          requiredHoodieSchema,
-          internalSchemaManager,
-          conf.get(FlinkOptions.MERGE_TYPE),
-          predicates,
-          emitDelete);
+      splitReaderFunctionSupplier = () -> new HoodieSplitReaderFunction(
+          conf, tableSchema, requiredHoodieSchema, internalSchemaManager,
+          conf.get(FlinkOptions.MERGE_TYPE), predicates, emitDelete);
     }
-    return new HoodieSource<>(context, splitReaderFunction, new 
HoodieSourceSplitComparator(), metaClient, new HoodieRecordEmitter<>());
+    return new HoodieSource<>(
+        context, splitReaderFunctionSupplier, new 
HoodieSourceSplitComparator(), metaClient,
+        new HoodieRecordEmitter<>());
   }
 
   /**
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestHoodieSource.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestHoodieSource.java
index 5714484d491f..1f975d078f21 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestHoodieSource.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestHoodieSource.java
@@ -465,18 +465,18 @@ public class TestHoodieSource {
         .build();
     HoodieSchema schema = HoodieSchemaConverter.convertToSchema(rowType);
     HadoopStorageConfiguration hadoopConf = new 
HadoopStorageConfiguration(HadoopConfigurations.getHadoopConf(conf));
-    HoodieSplitReaderFunction splitReaderFunction = new 
HoodieSplitReaderFunction(
-        conf,
-        schema, // schema will be resolved from table
-        schema, // required schema
-        InternalSchemaManager.get(hadoopConf, this.metaClient),
-        conf.get(FlinkOptions.MERGE_TYPE),
-        Collections.emptyList(),
-            false);
+    InternalSchemaManager internalSchemaManager = 
InternalSchemaManager.get(hadoopConf, this.metaClient);
 
     return new HoodieSource<>(
         scanContext,
-        splitReaderFunction,
+        () -> new HoodieSplitReaderFunction(
+            conf,
+            schema, // schema will be resolved from table
+            schema, // required schema
+            internalSchemaManager,
+            conf.get(FlinkOptions.MERGE_TYPE),
+            Collections.emptyList(),
+            false),
         new HoodieSourceSplitComparator(),
         metaClient,
         new HoodieRecordEmitter<>());
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/reader/TestHoodieSourceSplitReader.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/reader/TestHoodieSourceSplitReader.java
index ccb6c2bbf893..11e2cd940ea7 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/reader/TestHoodieSourceSplitReader.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/reader/TestHoodieSourceSplitReader.java
@@ -18,6 +18,7 @@
 
 package org.apache.hudi.source.reader;
 
+import org.apache.hudi.common.function.SerializableSupplier;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.source.reader.function.SplitReaderFunction;
 import org.apache.hudi.source.split.HoodieSourceSplit;
@@ -45,6 +46,7 @@ import java.util.stream.Collectors;
 import java.util.stream.IntStream;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -67,7 +69,7 @@ public class TestHoodieSourceSplitReader {
   public void testFetchWithNoSplits() throws IOException {
     TestSplitReaderFunction readerFunction = new TestSplitReaderFunction();
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
 
     RecordsWithSplitIds<HoodieRecordWithPosition<String>> result = 
reader.fetch();
 
@@ -81,7 +83,7 @@ public class TestHoodieSourceSplitReader {
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
 
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
 
     HoodieSourceSplit split = createTestSplit(1, "file1");
     SplitsAddition<HoodieSourceSplit> splitsChange = new 
SplitsAddition<>(Collections.singletonList(split));
@@ -99,7 +101,7 @@ public class TestHoodieSourceSplitReader {
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
 
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
 
     HoodieSourceSplit split1 = createTestSplit(1, "file1");
     HoodieSourceSplit split2 = createTestSplit(2, "file2");
@@ -127,7 +129,7 @@ public class TestHoodieSourceSplitReader {
         (s1, s2) -> s2.getFileId().compareTo(s1.getFileId());
 
     HoodieSourceSplitReader<String> reader =
-            new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, comparator, Option.empty());
+            new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, comparator, Option.empty());
 
     HoodieSourceSplit split1 = createTestSplit(1, "file1");
     HoodieSourceSplit split2 = createTestSplit(2, "file2");
@@ -150,7 +152,7 @@ public class TestHoodieSourceSplitReader {
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
 
     HoodieSourceSplitReader<String> reader =
-            new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+            new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
 
     // First batch
     HoodieSourceSplit split1 = createTestSplit(1, "file1");
@@ -172,7 +174,7 @@ public class TestHoodieSourceSplitReader {
   public void testClose() throws Exception {
     TestSplitReaderFunction readerFunction = new TestSplitReaderFunction();
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
 
     HoodieSourceSplit split = createTestSplit(1, "file1");
     reader.handleSplitsChanges(new 
SplitsAddition<>(Collections.singletonList(split)));
@@ -191,7 +193,7 @@ public class TestHoodieSourceSplitReader {
   public void testWakeUp() throws IOException {
     TestSplitReaderFunction readerFunction = new TestSplitReaderFunction();
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
 
     // wakeUp() now sets a flag but must not throw, and the flag is reset at 
the top of each fetch()
     // so a wakeUp with no in-flight drain leaves the next fetch() unaffected.
@@ -213,7 +215,7 @@ public class TestHoodieSourceSplitReader {
     List<String> testData = Arrays.asList("r1", "r2", "r3", "r4", "r5");
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
     boolean[] fired = {false};
     readerFunction.setDrainProbe((buffered, hasNext) -> {
       if (buffered == 1 && !fired[0]) {
@@ -260,7 +262,7 @@ public class TestHoodieSourceSplitReader {
     List<String> testData = Arrays.asList("r1", "r2", "r3");
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
     boolean[] fired = {false};
     readerFunction.setDrainProbe((buffered, hasNext) -> {
       // Wake at the start of the drain (count 0) while data is still 
available.
@@ -292,7 +294,7 @@ public class TestHoodieSourceSplitReader {
     List<String> testData = Arrays.asList("r1", "r2");
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
     boolean[] fired = {false};
     readerFunction.setDrainProbe((buffered, hasNext) -> {
       // Wake only at the start of a drain that finds the cursor already 
exhausted.
@@ -340,7 +342,7 @@ public class TestHoodieSourceSplitReader {
       }
     };
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.of(wakingLimiter));
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.of(wakingLimiter));
     holder.set(reader);
 
     HoodieSourceSplit split = createTestSplit(1, "file1");
@@ -364,7 +366,7 @@ public class TestHoodieSourceSplitReader {
     List<String> testData = Arrays.asList("r1", "r2", "r3", "r4", "r5");
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.of(new RecordLimiter(3L)));
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.of(new RecordLimiter(3L)));
     boolean[] fired = {false};
     readerFunction.setDrainProbe((buffered, hasNext) -> {
       if (buffered == 1 && !fired[0]) {
@@ -397,7 +399,7 @@ public class TestHoodieSourceSplitReader {
     List<String> testData = Arrays.asList("r1", "r2", "r3");
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
     boolean[] fired = {false};
     readerFunction.setDrainProbe((buffered, hasNext) -> {
       if (buffered == 1 && !fired[0]) {
@@ -423,7 +425,7 @@ public class TestHoodieSourceSplitReader {
   public void testPauseOrResumeSplits() {
     TestSplitReaderFunction readerFunction = new TestSplitReaderFunction();
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
 
     HoodieSourceSplit split1 = createTestSplit(1, "file1");
     HoodieSourceSplit split2 = createTestSplit(2, "file2");
@@ -441,7 +443,7 @@ public class TestHoodieSourceSplitReader {
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
 
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
 
     HoodieSourceSplit split = createTestSplit(1, "file1");
     reader.handleSplitsChanges(new 
SplitsAddition<>(Collections.singletonList(split)));
@@ -456,18 +458,42 @@ public class TestHoodieSourceSplitReader {
   public void testReaderFunctionClosedOnReaderClose() throws Exception {
     TestSplitReaderFunction readerFunction = new TestSplitReaderFunction();
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
 
     reader.close();
 
     assertTrue(readerFunction.isClosed(), "Reader function should be closed");
   }
 
+  @Test
+  public void testEachSplitReaderUsesIndependentReaderFunction() throws 
Exception {
+    List<TestSplitReaderFunction> readerFunctions = new ArrayList<>();
+    SerializableSupplier<SplitReaderFunction<String>> readerFunctionSupplier = 
() -> {
+      TestSplitReaderFunction readerFunction = new TestSplitReaderFunction();
+      readerFunctions.add(readerFunction);
+      return readerFunction;
+    };
+
+    HoodieSourceSplitReader<String> firstReader = new 
HoodieSourceSplitReader<>(
+        TABLE_NAME, readerContext, readerFunctionSupplier, null, 
Option.empty());
+    HoodieSourceSplitReader<String> secondReader = new 
HoodieSourceSplitReader<>(
+        TABLE_NAME, readerContext, readerFunctionSupplier, null, 
Option.empty());
+
+    assertEquals(2, readerFunctions.size());
+    firstReader.close();
+    assertTrue(readerFunctions.get(0).isClosed());
+    assertFalse(readerFunctions.get(1).isClosed(),
+        "Closing an idle fetcher's split reader must not close the next 
fetcher's reader function");
+
+    secondReader.close();
+    assertTrue(readerFunctions.get(1).isClosed());
+  }
+
   @Test
   public void testFetchEmptyResultWhenNoSplitsAdded() throws IOException {
     TestSplitReaderFunction readerFunction = new TestSplitReaderFunction();
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
 
     RecordsWithSplitIds<HoodieRecordWithPosition<String>> result = 
reader.fetch();
 
@@ -483,7 +509,7 @@ public class TestHoodieSourceSplitReader {
 
     // No comparator - should preserve insertion order
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
 
     HoodieSourceSplit split3 = createTestSplit(3, "file3");
     HoodieSourceSplit split1 = createTestSplit(1, "file1");
@@ -505,7 +531,7 @@ public class TestHoodieSourceSplitReader {
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
 
     HoodieSourceSplitReader<String> reader =
-            new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+            new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
 
     HoodieSourceSplit split1 = createTestSplit(1, "file1");
     HoodieSourceSplit split2 = createTestSplit(2, "file2");
@@ -527,7 +553,7 @@ public class TestHoodieSourceSplitReader {
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
 
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.of(new RecordLimiter(2L)));
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.of(new RecordLimiter(2L)));
 
     HoodieSourceSplit split = createTestSplit(1, "file1");
     reader.handleSplitsChanges(new 
SplitsAddition<>(Collections.singletonList(split)));
@@ -555,7 +581,7 @@ public class TestHoodieSourceSplitReader {
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
 
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.of(new RecordLimiter(2L)));
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.of(new RecordLimiter(2L)));
 
     HoodieSourceSplit split1 = createTestSplit(1, "file1");
     HoodieSourceSplit split2 = createTestSplit(2, "file2");
@@ -583,7 +609,7 @@ public class TestHoodieSourceSplitReader {
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(Arrays.asList("r1", "r2"));
 
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.of(new RecordLimiter(0L)));
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.of(new RecordLimiter(0L)));
 
     HoodieSourceSplit split1 = createTestSplit(1, "file1");
     HoodieSourceSplit split2 = createTestSplit(2, "file2");
@@ -606,7 +632,7 @@ public class TestHoodieSourceSplitReader {
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
 
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.of(new RecordLimiter(3L)));
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.of(new RecordLimiter(3L)));
 
     HoodieSourceSplit split = createTestSplit(1, "file1");
     reader.handleSplitsChanges(new 
SplitsAddition<>(Collections.singletonList(split)));
@@ -632,7 +658,7 @@ public class TestHoodieSourceSplitReader {
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
 
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.of(new RecordLimiter(5L)));
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.of(new RecordLimiter(5L)));
 
     HoodieSourceSplit split1 = createTestSplit(1, "file1");
     HoodieSourceSplit split2 = createTestSplit(2, "file2");
@@ -665,7 +691,7 @@ public class TestHoodieSourceSplitReader {
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
 
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null,
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null,
             Option.empty());
 
     HoodieSourceSplit split = createTestSplit(1, "file1");
@@ -689,7 +715,7 @@ public class TestHoodieSourceSplitReader {
     List<String> testData = IntStream.range(0, n).mapToObj(i -> "r" + 
i).collect(Collectors.toList());
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
 
     HoodieSourceSplit split = createTestSplit(1, "file1");
     reader.handleSplitsChanges(new 
SplitsAddition<>(Collections.singletonList(split)));
@@ -732,7 +758,7 @@ public class TestHoodieSourceSplitReader {
     List<String> testData = Arrays.asList("r1", "r2", "r3", "r4", "r5");
     TestSplitReaderFunction readerFunction = new 
TestSplitReaderFunction(testData);
     HoodieSourceSplitReader<String> reader =
-        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, 
readerFunction, null, Option.empty());
+        new HoodieSourceSplitReader<>(TABLE_NAME, readerContext, () -> 
readerFunction, null, Option.empty());
 
     HoodieSourceSplit split = createTestSplit(1, "file1");
     split.updatePosition(0, 2L); // 2 records already consumed before recovery

Reply via email to