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 46cd13434308 fix(flink): improve read path coverage and CDC
reliability (#19402)
46cd13434308 is described below
commit 46cd134343081c0252e3019bdce0c136e7b8c1d2
Author: Shuo Cheng <[email protected]>
AuthorDate: Thu Jul 30 18:25:18 2026 +0800
fix(flink): improve read path coverage and CDC reliability (#19402)
---
.../hudi/table/format/cdc/CdcImageManager.java | 4 +
.../apache/hudi/table/format/cdc/CdcIterators.java | 16 +-
.../hudi/source/TestExpressionEvaluators.java | 57 ++++
.../org/apache/hudi/source/TestHoodieSource.java | 57 ++++
.../hudi/source/TestIncrementalInputSplits.java | 12 +
.../function/TestHoodieCdcSplitReaderFunction.java | 290 ++++++++++++++++++++-
.../function/TestHoodieSplitReaderFunction.java | 80 ++++++
.../hudi/table/format/TestFilePathUtils.java | 250 ++++++++++++++++++
.../table/format/TestHoodieRowDataLanceReader.java | 219 ++++++++++++++++
.../table/format/TestInternalSchemaManager.java | 204 +++++++++++++++
.../hudi/table/format/TestRecordIterators.java | 276 ++++++++++++++++++++
.../hudi/table/format/cdc/TestCdcImageManager.java | 180 +++++++++++++
.../hudi/table/format/cdc/TestCdcIterators.java | 229 ++++++++++++++++
.../format/cow/TestCopyOnWriteInputFormat.java | 199 ++++++++++++++
14 files changed, 2061 insertions(+), 12 deletions(-)
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcImageManager.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcImageManager.java
index 55a9445dc4c2..2854498c9bda 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcImageManager.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcImageManager.java
@@ -39,6 +39,7 @@ import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
import java.io.DataInputStream;
import java.io.DataOutputStream;
+import java.io.EOFException;
import java.io.IOException;
import java.util.Map;
import java.util.TreeMap;
@@ -170,6 +171,9 @@ public class CdcImageManager implements AutoCloseable {
public void skipBytesToRead(int numBytes) throws IOException {
while (numBytes > 0) {
int skipped = skipBytes(numBytes);
+ if (skipped == 0) {
+ throw new EOFException("Could not skip " + numBytes + " remaining
bytes");
+ }
numBytes -= skipped;
}
}
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcIterators.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcIterators.java
index 104276583ae0..725873329bcc 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcIterators.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/cdc/CdcIterators.java
@@ -124,19 +124,19 @@ public final class CdcIterators {
@Override
public boolean hasNext() {
- if (recordIterator != null) {
- if (recordIterator.hasNext()) {
- return true;
- } else {
+ while (true) {
+ if (recordIterator != null) {
+ if (recordIterator.hasNext()) {
+ return true;
+ }
recordIterator.close();
recordIterator = null;
}
- }
- if (fileSplitIterator.hasNext()) {
+ if (!fileSplitIterator.hasNext()) {
+ return false;
+ }
recordIterator = recordIteratorFunc.apply(fileSplitIterator.next());
- return recordIterator.hasNext();
}
- return false;
}
@Override
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestExpressionEvaluators.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestExpressionEvaluators.java
index 2c077ab8a5e3..4ae717306e22 100644
---
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestExpressionEvaluators.java
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestExpressionEvaluators.java
@@ -28,6 +28,7 @@ import org.apache.flink.table.data.StringData;
import org.apache.flink.table.data.TimestampData;
import org.apache.flink.table.expressions.CallExpression;
import org.apache.flink.table.expressions.FieldReferenceExpression;
+import org.apache.flink.table.expressions.ResolvedExpression;
import org.apache.flink.table.expressions.ValueLiteralExpression;
import org.apache.flink.table.functions.BuiltInFunctionDefinition;
import org.apache.flink.table.functions.BuiltInFunctionDefinitions;
@@ -45,7 +46,9 @@ import java.util.stream.Stream;
import static org.apache.hudi.source.ExpressionEvaluators.fromExpression;
import static org.apache.hudi.source.prune.ColumnStatsProbe.convertColumnStats;
+import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
@@ -393,6 +396,55 @@ public class TestExpressionEvaluators {
}
}
+ @Test
+ void testFromExpressionBuildsCompositeEvaluators() {
+ FieldReferenceExpression ref = new FieldReferenceExpression("f_int",
DataTypes.INT(), 2, 2);
+ ValueLiteralExpression twelve = new ValueLiteralExpression(12);
+ ValueLiteralExpression thirteen = new ValueLiteralExpression(13);
+ CallExpression equals = call(BuiltInFunctionDefinitions.EQUALS, ref,
twelve);
+ CallExpression greaterThan = call(BuiltInFunctionDefinitions.GREATER_THAN,
ref, thirteen);
+ Map<String, ColumnStats> stats = convertColumnStats(intIndexRow(12, 12,
0L), queryFields(2));
+
+ ExpressionEvaluators.Evaluator not =
fromExpression(CallExpression.permanent(
+ BuiltInFunctionDefinitions.NOT, List.of(equals), DataTypes.BOOLEAN()));
+ ExpressionEvaluators.Evaluator and =
fromExpression(CallExpression.permanent(
+ BuiltInFunctionDefinitions.AND, Arrays.asList(equals, greaterThan),
DataTypes.BOOLEAN()));
+ ExpressionEvaluators.Evaluator or =
fromExpression(CallExpression.permanent(
+ BuiltInFunctionDefinitions.OR, Arrays.asList(equals, greaterThan),
DataTypes.BOOLEAN()));
+
+ assertFalse(not.eval(stats));
+ assertFalse(and.eval(stats));
+ assertTrue(or.eval(stats));
+ assertEquals(2, ((ExpressionEvaluators.Or) or).getEvaluators().length);
+ assertEquals(2, fromExpression(Arrays.<ResolvedExpression>asList(equals,
greaterThan)).size());
+ }
+
+ @Test
+ void testFromExpressionHandlesUnaryInAndReversedComparison() {
+ FieldReferenceExpression ref = new FieldReferenceExpression("f_int",
DataTypes.INT(), 2, 2);
+ ValueLiteralExpression eleven = new ValueLiteralExpression(11);
+ ValueLiteralExpression twelve = new ValueLiteralExpression(12);
+ Map<String, ColumnStats> stats = convertColumnStats(intIndexRow(12, 13),
queryFields(2));
+
+ assertInstanceOf(ExpressionEvaluators.IsNull.class,
fromExpression(CallExpression.permanent(
+ BuiltInFunctionDefinitions.IS_NULL, List.of(ref),
DataTypes.BOOLEAN())));
+ assertInstanceOf(ExpressionEvaluators.IsNotNull.class,
fromExpression(CallExpression.permanent(
+ BuiltInFunctionDefinitions.IS_NOT_NULL, List.of(ref),
DataTypes.BOOLEAN())));
+
+ ExpressionEvaluators.Evaluator in =
fromExpression(CallExpression.permanent(
+ BuiltInFunctionDefinitions.IN, Arrays.asList(ref, eleven, twelve),
DataTypes.BOOLEAN()));
+ assertTrue(in.eval(stats));
+
+ assertInstanceOf(ExpressionEvaluators.GreaterThan.class,
fromExpression(CallExpression.permanent(
+ BuiltInFunctionDefinitions.LESS_THAN, Arrays.asList(eleven, ref),
DataTypes.BOOLEAN())));
+ assertInstanceOf(ExpressionEvaluators.LessThan.class,
fromExpression(CallExpression.permanent(
+ BuiltInFunctionDefinitions.GREATER_THAN, Arrays.asList(eleven, ref),
DataTypes.BOOLEAN())));
+ assertInstanceOf(ExpressionEvaluators.GreaterThanOrEqual.class,
fromExpression(CallExpression.permanent(
+ BuiltInFunctionDefinitions.LESS_THAN_OR_EQUAL, Arrays.asList(eleven,
ref), DataTypes.BOOLEAN())));
+ assertInstanceOf(ExpressionEvaluators.LessThanOrEqual.class,
fromExpression(CallExpression.permanent(
+ BuiltInFunctionDefinitions.GREATER_THAN_OR_EQUAL,
Arrays.asList(eleven, ref), DataTypes.BOOLEAN())));
+ }
+
@ParameterizedTest
@MethodSource("twelveObjects")
void testAllNumericDataTypes(Object twelve) {
@@ -410,6 +462,11 @@ public class TestExpressionEvaluators {
return Stream.of((byte) 12, (short) 12, 12, 12L, new BigDecimal(12), 12f,
12d);
}
+ private static CallExpression call(
+ BuiltInFunctionDefinition definition, ResolvedExpression left,
ResolvedExpression right) {
+ return CallExpression.permanent(definition, Arrays.asList(left, right),
DataTypes.BOOLEAN());
+ }
+
private static RowData intIndexRow(Integer minVal, Integer maxVal) {
return intIndexRow(minVal, maxVal, 2L);
}
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 1f975d078f21..ee8f7e23f251 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
@@ -19,21 +19,28 @@
package org.apache.hudi.source;
import org.apache.hudi.common.config.HoodieMetadataConfig;
+import org.apache.hudi.common.function.SerializableSupplier;
import org.apache.hudi.common.model.HoodieTableType;
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.testutils.HoodieTestUtils;
+import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.PartitionPathEncodeUtils;
import org.apache.hudi.configuration.FlinkOptions;
import org.apache.hudi.configuration.HadoopConfigurations;
import org.apache.hudi.index.HoodieIndex;
import org.apache.hudi.index.bucket.BucketIdentifier;
+import org.apache.hudi.source.enumerator.HoodieSplitEnumeratorState;
+import org.apache.hudi.source.enumerator.HoodieStaticSplitEnumerator;
import org.apache.hudi.source.prune.ColumnStatsProbe;
import org.apache.hudi.source.prune.PartitionPruners;
import org.apache.hudi.source.reader.HoodieRecordEmitter;
import org.apache.hudi.source.reader.function.HoodieSplitReaderFunction;
+import org.apache.hudi.source.reader.function.SplitReaderFunction;
import org.apache.hudi.source.split.HoodieSourceSplit;
import org.apache.hudi.source.split.HoodieSourceSplitComparator;
+import org.apache.hudi.source.split.SerializableComparator;
import org.apache.hudi.storage.StoragePath;
import org.apache.hudi.storage.hadoop.HadoopStorageConfiguration;
import org.apache.hudi.table.format.InternalSchemaManager;
@@ -43,6 +50,8 @@ import org.apache.hudi.utils.TestConfigurations;
import org.apache.hudi.utils.TestData;
import org.apache.flink.api.connector.source.Boundedness;
+import org.apache.flink.api.connector.source.SplitEnumerator;
+import org.apache.flink.api.connector.source.SplitEnumeratorContext;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.table.api.DataTypes;
import org.apache.flink.table.data.RowData;
@@ -66,8 +75,12 @@ import java.util.function.Function;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
/**
* Test cases for {@link HoodieSource}.
@@ -418,6 +431,50 @@ public class TestHoodieSource {
assertNotNull(splits, "Incremental splits with pruner should not be null");
}
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testConstructorRejectsNullCollaborators() {
+ HoodieScanContext scanContext = mock(HoodieScanContext.class);
+ SerializableSupplier<SplitReaderFunction<RowData>> readerSupplier =
mock(SerializableSupplier.class);
+ SerializableComparator<HoodieSourceSplit> comparator =
mock(SerializableComparator.class);
+ HoodieTableMetaClient client = mock(HoodieTableMetaClient.class);
+ HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
+ HoodieRecordEmitter<RowData> emitter = mock(HoodieRecordEmitter.class);
+ when(client.getTableConfig()).thenReturn(tableConfig);
+ when(tableConfig.getTableName()).thenReturn("test_table");
+
+ assertThrows(IllegalArgumentException.class,
+ () -> new HoodieSource<>(null, readerSupplier, comparator, client,
emitter));
+ assertThrows(IllegalArgumentException.class,
+ () -> new HoodieSource<>(scanContext, null, comparator, client,
emitter));
+ assertThrows(IllegalArgumentException.class,
+ () -> new HoodieSource<>(scanContext, readerSupplier, null, client,
emitter));
+ assertThrows(IllegalArgumentException.class,
+ () -> new HoodieSource<>(scanContext, readerSupplier, comparator,
null, emitter));
+ assertThrows(IllegalArgumentException.class,
+ () -> new HoodieSource<>(scanContext, readerSupplier, comparator,
client, null));
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testCreateAndRestoreStaticEnumerator() throws Exception {
+ metaClient = HoodieTestUtils.init(tempDir.getAbsolutePath(),
HoodieTableType.COPY_ON_WRITE);
+ conf.set(FlinkOptions.TABLE_TYPE, HoodieTableType.COPY_ON_WRITE.name());
+ HoodieSource<RowData> source = createHoodieSource(conf, metaClient);
+ SplitEnumeratorContext<HoodieSourceSplit> context =
mock(SplitEnumeratorContext.class);
+ when(context.currentParallelism()).thenReturn(1);
+
+ SplitEnumerator<HoodieSourceSplit, HoodieSplitEnumeratorState> created =
+ source.createEnumerator(context);
+ HoodieSplitEnumeratorState state = new HoodieSplitEnumeratorState(
+ Collections.emptyList(), Option.empty(), Option.empty());
+ SplitEnumerator<HoodieSourceSplit, HoodieSplitEnumeratorState> restored =
+ source.restoreEnumerator(context, state);
+
+ assertInstanceOf(HoodieStaticSplitEnumerator.class, created);
+ assertInstanceOf(HoodieStaticSplitEnumerator.class, restored);
+ }
+
// Helper methods
private HoodieSource<RowData> createHoodieSource(Configuration conf,
HoodieTableMetaClient metaClient) {
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestIncrementalInputSplits.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestIncrementalInputSplits.java
index 472c8d74c3c8..9c8e0a7f7db1 100644
---
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestIncrementalInputSplits.java
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/TestIncrementalInputSplits.java
@@ -66,6 +66,7 @@ import org.junit.jupiter.params.provider.EnumSource;
import org.junit.jupiter.params.provider.MethodSource;
import java.io.IOException;
+import java.lang.reflect.Method;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
@@ -811,4 +812,15 @@ public class TestIncrementalInputSplits extends
HoodieCommonTestHarness {
assertNotNull(result.getSplits(), "Batch splits should not be null for
table type: " + tableType);
assertFalse(result.getSplits().isEmpty(), "Batch splits should not be
empty for table type: " + tableType);
}
+
+ @Test
+ @SuppressWarnings("unchecked")
+ void testMergeListHandlesEmptyAndPopulatedInputs() throws Exception {
+ Method mergeList =
IncrementalInputSplits.class.getDeclaredMethod("mergeList", List.class,
List.class);
+ mergeList.setAccessible(true);
+
+ assertEquals(List.of(1), mergeList.invoke(null, Collections.emptyList(),
List.of(1)));
+ assertEquals(List.of(1), mergeList.invoke(null, List.of(1),
Collections.emptyList()));
+ assertEquals(List.of(1, 2), mergeList.invoke(null, List.of(1),
List.of(2)));
+ }
}
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/reader/function/TestHoodieCdcSplitReaderFunction.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/reader/function/TestHoodieCdcSplitReaderFunction.java
index d5dc3eca5770..b4d11e178620 100644
---
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/reader/function/TestHoodieCdcSplitReaderFunction.java
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/reader/function/TestHoodieCdcSplitReaderFunction.java
@@ -18,34 +18,66 @@
package org.apache.hudi.source.reader.function;
+import org.apache.hudi.common.model.FileSlice;
+import org.apache.hudi.common.model.HoodieFileGroupId;
import org.apache.hudi.common.schema.HoodieSchema;
+import org.apache.hudi.common.table.HoodieTableMetaClient;
import org.apache.hudi.common.table.cdc.HoodieCDCFileSplit;
import org.apache.hudi.common.table.cdc.HoodieCDCInferenceCase;
+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.configuration.FlinkOptions;
+import org.apache.hudi.exception.HoodieIOException;
import org.apache.hudi.source.split.HoodieCdcSourceSplit;
import org.apache.hudi.source.split.HoodieSourceSplit;
+import org.apache.hudi.table.format.FormatUtils;
import org.apache.hudi.table.format.InternalSchemaManager;
+import org.apache.hudi.table.format.RecordIterators;
import org.apache.hudi.table.format.mor.MergeOnReadTableState;
+import org.apache.hudi.util.FlinkWriteClients;
import org.apache.hudi.util.HoodieSchemaConverter;
+import org.apache.hudi.util.StreamerUtil;
import org.apache.hudi.utils.TestConfigurations;
import org.apache.flink.configuration.Configuration;
+import org.apache.flink.table.data.GenericRowData;
+import org.apache.flink.table.data.RowData;
+import org.apache.flink.table.types.DataType;
import org.apache.flink.table.types.logical.RowType;
+import org.apache.flink.types.RowKind;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
+import org.mockito.MockedStatic;
import java.io.File;
+import java.io.IOException;
import java.util.ArrayList;
import java.util.Collections;
+import java.util.List;
import static org.apache.hudi.util.StreamerUtil.EMPTY_PARTITION_PATH;
import static org.apache.hudi.utils.TestConfigurations.ROW_DATA_TYPE;
import static org.apache.hudi.utils.TestConfigurations.ROW_TYPE;
-import static org.apache.hudi.utils.TestConfigurations.TABLE_SCHEMA;
+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.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.anyInt;
+import static org.mockito.ArgumentMatchers.anyList;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyMap;
+import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.when;
/**
* Test cases for {@link HoodieCdcSplitReaderFunction}.
@@ -67,7 +99,8 @@ public class TestHoodieCdcSplitReaderFunction {
tableSchema = mock(HoodieSchema.class);
requiredSchema = mock(HoodieSchema.class);
internalSchemaManager = mock(InternalSchemaManager.class);
- tableState = new MergeOnReadTableState(ROW_TYPE, ROW_TYPE,
TABLE_SCHEMA.toString(), TABLE_SCHEMA.toString(), new ArrayList<>());
+ String schema = HoodieSchemaConverter.convertToSchema(ROW_TYPE).toString();
+ tableState = new MergeOnReadTableState(ROW_TYPE, ROW_TYPE, schema, schema,
new ArrayList<>());
}
private HoodieCdcSplitReaderFunction createFunction() {
@@ -77,7 +110,7 @@ public class TestHoodieCdcSplitReaderFunction {
internalSchemaManager,
ROW_DATA_TYPE.getChildren(),
Collections.emptyList(),
- false);
+ false);
}
// -------------------------------------------------------------------------
@@ -98,7 +131,12 @@ public class TestHoodieCdcSplitReaderFunction {
.map(f -> new RowType.RowField(f.getName(), f.getType()))
.collect(java.util.stream.Collectors.toList()));
- tableState = new MergeOnReadTableState(ROW_TYPE, projectedRowType,
TABLE_SCHEMA.toString(),
HoodieSchemaConverter.convertToSchema(projectedRowType.copy()).toString(), new
ArrayList<>());
+ tableState = new MergeOnReadTableState(
+ ROW_TYPE,
+ projectedRowType,
+ HoodieSchemaConverter.convertToSchema(ROW_TYPE).toString(),
+
HoodieSchemaConverter.convertToSchema(projectedRowType.copy()).toString(),
+ new ArrayList<>());
HoodieCdcSplitReaderFunction function = new HoodieCdcSplitReaderFunction(
conf,
tableState,
@@ -241,4 +279,248 @@ public class TestHoodieCdcSplitReaderFunction {
function.open(cdcSplit);
function.close();
}
+
+ @Test
+ public void testConstructorValidationAndProducedRowType() {
+ assertThrows(IllegalArgumentException.class, () -> new
HoodieCdcSplitReaderFunction(
+ conf, null, internalSchemaManager, ROW_DATA_TYPE.getChildren(),
Collections.emptyList(), false));
+ assertThrows(IllegalArgumentException.class, () -> new
HoodieCdcSplitReaderFunction(
+ conf, tableState, null, ROW_DATA_TYPE.getChildren(),
Collections.emptyList(), false));
+ assertEquals(ROW_TYPE, createFunction().producedRowType());
+ }
+
+ @Test
+ public void testBaseFileInsertFromLance() {
+ GenericRowData row = new GenericRowData(ROW_TYPE.getFieldCount());
+ ClosableIterator<RowData> nested =
ClosableIterator.wrap(List.<RowData>of(row).iterator());
+ HoodieCdcSourceSplit split = cdcSplit(new HoodieCDCFileSplit(
+ "20230101000000000", HoodieCDCInferenceCase.BASE_FILE_INSERT,
"insert.lance"));
+
+ try (MockedStatic<FormatUtils> mocked = mockStatic(FormatUtils.class)) {
+ mocked.when(() -> FormatUtils.getLanceRecordIterator(
+ anyString(), anyList(), anyList(), any(int[].class), any()))
+ .thenReturn(nested);
+ try (ClosableIterator<RowData> iterator =
createFunction().createRecordIterator(split)) {
+ assertTrue(iterator.hasNext());
+ assertSame(row, iterator.next());
+ assertEquals(RowKind.INSERT, row.getRowKind());
+ assertFalse(iterator.hasNext());
+ }
+ }
+ }
+
+ @Test
+ public void testBaseFileInsertFromParquet() {
+ GenericRowData row = new GenericRowData(ROW_TYPE.getFieldCount());
+ ClosableIterator<RowData> nested =
ClosableIterator.wrap(List.<RowData>of(row).iterator());
+ HoodieCdcSourceSplit split = cdcSplit(new HoodieCDCFileSplit(
+ "20230101000000000", HoodieCDCInferenceCase.BASE_FILE_INSERT,
"insert.parquet"));
+
+ try (MockedStatic<RecordIterators> mocked =
mockStatic(RecordIterators.class)) {
+ mocked.when(() -> RecordIterators.getParquetRecordIterator(
+ any(), anyBoolean(), anyBoolean(), any(), any(String[].class),
any(DataType[].class),
+ anyMap(), any(int[].class), anyInt(), any(), anyLong(), anyLong(),
anyList()))
+ .thenReturn(nested);
+ try (ClosableIterator<RowData> iterator =
createFunction().createRecordIterator(split)) {
+ assertTrue(iterator.hasNext());
+ assertSame(row, iterator.next());
+ assertFalse(iterator.hasNext());
+ }
+ }
+ }
+
+ @Test
+ public void testBaseFileInsertRequiresExactlyOneFile() {
+ HoodieCDCFileSplit change = new HoodieCDCFileSplit(
+ "20230101000000000", HoodieCDCInferenceCase.BASE_FILE_INSERT,
Collections.emptyList());
+ try (ClosableIterator<RowData> iterator =
createFunction().createRecordIterator(cdcSplit(change))) {
+ assertThrows(IllegalStateException.class, iterator::hasNext);
+ }
+ }
+
+ @Test
+ public void testAsIsBeforeAfterWithNoCdcFilesIsEmpty() {
+ conf.set(FlinkOptions.SUPPLEMENTAL_LOGGING_MODE, "DATA_BEFORE_AFTER");
+ HoodieCDCFileSplit change = new HoodieCDCFileSplit(
+ "20230101000000000", HoodieCDCInferenceCase.AS_IS,
Collections.emptyList());
+ try (ClosableIterator<RowData> iterator =
createFunction().createRecordIterator(cdcSplit(change))) {
+ assertFalse(iterator.hasNext());
+ }
+ }
+
+ @Test
+ public void testBaseFileDeleteReadsBeforeFileSlice() throws Exception {
+ FileSlice before = fileSlice("001");
+ HoodieCDCFileSplit change = new HoodieCDCFileSplit(
+ "20230101000000000", HoodieCDCInferenceCase.BASE_FILE_DELETE,
+ Collections.emptyList(), Option.of(before), Option.empty());
+ GenericRowData row = new GenericRowData(ROW_TYPE.getFieldCount());
+ HoodieRecordReader<RowData> recordReader = mock(HoodieRecordReader.class);
+ when(recordReader.getClosableIterator())
+ .thenReturn(ClosableIterator.wrap(List.<RowData>of(row).iterator()));
+ HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+ HoodieWriteConfig writeConfig = mockWriteConfig();
+
+ try (MockedStatic<FormatUtils> mockedFormatUtils =
mockStatic(FormatUtils.class);
+ MockedStatic<StreamerUtil> mockedStreamerUtil =
mockStatic(StreamerUtil.class)) {
+ mockRecordReader(mockedFormatUtils, recordReader);
+ mockedStreamerUtil.when(() -> StreamerUtil.metaClientForReader(any(),
any()))
+ .thenReturn(metaClient);
+
+ try (ClosableIterator<RowData> iterator =
+
createFunction(writeConfig).createRecordIterator(cdcSplit(change))) {
+ assertTrue(iterator.hasNext());
+ assertEquals(RowKind.DELETE, iterator.next().getRowKind());
+ assertFalse(iterator.hasNext());
+ }
+ }
+ }
+
+ @Test
+ public void testReplaceCommitReadsBeforeFileSlice() throws Exception {
+ FileSlice before = fileSlice("001");
+ HoodieCDCFileSplit change = new HoodieCDCFileSplit(
+ "20230101000000000", HoodieCDCInferenceCase.REPLACE_COMMIT,
+ Collections.emptyList(), Option.of(before), Option.empty());
+ GenericRowData row = new GenericRowData(ROW_TYPE.getFieldCount());
+ HoodieRecordReader<RowData> recordReader = mock(HoodieRecordReader.class);
+ when(recordReader.getClosableIterator())
+ .thenReturn(ClosableIterator.wrap(List.<RowData>of(row).iterator()));
+ HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+ HoodieWriteConfig writeConfig = mockWriteConfig();
+
+ try (MockedStatic<FormatUtils> mockedFormatUtils =
mockStatic(FormatUtils.class);
+ MockedStatic<StreamerUtil> mockedStreamerUtil =
mockStatic(StreamerUtil.class)) {
+ mockRecordReader(mockedFormatUtils, recordReader);
+ mockedStreamerUtil.when(() -> StreamerUtil.metaClientForReader(any(),
any()))
+ .thenReturn(metaClient);
+
+ try (ClosableIterator<RowData> iterator =
+
createFunction(writeConfig).createRecordIterator(cdcSplit(change))) {
+ assertTrue(iterator.hasNext());
+ assertEquals(RowKind.DELETE, iterator.next().getRowKind());
+ assertFalse(iterator.hasNext());
+ }
+ }
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ public void testAsIsModesLoadRequiredImages() throws Exception {
+ HoodieRecordReader<RowData> recordReader = mock(HoodieRecordReader.class);
+ when(recordReader.getClosableIterator()).thenAnswer(
+ invocation ->
ClosableIterator.wrap(Collections.<RowData>emptyList().iterator()));
+ ExternalSpillableMap<String, byte[]> images =
mock(ExternalSpillableMap.class);
+ HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+ HoodieWriteConfig writeConfig = mockWriteConfig();
+ FileSlice before = fileSlice("001");
+ FileSlice after = fileSlice("002");
+
+ try (MockedStatic<FormatUtils> mockedFormatUtils =
mockStatic(FormatUtils.class);
+ MockedStatic<StreamerUtil> mockedStreamerUtil =
mockStatic(StreamerUtil.class)) {
+ mockRecordReader(mockedFormatUtils, recordReader);
+ mockedFormatUtils.when(() -> FormatUtils.spillableMap(any(), anyLong(),
anyString()))
+ .thenReturn(images);
+ mockedStreamerUtil.when(() -> StreamerUtil.metaClientForReader(any(),
any()))
+ .thenReturn(metaClient);
+
+ for (String mode : List.of("DATA_BEFORE", "OP_KEY_ONLY")) {
+ conf.set(FlinkOptions.SUPPLEMENTAL_LOGGING_MODE, mode);
+ HoodieCDCFileSplit change = new HoodieCDCFileSplit(
+ "20230101000000000", HoodieCDCInferenceCase.AS_IS,
+ Collections.emptyList(), Option.of(before), Option.of(after));
+ try (ClosableIterator<RowData> iterator =
+
createFunction(writeConfig).createRecordIterator(cdcSplit(change))) {
+ assertFalse(iterator.hasNext());
+ }
+ }
+ }
+ }
+
+ @Test
+ public void testFileSliceReaderWrapsInitializationFailure() throws Exception
{
+ HoodieRecordReader<RowData> recordReader = mock(HoodieRecordReader.class);
+ when(recordReader.getClosableIterator()).thenThrow(new
IOException("failed"));
+ HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+ HoodieWriteConfig writeConfig = mockWriteConfig();
+ HoodieCDCFileSplit change = new HoodieCDCFileSplit(
+ "20230101000000000", HoodieCDCInferenceCase.BASE_FILE_DELETE,
+ Collections.emptyList(), Option.of(fileSlice("001")), Option.empty());
+
+ try (MockedStatic<FormatUtils> mockedFormatUtils =
mockStatic(FormatUtils.class);
+ MockedStatic<StreamerUtil> mockedStreamerUtil =
mockStatic(StreamerUtil.class)) {
+ mockRecordReader(mockedFormatUtils, recordReader);
+ mockedStreamerUtil.when(() -> StreamerUtil.metaClientForReader(any(),
any()))
+ .thenReturn(metaClient);
+ try (ClosableIterator<RowData> iterator =
+
createFunction(writeConfig).createRecordIterator(cdcSplit(change))) {
+ assertThrows(HoodieIOException.class, iterator::hasNext);
+ }
+ }
+ }
+
+ @Test
+ public void testNonCdcSplitUsesFallbackIterator() throws Exception {
+ HoodieRecordReader<RowData> recordReader = mock(HoodieRecordReader.class);
+ when(recordReader.getClosableIterator()).thenReturn(
+ ClosableIterator.wrap(Collections.<RowData>emptyList().iterator()));
+ HoodieTableMetaClient metaClient = mock(HoodieTableMetaClient.class);
+ HoodieWriteConfig writeConfig = mockWriteConfig();
+ HoodieSourceSplit split = new HoodieSourceSplit(
+ 1, null, Option.of(Collections.emptyList()),
tempDir.getAbsolutePath(), "",
+ "read_optimized", "20230101000000000", "file-1", Option.empty());
+
+ try (MockedStatic<FormatUtils> mockedFormatUtils =
mockStatic(FormatUtils.class);
+ MockedStatic<StreamerUtil> mockedStreamerUtil =
mockStatic(StreamerUtil.class);
+ MockedStatic<FlinkWriteClients> mockedWriteClients =
mockStatic(FlinkWriteClients.class)) {
+ mockRecordReader(mockedFormatUtils, recordReader);
+ mockedStreamerUtil.when(() -> StreamerUtil.metaClientForReader(any(),
any()))
+ .thenReturn(metaClient);
+ mockedWriteClients.when(() ->
FlinkWriteClients.getHoodieClientConfig(any()))
+ .thenReturn(writeConfig);
+ try (ClosableIterator<RowData> iterator =
createFunction().createRecordIterator(split)) {
+ assertFalse(iterator.hasNext());
+ }
+ }
+ }
+
+ private static void mockRecordReader(
+ MockedStatic<FormatUtils> mockedFormatUtils,
+ HoodieRecordReader<RowData> recordReader) {
+ mockedFormatUtils.when(() -> FormatUtils.createRecordReader(
+ any(), any(), any(), any(), any(), any(), anyString(), anyString(),
+ anyBoolean(), anyList(), any())).thenReturn(recordReader);
+ }
+
+ private static FileSlice fileSlice(String instant) {
+ return new FileSlice(
+ new HoodieFileGroupId("partition", "file"), instant, null,
Collections.emptyList());
+ }
+
+ private HoodieWriteConfig mockWriteConfig() {
+ HoodieWriteConfig writeConfig = mock(HoodieWriteConfig.class);
+ when(writeConfig.getBasePath()).thenReturn(tempDir.getAbsolutePath());
+ return writeConfig;
+ }
+
+ private HoodieCdcSplitReaderFunction createFunction(HoodieWriteConfig
writeConfig) {
+ return new HoodieCdcSplitReaderFunction(
+ conf,
+ tableState,
+ internalSchemaManager,
+ ROW_DATA_TYPE.getChildren(),
+ Collections.emptyList(),
+ false) {
+ @Override
+ protected HoodieWriteConfig getWriteConfig() {
+ return writeConfig;
+ }
+ };
+ }
+
+ private HoodieCdcSourceSplit cdcSplit(HoodieCDCFileSplit... changes) {
+ return new HoodieCdcSourceSplit(
+ 1, tempDir.getAbsolutePath(), 128 * 1024 * 1024L, "file-cdc",
+ EMPTY_PARTITION_PATH, changes, "read_optimized", "20230101000000000");
+ }
}
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/reader/function/TestHoodieSplitReaderFunction.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/reader/function/TestHoodieSplitReaderFunction.java
index 3f2d3de07f77..d09eb29d7fbf 100644
---
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/reader/function/TestHoodieSplitReaderFunction.java
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/source/reader/function/TestHoodieSplitReaderFunction.java
@@ -26,10 +26,13 @@ import org.apache.hudi.common.table.HoodieTableMetaClient;
import org.apache.hudi.common.table.read.HoodieFileGroupReader;
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.exception.HoodieIOException;
import org.apache.hudi.source.ExpressionPredicates;
import org.apache.hudi.source.split.HoodieSourceSplit;
+import org.apache.hudi.table.format.FormatUtils;
import org.apache.hudi.table.format.InternalSchemaManager;
+import org.apache.hudi.util.HoodieSchemaConverter;
import org.apache.hudi.util.StreamerUtil;
import org.apache.hudi.utils.TestConfigurations;
@@ -38,6 +41,7 @@ import org.apache.flink.table.data.RowData;
import org.apache.flink.table.expressions.FieldReferenceExpression;
import org.apache.flink.table.expressions.ValueLiteralExpression;
import org.apache.flink.table.types.AtomicDataType;
+import org.apache.flink.table.types.logical.RowType;
import org.apache.flink.table.types.logical.VarCharType;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -54,9 +58,11 @@ import static org.junit.jupiter.api.Assertions.assertNotNull;
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.mock;
import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -514,6 +520,80 @@ public class TestHoodieSplitReaderFunction {
}
}
+ @Test
+ public void testCreateRecordIteratorReturnsInitializedIterator() throws
Exception {
+ HoodieFileGroupReader<RowData> reader = mockReader();
+ @SuppressWarnings("unchecked")
+ ClosableIterator<RowData> iterator = mock(ClosableIterator.class);
+ when(reader.getClosableIterator()).thenReturn(iterator);
+
+ HoodieSplitReaderFunction function = readerFunctionReturning(reader);
+ try (MockedStatic<StreamerUtil> mockedStreamerUtil =
mockStatic(StreamerUtil.class)) {
+ mockedStreamerUtil.when(() -> StreamerUtil.metaClientForReader(any(),
any()))
+ .thenReturn(mockMetaClient);
+ assertSame(iterator, function.createRecordIterator(createSplit()));
+ }
+ verify(reader, never()).close();
+ }
+
+ @Test
+ public void testCreateRecordIteratorClosesReaderOnRuntimeFailure() throws
Exception {
+ HoodieFileGroupReader<RowData> reader = mockReader();
+ IllegalStateException failure = new IllegalStateException("init failed");
+ when(reader.getClosableIterator()).thenThrow(failure);
+
+ HoodieSplitReaderFunction function = readerFunctionReturning(reader);
+ try (MockedStatic<StreamerUtil> mockedStreamerUtil =
mockStatic(StreamerUtil.class)) {
+ mockedStreamerUtil.when(() -> StreamerUtil.metaClientForReader(any(),
any()))
+ .thenReturn(mockMetaClient);
+ assertSame(failure, assertThrows(
+ IllegalStateException.class, () ->
function.createRecordIterator(createSplit())));
+ }
+ verify(reader).close();
+ }
+
+ @Test
+ public void testProducedRowTypeUsesRequiredSchema() {
+ HoodieSchema required =
HoodieSchemaConverter.convertToSchema(TestConfigurations.ROW_TYPE);
+ class ExposedReaderFunction extends HoodieSplitReaderFunction {
+ ExposedReaderFunction() {
+ super(TestHoodieSplitReaderFunction.this.conf,
+ required, required,
+ mockInternalSchemaManager, "AVRO_PAYLOAD",
Collections.emptyList(), false);
+ }
+
+ RowType getProducedRowType() {
+ return producedRowType();
+ }
+ }
+
+ assertEquals(HoodieSchemaConverter.convertToRowType(required),
+ new ExposedReaderFunction().getProducedRowType());
+ }
+
+ @Test
+ public void testCreateRecordReaderBuildsFileSlice() {
+ HoodieRecordReader<RowData> reader = mock(HoodieRecordReader.class);
+ HoodieSplitReaderFunction function = new HoodieSplitReaderFunction(
+ conf,
+ HoodieSchemaConverter.convertToSchema(TestConfigurations.ROW_TYPE),
+ HoodieSchemaConverter.convertToSchema(TestConfigurations.ROW_TYPE),
+ mockInternalSchemaManager,
+ "AVRO_PAYLOAD",
+ Collections.emptyList(),
+ false);
+
+ try (MockedStatic<FormatUtils> mockedFormatUtils =
mockStatic(FormatUtils.class)) {
+ mockedFormatUtils.when(() -> FormatUtils.createRecordReader(
+ any(), any(), any(), any(), any(), any(), any(), any(), anyBoolean(),
+ any(), any())).thenReturn(reader);
+ HoodieSourceSplit split = new HoodieSourceSplit(
+ 1, null, Option.of(Collections.emptyList()), "/tbl", "/part",
+ "read_optimized", "19700101000000000", "file1", Option.empty());
+ assertSame(reader, function.createRecordReader(split, mockMetaClient));
+ }
+ }
+
@SuppressWarnings("unchecked")
private static HoodieFileGroupReader<RowData> mockReader() {
return mock(HoodieFileGroupReader.class);
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestFilePathUtils.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestFilePathUtils.java
new file mode 100644
index 000000000000..54d51e65740f
--- /dev/null
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestFilePathUtils.java
@@ -0,0 +1,250 @@
+/*
+ * 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.table.format;
+
+import org.apache.hudi.configuration.FlinkOptions;
+import org.apache.hudi.storage.StoragePath;
+
+import org.apache.flink.api.java.tuple.Tuple2;
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.table.api.DataTypes;
+import org.apache.flink.table.api.TableException;
+import org.apache.flink.table.types.DataType;
+import org.apache.hadoop.fs.FileStatus;
+import org.apache.hadoop.fs.FileSystem;
+import org.apache.hadoop.fs.Path;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import java.io.IOException;
+import java.nio.file.Files;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Tests {@link FilePathUtils}.
+ */
+class TestFilePathUtils {
+
+ @TempDir
+ java.nio.file.Path tempDir;
+
+ @Test
+ void testGenerateAndUnescapePartitionPath() {
+ LinkedHashMap<String, String> partitionSpec = new LinkedHashMap<>();
+ partitionSpec.put("region", "us/west");
+ partitionSpec.put("day", "2026-07-28");
+
+ assertEquals(
+ "region=us%2Fwest/day=2026-07-28/",
+ FilePathUtils.generatePartitionPath(partitionSpec, true, true));
+ assertEquals(
+ "us%2Fwest/2026-07-28",
+ FilePathUtils.generatePartitionPath(partitionSpec, false, false));
+ assertEquals("", FilePathUtils.generatePartitionPath(new
LinkedHashMap<>(), true, true));
+
+ assertEquals("a/b=c%invalid",
FilePathUtils.unescapePathName("a%2Fb%3Dc%invalid"));
+ assertEquals("trailing%", FilePathUtils.unescapePathName("trailing%"));
+ assertThrows(
+ TableException.class,
+ () -> FilePathUtils.generatePartitionPath(
+ new LinkedHashMap<>(Collections.singletonMap("region", "")), true,
false));
+ }
+
+ @Test
+ void testExtractPartitionKeyValues() {
+ assertEquals(
+ linkedMap("region", "us/west", "day", "2026-07-28"),
+ FilePathUtils.extractPartitionKeyValues(
+ new Path("/table/region=us%2Fwest/day=2026-07-28"),
+ true,
+ new String[] {"region", "day"}));
+ assertEquals(
+ linkedMap("region", "us/west", "day", "2026-07-28"),
+ FilePathUtils.extractPartitionKeyValues(
+ new Path("/table/us%2Fwest/2026-07-28"),
+ false,
+ new String[] {"region", "day"}));
+ assertTrue(FilePathUtils.extractPartitionKeyValues(
+ new Path("/table"), true, new String[0]).isEmpty());
+ }
+
+ @Test
+ void testGeneratePartitionSpecs() {
+ List<String> fieldNames = Arrays.asList("region", "day", "event_time");
+ List<DataType> fieldTypes = Arrays.asList(
+ DataTypes.STRING(), DataTypes.INT(), DataTypes.TIMESTAMP(3));
+
+ assertEquals(
+ objectMap("region", "us", "day", 28),
+ FilePathUtils.generatePartitionSpecs(
+
"/table/region=us/day=28/event_time=2026-07-28%2010%3A15/data.parquet",
+ fieldNames,
+ fieldTypes,
+ FlinkOptions.PARTITION_DEFAULT_NAME.defaultValue(),
+ "region,day,event_time",
+ true));
+
+ assertTrue(FilePathUtils.generatePartitionSpecs(
+ "/table/data.parquet",
+ fieldNames,
+ fieldTypes,
+ FlinkOptions.PARTITION_DEFAULT_NAME.defaultValue(),
+ FlinkOptions.PARTITION_PATH_FIELD.defaultValue(),
+ true).isEmpty());
+ }
+
+ @Test
+ void testRecursivePartitionDiscoveryAndHiddenPaths() throws IOException {
+ Files.createDirectories(tempDir.resolve("region=us/day=28"));
+ Files.createDirectories(tempDir.resolve("region=eu/day=29"));
+ Files.createDirectories(tempDir.resolve("_temporary/day=30"));
+ Files.createDirectories(tempDir.resolve(".hidden/day=31"));
+ Files.createDirectories(tempDir.resolve(".file.log.1/day=32"));
+
+ Path root = new Path(tempDir.toUri());
+ FileSystem fs = root.getFileSystem(new
org.apache.hadoop.conf.Configuration());
+ FileStatus[] statuses = FilePathUtils.getFileStatusRecursively(root, 2,
fs);
+ assertEquals(3, statuses.length);
+
+ List<Tuple2<LinkedHashMap<String, String>, Path>> partitions =
+ FilePathUtils.searchPartKeyValueAndPaths(
+ fs, root, true, new String[] {"region", "day"});
+ assertEquals(3, partitions.size());
+ assertTrue(partitions.stream().anyMatch(
+ tuple -> "us".equals(tuple.f0.get("region")) &&
"28".equals(tuple.f0.get("day"))));
+
+ assertEquals(0, FilePathUtils.getFileStatusRecursively(
+ new Path(root, "missing"), 1, fs).length);
+ }
+
+ @Test
+ void testGetPartitionsResolvesDefaultValue() throws IOException {
+ String defaultPartition =
FlinkOptions.PARTITION_DEFAULT_NAME.defaultValue();
+ Files.createDirectories(tempDir.resolve("region=" + defaultPartition));
+ Files.createDirectories(tempDir.resolve("region=us"));
+
+ List<Map<String, String>> partitions = FilePathUtils.getPartitions(
+ new Path(tempDir.toUri()),
+ new org.apache.hadoop.conf.Configuration(),
+ Collections.singletonList("region"),
+ defaultPartition,
+ true);
+
+ assertEquals(2, partitions.size());
+ assertTrue(partitions.stream().anyMatch(partition ->
partition.containsKey("region")
+ && partition.get("region") == null));
+ assertTrue(partitions.stream().anyMatch(partition ->
"us".equals(partition.get("region"))));
+ }
+
+ @Test
+ void testValidateAndConvertPartitionPaths() {
+ Map<String, String> unordered = new LinkedHashMap<>();
+ unordered.put("day", "28");
+ unordered.put("region", "us");
+
+ assertEquals(
+ linkedMap("region", "us", "day", "28"),
+ FilePathUtils.validateAndReorderPartitions(
+ unordered, Arrays.asList("region", "day")));
+ assertEquals(
+ unordered,
+ FilePathUtils.validateAndReorderPartitions(unordered,
Collections.emptyList()));
+ assertThrows(
+ TableException.class,
+ () -> FilePathUtils.validateAndReorderPartitions(
+ Collections.singletonMap("region", "us"), Arrays.asList("region",
"day")));
+
+ List<Map<String, String>> partitionPaths =
+ Collections.singletonList(unordered);
+ assertArrayEquals(
+ new Path[] {new Path("/table/region=us/day=28/")},
+ FilePathUtils.partitionPath2ReadPath(
+ new Path("/table"), Arrays.asList("region", "day"),
partitionPaths, true));
+ assertEquals(
+ Collections.singleton("us/28"),
+ FilePathUtils.toRelativePartitionPaths(
+ Arrays.asList("region", "day"), partitionPaths, false));
+ }
+
+ @Test
+ void testReadPathsAndPathConversions() throws IOException {
+ Path root = new Path(tempDir.toUri());
+ Configuration flinkConf = new Configuration();
+ org.apache.hadoop.conf.Configuration hadoopConf = new
org.apache.hadoop.conf.Configuration();
+
+ assertArrayEquals(
+ new Path[] {root},
+ FilePathUtils.getReadPaths(
+ root, flinkConf, hadoopConf, Collections.emptyList()));
+
+ Path[] hadoopPaths = {new Path("/table/a"), new Path("/table/b")};
+ org.apache.flink.core.fs.Path[] flinkPaths =
FilePathUtils.toFlinkPaths(hadoopPaths);
+ assertEquals(hadoopPaths[0].toUri(), flinkPaths[0].toUri());
+ assertEquals(
+ new StoragePath("/table/c").toUri(),
+ FilePathUtils.toFlinkPath(new StoragePath("/table/c")).toUri());
+ }
+
+ @Test
+ void testExtractPartitionConfiguration() {
+ Configuration conf = new Configuration();
+ assertArrayEquals(new String[0], FilePathUtils.extractPartitionKeys(conf));
+ assertArrayEquals(new String[0],
FilePathUtils.extractHivePartitionFields(conf));
+
+ conf.set(FlinkOptions.PARTITION_PATH_FIELD, "region,day");
+ assertArrayEquals(
+ new String[] {"region", "day"},
FilePathUtils.extractPartitionKeys(conf));
+ assertArrayEquals(
+ new String[] {"region", "day"},
FilePathUtils.extractHivePartitionFields(conf));
+
+ conf.set(FlinkOptions.HIVE_SYNC_PARTITION_FIELDS, "country,date");
+ assertArrayEquals(
+ new String[] {"country", "date"},
FilePathUtils.extractHivePartitionFields(conf));
+ assertTrue(FilePathUtils.isHiveStylePartitioning("region=us"));
+ assertFalse(FilePathUtils.isHiveStylePartitioning("us"));
+ assertFalse(FilePathUtils.isHiveStylePartitioning("region=us/day=28"));
+ }
+
+ private static LinkedHashMap<String, String> linkedMap(String... entries) {
+ LinkedHashMap<String, String> result = new LinkedHashMap<>();
+ for (int index = 0; index < entries.length; index += 2) {
+ result.put(entries[index], entries[index + 1]);
+ }
+ return result;
+ }
+
+ private static LinkedHashMap<String, Object> objectMap(Object... entries) {
+ LinkedHashMap<String, Object> result = new LinkedHashMap<>();
+ for (int index = 0; index < entries.length; index += 2) {
+ result.put((String) entries[index], entries[index + 1]);
+ }
+ return result;
+ }
+}
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestHoodieRowDataLanceReader.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestHoodieRowDataLanceReader.java
new file mode 100644
index 000000000000..453cdbe4d508
--- /dev/null
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestHoodieRowDataLanceReader.java
@@ -0,0 +1,219 @@
+/*
+ * 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.table.format;
+
+import org.apache.hudi.common.bloom.SimpleBloomFilter;
+import org.apache.hudi.common.config.HoodieConfig;
+import org.apache.hudi.common.schema.HoodieSchema;
+import org.apache.hudi.common.schema.HoodieSchemaUtils;
+import org.apache.hudi.common.schema.internal.InternalSchema;
+import org.apache.hudi.common.util.collection.ClosableIterator;
+import org.apache.hudi.common.util.collection.Pair;
+import org.apache.hudi.exception.HoodieException;
+import org.apache.hudi.exception.HoodieValidationException;
+import org.apache.hudi.storage.StoragePath;
+import org.apache.hudi.util.RowDataQueryContexts;
+
+import org.apache.arrow.memory.BufferAllocator;
+import org.apache.arrow.memory.RootAllocator;
+import org.apache.arrow.vector.VarCharVector;
+import org.apache.arrow.vector.VectorSchemaRoot;
+import org.apache.arrow.vector.ipc.ArrowReader;
+import org.apache.arrow.vector.types.pojo.Schema;
+import org.apache.flink.table.types.DataType;
+import org.apache.flink.table.types.logical.RowType;
+import org.junit.jupiter.api.Test;
+import org.lance.file.LanceFileReader;
+import org.mockito.MockedStatic;
+
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import static
org.apache.hudi.common.avro.HoodieBloomFilterWriteSupport.HOODIE_AVRO_BLOOM_FILTER_METADATA_KEY;
+import static
org.apache.hudi.common.avro.HoodieBloomFilterWriteSupport.HOODIE_MAX_RECORD_KEY_FOOTER;
+import static
org.apache.hudi.common.avro.HoodieBloomFilterWriteSupport.HOODIE_MIN_RECORD_KEY_FOOTER;
+import static org.apache.hudi.common.util.hash.Hash.MURMUR_HASH;
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNull;
+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.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Tests for {@link HoodieRowDataLanceReader}.
+ */
+class TestHoodieRowDataLanceReader {
+ private static final StoragePath PATH = new StoragePath("/tmp/test.lance");
+
+ @Test
+ void testReadsMetadataAndClosesIdempotently() throws Exception {
+ SimpleBloomFilter bloomFilter = new SimpleBloomFilter(100, 0.01,
MURMUR_HASH);
+ bloomFilter.add("key1");
+ Map<String, String> metadata = new HashMap<>();
+ metadata.put(HOODIE_MIN_RECORD_KEY_FOOTER, "key1");
+ metadata.put(HOODIE_MAX_RECORD_KEY_FOOTER, "key9");
+ metadata.put(HOODIE_AVRO_BLOOM_FILTER_METADATA_KEY,
bloomFilter.serializeToString());
+ LanceFileReader metadataReader = mock(LanceFileReader.class);
+ when(metadataReader.schema()).thenReturn(new
Schema(Collections.emptyList(), metadata));
+ when(metadataReader.numRows()).thenReturn(9L);
+
+ try (MockedStatic<LanceFileReader> mocked = mockLanceOpen(metadataReader))
{
+ HoodieRowDataLanceReader reader = new HoodieRowDataLanceReader(PATH, new
HoodieConfig());
+ assertArrayEquals(new String[] {"key1", "key9"},
reader.readMinMaxRecordKeys());
+ assertInstanceOf(SimpleBloomFilter.class, reader.readBloomFilter());
+ assertTrue(reader.readBloomFilter().mightContain("key1"));
+ assertEquals(9L, reader.getTotalRecords());
+ reader.close();
+ reader.close();
+ }
+ verify(metadataReader, times(1)).close();
+ }
+
+ @Test
+ void testMissingMetadataAndRowCountFailure() throws Exception {
+ LanceFileReader metadataReader = mock(LanceFileReader.class);
+ when(metadataReader.schema()).thenReturn(new
Schema(Collections.emptyList(), null));
+ when(metadataReader.numRows()).thenThrow(new IOException("failed"));
+
+ try (MockedStatic<LanceFileReader> mocked = mockLanceOpen(metadataReader))
{
+ HoodieRowDataLanceReader reader = new HoodieRowDataLanceReader(PATH, new
HoodieConfig());
+ assertNull(reader.readBloomFilter());
+ assertThrows(HoodieException.class, reader::readMinMaxRecordKeys);
+ assertThrows(HoodieException.class, reader::getTotalRecords);
+ reader.close();
+ }
+ }
+
+ @Test
+ void testFilterRowKeysTracksPhysicalPositions() throws Exception {
+ LanceFileReader metadataReader = metadataReader();
+ try (MockedStatic<LanceFileReader> mocked = mockLanceOpen(metadataReader))
{
+ HoodieRowDataLanceReader reader = new HoodieRowDataLanceReader(PATH, new
HoodieConfig()) {
+ @Override
+ public ClosableIterator<String> getRecordKeyIterator() {
+ return ClosableIterator.wrap(List.of("key1", "key2",
"key3").iterator());
+ }
+ };
+
+ assertEquals(
+ Set.of(Pair.of("key2", 1L)),
+ reader.filterRowKeys(Set.of("key2")));
+ assertEquals(3, reader.filterRowKeys(Collections.emptySet()).size());
+ reader.close();
+ }
+ }
+
+ @Test
+ void testRejectsSchemaEvolution() throws Exception {
+ LanceFileReader metadataReader = metadataReader();
+ InternalSchemaManager schemaManager = mock(InternalSchemaManager.class);
+ InternalSchema mergeSchema = mock(InternalSchema.class);
+ when(mergeSchema.isEmptySchema()).thenReturn(false);
+ when(schemaManager.getMergeSchema(PATH.getName())).thenReturn(mergeSchema);
+ HoodieSchema schema = HoodieSchemaUtils.getRecordKeySchema();
+
+ try (MockedStatic<LanceFileReader> mocked = mockLanceOpen(metadataReader))
{
+ HoodieRowDataLanceReader reader = new HoodieRowDataLanceReader(PATH, new
HoodieConfig());
+ assertThrows(HoodieValidationException.class, () ->
reader.getRowDataIterator(
+ schema, schema, schemaManager, Collections.emptyList()));
+ reader.close();
+ }
+ }
+
+ @Test
+ void testRecordKeyIteratorReadsProjectedBatch() throws Exception {
+ HoodieSchema schema = HoodieSchemaUtils.getRecordKeySchema();
+ DataType dataType = RowDataQueryContexts.fromSchema(schema).getRowType();
+ RowType rowType = (RowType) dataType.getLogicalType();
+ String fieldName = rowType.getFieldNames().get(0);
+ LanceFileReader metadataReader = metadataReader();
+ LanceFileReader dataReader = mock(LanceFileReader.class);
+ ArrowReader arrowReader = mock(ArrowReader.class);
+ VectorSchemaRoot batch = mock(VectorSchemaRoot.class);
+
+ try (RootAllocator vectorAllocator = new RootAllocator();
+ VarCharVector vector = new VarCharVector(fieldName, vectorAllocator);
+ MockedStatic<LanceFileReader> mocked = mockLanceOpen(metadataReader,
dataReader)) {
+ vector.allocateNew();
+ vector.setSafe(0, "key1".getBytes(StandardCharsets.UTF_8));
+ vector.setValueCount(1);
+ when(dataReader.readAll(eq(List.of(fieldName)), eq(null),
eq(512))).thenReturn(arrowReader);
+ when(arrowReader.loadNextBatch()).thenReturn(true, false);
+ when(arrowReader.getVectorSchemaRoot()).thenReturn(batch);
+ when(batch.getFieldVectors()).thenReturn(List.of(vector));
+ when(batch.getRowCount()).thenReturn(1);
+
+ HoodieRowDataLanceReader reader = new HoodieRowDataLanceReader(PATH, new
HoodieConfig());
+ try (ClosableIterator<String> iterator = reader.getRecordKeyIterator()) {
+ assertTrue(iterator.hasNext());
+ assertEquals("key1", iterator.next());
+ assertFalse(iterator.hasNext());
+ }
+ verify(arrowReader).close();
+ verify(dataReader).close();
+ verify(metadataReader).close();
+ }
+ }
+
+ @Test
+ void testIteratorCreationFailureClosesDataReader() throws Exception {
+ HoodieSchema schema = HoodieSchemaUtils.getRecordKeySchema();
+ LanceFileReader metadataReader = metadataReader();
+ LanceFileReader dataReader = mock(LanceFileReader.class);
+ when(dataReader.readAll(any(), eq(null), eq(512))).thenThrow(new
IOException("failed"));
+
+ try (MockedStatic<LanceFileReader> mocked = mockLanceOpen(metadataReader,
dataReader)) {
+ HoodieRowDataLanceReader reader = new HoodieRowDataLanceReader(PATH, new
HoodieConfig());
+ assertThrows(HoodieException.class, () -> reader.getRowDataIterator(
+ RowDataQueryContexts.fromSchema(schema).getRowType(), schema));
+ verify(dataReader).close();
+ reader.close();
+ }
+ }
+
+ private static LanceFileReader metadataReader() throws Exception {
+ LanceFileReader reader = mock(LanceFileReader.class);
+ when(reader.schema()).thenReturn(new Schema(Collections.emptyList()));
+ return reader;
+ }
+
+ private static MockedStatic<LanceFileReader>
mockLanceOpen(LanceFileReader... readers) {
+ MockedStatic<LanceFileReader> mocked = mockStatic(LanceFileReader.class);
+ AtomicInteger readerIndex = new AtomicInteger();
+ mocked.when(() -> LanceFileReader.open(eq(PATH.toString()),
any(BufferAllocator.class)))
+ .thenAnswer(invocation -> readers[readerIndex.getAndIncrement()]);
+ return mocked;
+ }
+}
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestInternalSchemaManager.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestInternalSchemaManager.java
new file mode 100644
index 000000000000..0c119a9a4a02
--- /dev/null
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestInternalSchemaManager.java
@@ -0,0 +1,204 @@
+/*
+ * 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.table.format;
+
+import org.apache.hudi.common.config.HoodieCommonConfig;
+import org.apache.hudi.common.schema.internal.InternalSchema;
+import org.apache.hudi.common.schema.internal.Types;
+import org.apache.hudi.common.util.HoodieStorageUtils;
+import org.apache.hudi.common.util.InternalSchemaCache;
+import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.hadoop.HadoopStorageConfiguration;
+
+import org.apache.flink.table.api.DataTypes;
+import org.apache.flink.table.types.DataType;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+
+/**
+ * Tests read-time schema reconciliation in {@link InternalSchemaManager}.
+ */
+class TestInternalSchemaManager {
+
+ @Test
+ void testDisabledManagerUsesEmptySchema() {
+
assertTrue(InternalSchemaManager.DISABLED.getQuerySchema().isEmptySchema());
+ assertTrue(InternalSchemaManager.DISABLED
+ .getMergeSchema("file-id_1-0-1_001.parquet").isEmptySchema());
+
+ org.apache.hadoop.conf.Configuration hadoopConf =
+ new org.apache.hadoop.conf.Configuration();
+ hadoopConf.setBoolean(HoodieCommonConfig.SCHEMA_EVOLUTION_ENABLE.key(),
false);
+ assertSame(
+ InternalSchemaManager.DISABLED,
+ InternalSchemaManager.get(
+ new HadoopStorageConfiguration(hadoopConf),
+ mock(org.apache.hudi.common.table.HoodieTableMetaClient.class)));
+ }
+
+ @Test
+ void testGetCastMapForChangedSelectedField() {
+ InternalSchema querySchema = schema(
+ Types.Field.get(0, false, "id", Types.LongType.get()),
+ Types.Field.get(1, true, "new_name", Types.StringType.get()));
+ InternalSchema fileSchema = schema(
+ Types.Field.get(0, false, "id", Types.IntType.get()),
+ Types.Field.get(1, true, "old_name", Types.StringType.get()));
+ InternalSchemaManager manager = manager(querySchema);
+ DataType[] queryTypes = {DataTypes.BIGINT(), DataTypes.STRING(),
DataTypes.BOOLEAN()};
+
+ CastMap castMap = manager.getCastMap(
+ fileSchema,
+ new String[] {"id", "new_name", "extra"},
+ queryTypes,
+ new int[] {0, 1});
+
+ assertEquals(DataTypes.INT().notNull(), castMap.getFileFieldTypes()[0]);
+ assertEquals(DataTypes.STRING(), castMap.getFileFieldTypes()[1]);
+ assertEquals(DataTypes.BOOLEAN(), castMap.getFileFieldTypes()[2]);
+ assertEquals(7L, castMap.castIfNeeded(0, 7));
+ assertTrue(castMap.toRowDataProjection(new int[] {0, 1}).isPresent());
+ }
+
+ @Test
+ void testGetCastMapWithoutSelectedChangedField() {
+ InternalSchema querySchema = schema(
+ Types.Field.get(0, false, "id", Types.LongType.get()),
+ Types.Field.get(1, true, "name", Types.StringType.get()));
+ InternalSchema fileSchema = schema(
+ Types.Field.get(0, false, "id", Types.IntType.get()),
+ Types.Field.get(1, true, "name", Types.StringType.get()));
+ InternalSchemaManager manager = manager(querySchema);
+
+ CastMap castMap = manager.getCastMap(
+ fileSchema,
+ new String[] {"id", "name"},
+ new DataType[] {DataTypes.BIGINT(), DataTypes.STRING()},
+ new int[] {1});
+
+ assertEquals(DataTypes.INT().notNull(), castMap.getFileFieldTypes()[0]);
+ assertFalse(castMap.toRowDataProjection(new int[] {1}).isPresent());
+ }
+
+ @Test
+ void testGetCastMapWhenTypesAreUnchanged() {
+ InternalSchema schema = schema(
+ Types.Field.get(0, false, "id", Types.IntType.get()),
+ Types.Field.get(1, true, "name", Types.StringType.get()));
+ InternalSchemaManager manager = manager(schema);
+ DataType[] queryTypes = {DataTypes.INT(), DataTypes.STRING()};
+
+ CastMap castMap = manager.getCastMap(
+ schema,
+ new String[] {"id", "name"},
+ queryTypes,
+ new int[] {0, 1});
+
+ assertArrayEquals(queryTypes, castMap.getFileFieldTypes());
+ assertFalse(castMap.toRowDataProjection(new int[] {0, 1}).isPresent());
+ }
+
+ @Test
+ void testGetMergeFieldNamesResolvesRenames() {
+ InternalSchema querySchema = schema(
+ Types.Field.get(0, false, "id", Types.IntType.get()),
+ Types.Field.get(1, true, "new_name", Types.StringType.get()));
+ InternalSchema fileSchema = schema(
+ Types.Field.get(0, false, "id", Types.IntType.get()),
+ Types.Field.get(1, true, "old_name", Types.StringType.get()));
+ InternalSchemaManager manager = manager(querySchema);
+
+ assertArrayEquals(
+ new String[] {"id", "old_name", "extra"},
+ manager.getMergeFieldNames(
+ fileSchema, new String[] {"id", "new_name", "extra"}));
+ assertArrayEquals(
+ new String[] {"id", "new_name"},
+ manager.getMergeFieldNames(
+ querySchema, new String[] {"id", "new_name"}));
+ }
+
+ @Test
+ void testGetMergeSchemaLoadsSchemaForFileVersion() {
+ InternalSchema querySchema =
+ schema(Types.Field.get(0, false, "id", Types.IntType.get()));
+ InternalSchemaManager manager = manager(querySchema);
+ HoodieStorage storage = mock(HoodieStorage.class);
+
+ try (MockedStatic<HoodieStorageUtils> storageUtils =
mockStatic(HoodieStorageUtils.class);
+ MockedStatic<InternalSchemaCache> schemaCache =
mockStatic(InternalSchemaCache.class)) {
+ storageUtils.when(
+ () -> HoodieStorageUtils.getStorage((String) null,
null)).thenReturn(storage);
+ schemaCache.when(
+ () -> InternalSchemaCache.getInternalSchemaByVersionId(
+ 1L, null, storage, null, null, null)).thenReturn(querySchema);
+
+ assertTrue(
+ manager.getMergeSchema("file-id_1-0-1_001.parquet").isEmptySchema());
+ }
+ }
+
+ @Test
+ void testSchemaArgumentsMustBeNonEmpty() {
+ InternalSchemaManager emptyManager =
+ manager(InternalSchema.getEmptyInternalSchema());
+ DataType[] dataTypes = {DataTypes.INT()};
+
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> emptyManager.getCastMap(
+ schema(Types.Field.get(0, false, "id", Types.IntType.get())),
+ new String[] {"id"},
+ dataTypes,
+ new int[] {0}));
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> emptyManager.getMergeFieldNames(
+ schema(Types.Field.get(0, false, "id", Types.IntType.get())),
+ new String[] {"id"}));
+
+ InternalSchemaManager nonEmptyManager =
+ manager(schema(Types.Field.get(0, false, "id", Types.IntType.get())));
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> nonEmptyManager.getCastMap(
+ InternalSchema.getEmptyInternalSchema(),
+ new String[] {"id"},
+ dataTypes,
+ new int[] {0}));
+ }
+
+ private static InternalSchemaManager manager(InternalSchema querySchema) {
+ return new InternalSchemaManager(null, querySchema, null, null, null,
null);
+ }
+
+ private static InternalSchema schema(Types.Field... fields) {
+ return new InternalSchema(Types.RecordType.get(fields));
+ }
+}
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestRecordIterators.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestRecordIterators.java
new file mode 100644
index 000000000000..8086e1516972
--- /dev/null
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/TestRecordIterators.java
@@ -0,0 +1,276 @@
+/*
+ * 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.table.format;
+
+import org.apache.hudi.common.schema.internal.InternalSchema;
+import org.apache.hudi.common.schema.internal.Types;
+import org.apache.hudi.common.util.collection.ClosableIterator;
+import org.apache.hudi.configuration.FlinkOptions;
+import org.apache.hudi.source.ExpressionPredicates.Predicate;
+import org.apache.hudi.storage.StorageConfiguration;
+import org.apache.hudi.storage.StoragePath;
+import org.apache.hudi.storage.hadoop.HadoopStorageConfiguration;
+import org.apache.hudi.storage.inline.InLineFSUtils;
+import org.apache.hudi.table.format.cow.ParquetSplitReaderUtil;
+import
org.apache.hudi.table.format.cow.vector.reader.ParquetColumnarRowSplitReader;
+
+import org.apache.flink.table.api.DataTypes;
+import org.apache.flink.table.data.RowData;
+import org.apache.flink.table.types.DataType;
+import org.apache.hadoop.conf.Configurable;
+import org.apache.hadoop.conf.Configuration;
+import org.apache.parquet.column.ColumnReader;
+import org.apache.parquet.filter.RecordFilter;
+import org.apache.parquet.filter.UnboundRecordFilter;
+import org.apache.parquet.filter2.predicate.FilterPredicate;
+import org.apache.parquet.hadoop.BadConfigurationException;
+import org.apache.parquet.hadoop.util.SerializationUtil;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import java.lang.reflect.InvocationTargetException;
+import java.lang.reflect.Method;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+
+import static org.apache.parquet.filter2.predicate.FilterApi.eq;
+import static org.apache.parquet.filter2.predicate.FilterApi.intColumn;
+import static org.apache.parquet.hadoop.ParquetInputFormat.FILTER_PREDICATE;
+import static
org.apache.parquet.hadoop.ParquetInputFormat.UNBOUND_RECORD_FILTER;
+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.anyInt;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyMap;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Tests configuration and path handling used by {@link RecordIterators}.
+ */
+class TestRecordIterators {
+
+ @Test
+ void testGetFileNameHandlesRegularAndInlinePaths() throws Exception {
+ assertEquals(
+ "file-id_1-0-1_001.parquet",
+ invoke(
+ "getFileName",
+ new Class<?>[] {org.apache.flink.core.fs.Path.class},
+ new org.apache.flink.core.fs.Path(
+ "file:///table/file-id_1-0-1_001.parquet")));
+
+ StoragePath outerPath = new StoragePath(
+ "file:///table/region=us/file-id_1-0-1_001.parquet");
+ StoragePath inlinePath = InLineFSUtils.getInlineFilePath(
+ outerPath, "file", 10L, 20L);
+ assertEquals(
+ outerPath.getName(),
+ invoke(
+ "getFileName",
+ new Class<?>[] {org.apache.flink.core.fs.Path.class},
+ new org.apache.flink.core.fs.Path(inlinePath.toUri())));
+ }
+
+ @Test
+ void testFilterPredicateConfiguration() throws Exception {
+ Configuration conf = new Configuration();
+ assertNull(invoke(
+ "getFilterPredicate",
+ new Class<?>[] {Configuration.class},
+ conf));
+
+ FilterPredicate expected = eq(intColumn("id"), 7);
+ SerializationUtil.writeObjectToConfAsBase64(FILTER_PREDICATE, expected,
conf);
+ FilterPredicate actual = invoke(
+ "getFilterPredicate",
+ new Class<?>[] {Configuration.class},
+ conf);
+ assertEquals(expected.toString(), actual.toString());
+
+ conf.set(FILTER_PREDICATE, "not-base64");
+ InvocationTargetException exception = assertThrows(
+ InvocationTargetException.class,
+ () -> invoke(
+ "getFilterPredicate",
+ new Class<?>[] {Configuration.class},
+ conf));
+ assertTrue(exception.getCause() instanceof RuntimeException);
+ }
+
+ @Test
+ void testUnboundRecordFilterConfiguration() throws Exception {
+ Configuration conf = new Configuration();
+ assertNull(invoke(
+ "getUnboundRecordFilterInstance",
+ new Class<?>[] {Configuration.class},
+ conf));
+
+ conf.setClass(
+ UNBOUND_RECORD_FILTER,
+ ConfigurableRecordFilter.class,
+ UnboundRecordFilter.class);
+ ConfigurableRecordFilter filter = invoke(
+ "getUnboundRecordFilterInstance",
+ new Class<?>[] {Configuration.class},
+ conf);
+ assertSame(conf, filter.getConf());
+
+ conf.setClass(
+ UNBOUND_RECORD_FILTER,
+ InaccessibleRecordFilter.class,
+ UnboundRecordFilter.class);
+ InvocationTargetException exception = assertThrows(
+ InvocationTargetException.class,
+ () -> invoke(
+ "getUnboundRecordFilterInstance",
+ new Class<?>[] {Configuration.class},
+ conf));
+ assertTrue(exception.getCause() instanceof BadConfigurationException);
+ }
+
+ @Test
+ void testGetPartitionSpecUsesOuterPathForInlineFiles() throws Exception {
+ Configuration hadoopConf = new Configuration();
+ hadoopConf.set(
+ FlinkOptions.PARTITION_PATH_FIELD.key(),
+ "region");
+ hadoopConf.setBoolean(
+ FlinkOptions.HIVE_STYLE_PARTITIONING.key(),
+ true);
+ StorageConfiguration<?> storageConf =
+ new HadoopStorageConfiguration(hadoopConf);
+ StoragePath outerPath = new StoragePath(
+ "file:///table/region=us/file-id_1-0-1_001.parquet");
+ StoragePath inlinePath = InLineFSUtils.getInlineFilePath(
+ outerPath, "file", 10L, 20L);
+
+ LinkedHashMap<String, Object> partitionSpec = invoke(
+ "getPartitionSpec",
+ new Class<?>[] {
+ StorageConfiguration.class,
+ StoragePath.class,
+ java.util.List.class,
+ java.util.List.class
+ },
+ storageConf,
+ inlinePath,
+ Collections.singletonList("region"),
+ Collections.singletonList(DataTypes.STRING()));
+ assertEquals(Collections.singletonMap("region", "us"), partitionSpec);
+ }
+
+ @Test
+ void testParquetIteratorBuildsPredicateAndSchemaEvolutionReader() throws
Exception {
+ InternalSchema mergeSchema = new InternalSchema(Types.RecordType.get(
+ Types.Field.get(0, false, "id", Types.IntType.get())));
+ InternalSchemaManager schemaManager = mock(InternalSchemaManager.class);
+ CastMap castMap = new CastMap();
+ DataType[] fieldTypes = {DataTypes.INT()};
+ castMap.setFileFieldTypes(fieldTypes);
+ Predicate predicate = mock(Predicate.class);
+ when(predicate.filter()).thenReturn(eq(intColumn("id"), 7));
+ when(schemaManager.getMergeSchema(anyString())).thenReturn(mergeSchema);
+ when(schemaManager.getCastMap(
+ org.mockito.ArgumentMatchers.eq(mergeSchema),
+ any(String[].class),
+ any(DataType[].class),
+ any(int[].class))).thenReturn(castMap);
+ when(schemaManager.getMergeFieldNames(
+ org.mockito.ArgumentMatchers.eq(mergeSchema),
+ any(String[].class))).thenReturn(new String[] {"id"});
+ ParquetColumnarRowSplitReader reader =
mock(ParquetColumnarRowSplitReader.class);
+ when(reader.reachedEnd()).thenReturn(true);
+
+ try (MockedStatic<ParquetSplitReaderUtil> mocked =
mockStatic(ParquetSplitReaderUtil.class)) {
+ mocked.when(() -> ParquetSplitReaderUtil.genPartColumnarRowReader(
+ anyBoolean(), anyBoolean(), any(Configuration.class),
any(String[].class),
+ any(DataType[].class), anyMap(), any(int[].class), anyInt(), any(),
anyLong(),
+ anyLong(), any(), any())).thenReturn(reader);
+
+ ClosableIterator<RowData> iterator =
RecordIterators.getParquetRecordIterator(
+ schemaManager,
+ true,
+ true,
+ new Configuration(),
+ new String[] {"id"},
+ fieldTypes,
+ Collections.emptyMap(),
+ new int[] {0},
+ 16,
+ new org.apache.flink.core.fs.Path("file:///table.parquet"),
+ 0L,
+ 1L,
+ Collections.singletonList(predicate));
+ assertFalse(iterator.hasNext());
+ iterator.close();
+ verify(reader).close();
+ }
+ }
+
+ @SuppressWarnings("unchecked")
+ private static <T> T invoke(
+ String methodName,
+ Class<?>[] parameterTypes,
+ Object... arguments) throws Exception {
+ Method method = RecordIterators.class.getDeclaredMethod(methodName,
parameterTypes);
+ method.setAccessible(true);
+ return (T) method.invoke(null, arguments);
+ }
+
+ public static class ConfigurableRecordFilter
+ implements UnboundRecordFilter, Configurable {
+ private Configuration conf;
+
+ @Override
+ public RecordFilter bind(Iterable<ColumnReader> readers) {
+ return () -> true;
+ }
+
+ @Override
+ public void setConf(Configuration conf) {
+ this.conf = conf;
+ }
+
+ @Override
+ public Configuration getConf() {
+ return conf;
+ }
+ }
+
+ public static class InaccessibleRecordFilter implements UnboundRecordFilter {
+ private InaccessibleRecordFilter() {
+ }
+
+ @Override
+ public RecordFilter bind(Iterable<ColumnReader> readers) {
+ return () -> true;
+ }
+ }
+}
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/cdc/TestCdcImageManager.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/cdc/TestCdcImageManager.java
new file mode 100644
index 000000000000..e761a2a7945f
--- /dev/null
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/cdc/TestCdcImageManager.java
@@ -0,0 +1,180 @@
+/*
+ * 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.table.format.cdc;
+
+import org.apache.hudi.common.model.FileSlice;
+import org.apache.hudi.common.model.HoodieFileGroupId;
+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.table.format.FormatUtils;
+
+import org.apache.flink.table.data.GenericRowData;
+import org.apache.flink.table.data.RowData;
+import org.apache.flink.table.data.StringData;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.RowType;
+import org.apache.flink.table.types.logical.VarCharType;
+import org.apache.flink.types.RowKind;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import java.io.ByteArrayOutputStream;
+import java.io.EOFException;
+import java.io.IOException;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+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.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Tests {@link CdcImageManager}.
+ */
+class TestCdcImageManager {
+
+ @Test
+ void testImageRecordLifecycle() {
+ HoodieWriteConfig writeConfig = mock(HoodieWriteConfig.class);
+ CdcImageManager imageManager = new CdcImageManager(
+ rowType("value"),
+ writeConfig,
+ split -> {
+ throw new AssertionError("No split should be loaded by this test");
+ });
+ ExternalSpillableMap<String, byte[]> imageCache = mockImageCache();
+
+ GenericRowData original =
GenericRowData.of(StringData.fromString("before"));
+ imageManager.updateImageRecord("key-1", imageCache, original);
+
+ RowData image = imageManager.getImageRecord("key-1", imageCache,
RowKind.UPDATE_BEFORE);
+ assertEquals(RowKind.UPDATE_BEFORE, image.getRowKind());
+ assertEquals("before", image.getString(0).toString());
+
+ RowData removed = imageManager.removeImageRecord("key-1", imageCache);
+ assertEquals("before", removed.getString(0).toString());
+ assertNull(imageManager.removeImageRecord("key-1", imageCache));
+ assertThrows(
+ IllegalStateException.class,
+ () -> imageManager.getImageRecord("missing", imageCache,
RowKind.DELETE));
+ assertSame(writeConfig, imageManager.getWriteConfig());
+
+ imageManager.close();
+ }
+
+ @Test
+ void testDataViewAdapters() throws IOException {
+ ByteArrayOutputStream outputStream = new ByteArrayOutputStream();
+ CdcImageManager.BytesArrayOutputView outputView =
+ new CdcImageManager.BytesArrayOutputView(outputStream);
+ outputView.writeByte(7);
+ outputView.skipBytesToWrite(2);
+ outputView.write(new CdcImageManager.BytesArrayInputView(new byte[] {8,
9}), 2);
+ outputView.flush();
+
+ assertArrayEquals(new byte[] {7, 0, 0, 8, 9}, outputStream.toByteArray());
+
+ CdcImageManager.BytesArrayInputView inputView =
+ new CdcImageManager.BytesArrayInputView(outputStream.toByteArray());
+ inputView.skipBytesToRead(3);
+ assertEquals(8, inputView.readByte());
+ assertEquals(9, inputView.readUnsignedByte());
+
+ CdcImageManager.BytesArrayInputView truncatedInputView =
+ new CdcImageManager.BytesArrayInputView(new byte[] {1, 2});
+ assertThrows(EOFException.class, () ->
truncatedInputView.skipBytesToRead(3));
+ }
+
+ @Test
+ void testImageCacheReuseEvictionAndClose() throws IOException {
+ HoodieWriteConfig writeConfig = mock(HoodieWriteConfig.class);
+ when(writeConfig.getBasePath()).thenReturn("/table");
+ ExternalSpillableMap<String, byte[]> first = mockImageCache();
+ ExternalSpillableMap<String, byte[]> second = mockImageCache();
+ ExternalSpillableMap<String, byte[]> third = mockImageCache();
+ GenericRowData row = GenericRowData.of(
+ StringData.fromString("commit"),
+ StringData.fromString("seq"),
+ StringData.fromString("key-1"));
+ CdcImageManager imageManager = new CdcImageManager(
+ RowType.of(
+ new LogicalType[] {new VarCharType(), new VarCharType(), new
VarCharType()},
+ new String[] {"commit", "seq", "record_key"}),
+ writeConfig,
+ split -> ClosableIterator.wrap(List.<RowData>of(row).iterator()));
+
+ try (MockedStatic<FormatUtils> mockedFormatUtils =
mockStatic(FormatUtils.class)) {
+ mockedFormatUtils.when(() -> FormatUtils.spillableMap(
+ writeConfig, 1024L, CdcImageManager.class.getSimpleName()))
+ .thenReturn(first, second, third);
+
+ FileSlice slice1 = fileSlice("001");
+ FileSlice slice2 = fileSlice("002");
+ FileSlice slice3 = fileSlice("003");
+ assertSame(first, imageManager.getOrLoadImages(1024L, slice1));
+ assertSame(first, imageManager.getOrLoadImages(1024L, slice1));
+ assertSame(second, imageManager.getOrLoadImages(1024L, slice2));
+ assertSame(third, imageManager.getOrLoadImages(1024L, slice3));
+ verify(first).close();
+ verify(first).put(anyString(), any(byte[].class));
+
+ imageManager.close();
+ verify(second).close();
+ verify(third).close();
+ imageManager.close();
+ verify(second, times(1)).close();
+ }
+ }
+
+ @SuppressWarnings("unchecked")
+ private static ExternalSpillableMap<String, byte[]> mockImageCache() {
+ ExternalSpillableMap<String, byte[]> imageCache =
mock(ExternalSpillableMap.class);
+ Map<String, byte[]> records = new HashMap<>();
+ when(imageCache.get(anyString())).thenAnswer(
+ invocation -> records.get(invocation.getArgument(0)));
+ when(imageCache.put(anyString(), any(byte[].class))).thenAnswer(
+ invocation -> records.put(invocation.getArgument(0),
invocation.getArgument(1)));
+ when(imageCache.remove(anyString())).thenAnswer(
+ invocation -> records.remove(invocation.getArgument(0)));
+ return imageCache;
+ }
+
+ private static RowType rowType(String fieldName) {
+ return RowType.of(
+ new LogicalType[] {new VarCharType()},
+ new String[] {fieldName});
+ }
+
+ private static FileSlice fileSlice(String instant) {
+ return new FileSlice(
+ new HoodieFileGroupId("partition", "file"), instant, null,
java.util.Collections.emptyList());
+ }
+}
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/cdc/TestCdcIterators.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/cdc/TestCdcIterators.java
new file mode 100644
index 000000000000..2144f05c7c6e
--- /dev/null
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/cdc/TestCdcIterators.java
@@ -0,0 +1,229 @@
+/*
+ * 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.table.format.cdc;
+
+import org.apache.hudi.common.model.FileSlice;
+import org.apache.hudi.common.model.HoodieBaseFile;
+import org.apache.hudi.common.model.HoodieLogFile;
+import org.apache.hudi.common.table.cdc.HoodieCDCFileSplit;
+import org.apache.hudi.common.table.cdc.HoodieCDCInferenceCase;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.common.util.collection.ClosableIterator;
+import org.apache.hudi.storage.StoragePath;
+import org.apache.hudi.table.format.mor.MergeOnReadInputSplit;
+
+import org.apache.flink.table.data.GenericRowData;
+import org.apache.flink.table.data.RowData;
+import org.apache.flink.table.data.StringData;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.RowType;
+import org.apache.flink.table.types.logical.VarCharType;
+import org.apache.flink.types.RowKind;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Tests the lightweight iterator composition in {@link CdcIterators}.
+ */
+class TestCdcIterators {
+
+ @Test
+ void testAddBaseFileIteratorSetsInsertKindAndClosesNestedIterator() {
+ GenericRowData row = GenericRowData.of(StringData.fromString("value"));
+ ClosableIterator<RowData> nested = mockIterator();
+ when(nested.hasNext()).thenReturn(true, false);
+ when(nested.next()).thenReturn(row);
+
+ CdcIterators.AddBaseFileIterator iterator =
+ new CdcIterators.AddBaseFileIterator(nested);
+ assertTrue(iterator.hasNext());
+ assertSame(row, iterator.next());
+ assertEquals(RowKind.INSERT, row.getRowKind());
+ assertFalse(iterator.hasNext());
+
+ iterator.close();
+ verify(nested).close();
+ }
+
+ @Test
+ void testRemoveBaseFileIteratorProjectsAndSetsDeleteKind() {
+ GenericRowData row = GenericRowData.of(
+ StringData.fromString("ignored"), StringData.fromString("selected"));
+ ClosableIterator<RowData> nested = mockIterator();
+ when(nested.hasNext()).thenReturn(true);
+ when(nested.next()).thenReturn(row);
+
+ CdcIterators.RemoveBaseFileIterator iterator =
+ new CdcIterators.RemoveBaseFileIterator(
+ rowType("selected"),
+ new int[] {1},
+ nested);
+ assertTrue(iterator.hasNext());
+ RowData projected = iterator.next();
+ assertEquals(RowKind.DELETE, projected.getRowKind());
+ assertEquals("selected", projected.getString(0).toString());
+
+ iterator.close();
+ verify(nested).close();
+ }
+
+ @Test
+ void testCdcFileSplitsIteratorMovesAcrossEmptySplitsAndClosesResources() {
+ HoodieCDCFileSplit first = new HoodieCDCFileSplit(
+ "001", HoodieCDCInferenceCase.BASE_FILE_INSERT, "first.parquet");
+ HoodieCDCFileSplit second = new HoodieCDCFileSplit(
+ "002", HoodieCDCInferenceCase.BASE_FILE_INSERT, "second.parquet");
+ ClosableIterator<RowData> firstIterator = mockIterator();
+ ClosableIterator<RowData> secondIterator = mockIterator();
+ GenericRowData row = GenericRowData.of(StringData.fromString("row"));
+ when(firstIterator.hasNext()).thenReturn(false);
+ when(secondIterator.hasNext()).thenReturn(true, false);
+ when(secondIterator.next()).thenReturn(row);
+ CdcImageManager imageManager = mock(CdcImageManager.class);
+
+ CdcIterators.CdcFileSplitsIterator iterator =
+ new CdcIterators.CdcFileSplitsIterator(
+ new HoodieCDCFileSplit[] {first, second},
+ imageManager,
+ split -> split == first ? firstIterator : secondIterator);
+
+ assertTrue(iterator.hasNext());
+ assertSame(row, iterator.next());
+ assertFalse(iterator.hasNext());
+ iterator.close();
+
+ verify(firstIterator).close();
+ verify(secondIterator).close();
+ verify(imageManager).close();
+ }
+
+ @Test
+ void testReplaceCommitIteratorReadsBeforeSlice() {
+ FileSlice beforeSlice = fileSlice();
+ HoodieCDCFileSplit fileSplit = new HoodieCDCFileSplit(
+ "002",
+ HoodieCDCInferenceCase.REPLACE_COMMIT,
+ Collections.emptyList(),
+ Option.of(beforeSlice),
+ Option.empty());
+ GenericRowData row = GenericRowData.of(
+ StringData.fromString("ignored"), StringData.fromString("selected"));
+ ClosableIterator<RowData> nested = mockIterator();
+ when(nested.hasNext()).thenReturn(true);
+ when(nested.next()).thenReturn(row);
+ AtomicReference<MergeOnReadInputSplit> capturedSplit = new
AtomicReference<>();
+
+ CdcIterators.ReplaceCommitIterator iterator =
+ new CdcIterators.ReplaceCommitIterator(
+ "/table",
+ rowType("selected"),
+ new int[] {1},
+ 1024L,
+ fileSplit,
+ split -> {
+ capturedSplit.set(split);
+ return nested;
+ });
+
+ assertEquals("/table", capturedSplit.get().getTablePath());
+ assertTrue(iterator.hasNext());
+ RowData projected = iterator.next();
+ assertEquals(RowKind.DELETE, projected.getRowKind());
+ assertEquals("selected", projected.getString(0).toString());
+ iterator.close();
+ verify(nested).close();
+
+ HoodieCDCFileSplit missingBeforeSlice = new HoodieCDCFileSplit(
+ "003",
+ HoodieCDCInferenceCase.REPLACE_COMMIT,
+ Collections.emptyList());
+ assertThrows(
+ IllegalStateException.class,
+ () -> new CdcIterators.ReplaceCommitIterator(
+ "/table",
+ rowType("selected"),
+ new int[] {0},
+ 1024L,
+ missingBeforeSlice,
+ split -> nested));
+ }
+
+ @Test
+ void testFileSliceAndSingleLogFileSplitConversion() {
+ FileSlice fileSlice = fileSlice();
+ fileSlice.addLogFile(new HoodieLogFile(
+ new StoragePath("/table/region=us/.file-id_002.log.1_1-0-1")));
+ fileSlice.addLogFile(new HoodieLogFile(
+ new StoragePath("/table/region=us/.file-id_002.log.2_1-0-1.cdc")));
+
+ MergeOnReadInputSplit split =
+ CdcIterators.fileSlice2Split("/table", fileSlice, 4096L);
+ assertEquals(
+ "/table/region=us/file-id_1-0-1_001.parquet",
+ split.getBasePath().get());
+ assertEquals(
+ Collections.singletonList("/table/region=us/.file-id_002.log.1_1-0-1"),
+ split.getLogPaths().get());
+ assertEquals("file-id", split.getFileId());
+ assertEquals("region=us", split.getPartitionPath());
+ assertEquals(4096L, split.getMaxCompactionMemoryInBytes());
+
+ MergeOnReadInputSplit logSplit = CdcIterators.singleLogFile2Split(
+ "/table",
+ "/table/region=us/.file-id_003.log.1_1-0-1",
+ 8192L);
+ assertFalse(logSplit.getBasePath().isPresent());
+ assertEquals(
+ Collections.singletonList("/table/region=us/.file-id_003.log.1_1-0-1"),
+ logSplit.getLogPaths().get());
+ assertEquals("003", logSplit.getLatestCommit());
+ assertEquals("file-id", logSplit.getFileId());
+ assertEquals("region=us", logSplit.getPartitionPath());
+ }
+
+ private static FileSlice fileSlice() {
+ FileSlice fileSlice = new FileSlice("region=us", "001", "file-id");
+ fileSlice.setBaseFile(new HoodieBaseFile(
+ "/table/region=us/file-id_1-0-1_001.parquet"));
+ return fileSlice;
+ }
+
+ @SuppressWarnings("unchecked")
+ private static ClosableIterator<RowData> mockIterator() {
+ return mock(ClosableIterator.class);
+ }
+
+ private static RowType rowType(String fieldName) {
+ return RowType.of(
+ new LogicalType[] {new VarCharType()},
+ new String[] {fieldName});
+ }
+}
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/cow/TestCopyOnWriteInputFormat.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/cow/TestCopyOnWriteInputFormat.java
new file mode 100644
index 000000000000..c8210d170652
--- /dev/null
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/format/cow/TestCopyOnWriteInputFormat.java
@@ -0,0 +1,199 @@
+/*
+ * 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.table.format.cow;
+
+import org.apache.hudi.common.util.collection.ClosableIterator;
+import org.apache.hudi.configuration.FlinkOptions;
+import org.apache.hudi.table.format.InternalSchemaManager;
+import org.apache.hudi.util.HoodieSchemaConverter;
+import org.apache.hudi.utils.TestConfigurations;
+
+import org.apache.flink.api.common.io.FilePathFilter;
+import org.apache.flink.core.fs.FileInputSplit;
+import org.apache.flink.core.fs.Path;
+import org.apache.flink.table.data.GenericRowData;
+import org.apache.flink.table.data.RowData;
+import org.apache.flink.table.data.StringData;
+import org.apache.flink.table.types.DataType;
+import org.apache.hadoop.fs.FileStatus;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import java.io.IOException;
+import java.lang.reflect.Field;
+import java.nio.file.Files;
+import java.util.Collections;
+import java.util.List;
+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.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Tests split planning and iterator lifecycle in {@link
CopyOnWriteInputFormat}.
+ */
+class TestCopyOnWriteInputFormat {
+
+ @TempDir
+ java.nio.file.Path tempDir;
+
+ @Test
+ void testCreateInputSplitsFiltersHiddenFilesAndReadsNestedDirectories()
throws IOException {
+ java.nio.file.Path visible = tempDir.resolve("visible.parquet");
+ java.nio.file.Path empty = tempDir.resolve("empty.parquet");
+ java.nio.file.Path hidden = tempDir.resolve("_hidden.parquet");
+ java.nio.file.Path nested = tempDir.resolve("nested/nested.parquet");
+ Files.write(visible, new byte[32]);
+ Files.createFile(empty);
+ Files.write(hidden, new byte[8]);
+ Files.createDirectories(nested.getParent());
+ Files.write(nested, new byte[16]);
+
+ CopyOnWriteInputFormat inputFormat = inputFormat(
+ new Path[] {new Path(tempDir.toUri())}, Long.MAX_VALUE);
+ inputFormat.setNestedFileEnumeration(true);
+
+ assertThrows(IllegalArgumentException.class, () ->
inputFormat.createInputSplits(0));
+ FileInputSplit[] splits = inputFormat.createInputSplits(2);
+ assertTrue(splits.length >= 3);
+ assertTrue(inputFormat.supportsMultiPaths());
+ assertTrue(java.util.Arrays.stream(splits)
+ .anyMatch(split ->
split.getPath().getName().equals("visible.parquet")));
+ assertTrue(java.util.Arrays.stream(splits)
+ .anyMatch(split -> split.getPath().getName().equals("empty.parquet")));
+ assertTrue(java.util.Arrays.stream(splits)
+ .anyMatch(split ->
split.getPath().getName().equals("nested.parquet")));
+ assertFalse(java.util.Arrays.stream(splits)
+ .anyMatch(split ->
split.getPath().getName().equals("_hidden.parquet")));
+ }
+
+ @Test
+ void testUnsplittableLanceFileUsesWholeFileSplit() throws IOException {
+ java.nio.file.Path lance = tempDir.resolve("data.lance");
+ Files.write(lance, new byte[32]);
+ CopyOnWriteInputFormat inputFormat = inputFormat(
+ new Path[] {new Path(lance.toUri())}, Long.MAX_VALUE);
+
+ FileInputSplit[] splits = inputFormat.createInputSplits(8);
+ assertEquals(1, splits.length);
+ assertEquals(-1L, splits[0].getLength());
+ }
+
+ @Test
+ void testAcceptFileUsesBuiltInAndCustomFilters() {
+ CopyOnWriteInputFormat inputFormat = inputFormat(
+ new Path[] {new Path(tempDir.toUri())}, Long.MAX_VALUE);
+ assertFalse(inputFormat.acceptFile(fileStatus("_metadata")));
+ assertFalse(inputFormat.acceptFile(fileStatus(".hidden")));
+ assertTrue(inputFormat.acceptFile(fileStatus("data.parquet")));
+
+ inputFormat.setFilesFilter(new FilePathFilter() {
+ @Override
+ public boolean filterPath(Path filePath) {
+ return filePath.getName().endsWith(".skip");
+ }
+ });
+ assertFalse(inputFormat.acceptFile(fileStatus("data.skip")));
+ assertTrue(inputFormat.acceptFile(fileStatus("data.parquet")));
+ }
+
+ @Test
+ void testLimitAndIteratorLifecycle() throws Exception {
+ CopyOnWriteInputFormat inputFormat = inputFormat(
+ new Path[] {new Path(tempDir.toUri())}, 1L);
+ ClosableIterator<RowData> iterator = mockIterator();
+ GenericRowData row = GenericRowData.of(StringData.fromString("value"));
+ when(iterator.hasNext()).thenReturn(true);
+ when(iterator.next()).thenReturn(row);
+ setIterator(inputFormat, iterator);
+
+ assertFalse(inputFormat.reachedEnd());
+ assertSame(row, inputFormat.nextRecord(null));
+ assertTrue(inputFormat.reachedEnd());
+ verify(iterator).hasNext();
+
+ inputFormat.close();
+ verify(iterator).close();
+ inputFormat.close();
+ }
+
+ @Test
+ void testReachedEndDelegatesWhenLimitIsNotReached() throws Exception {
+ CopyOnWriteInputFormat inputFormat = inputFormat(
+ new Path[] {new Path(tempDir.toUri())}, Long.MAX_VALUE);
+ ClosableIterator<RowData> iterator = mockIterator();
+ when(iterator.hasNext()).thenReturn(false);
+ setIterator(inputFormat, iterator);
+
+ assertTrue(inputFormat.reachedEnd());
+ verify(iterator).hasNext();
+ verify(iterator, never()).next();
+ }
+
+ private static CopyOnWriteInputFormat inputFormat(Path[] paths, long limit) {
+ List<String> fieldNames = TestConfigurations.ROW_TYPE.getFieldNames();
+ List<DataType> fieldTypes = TestConfigurations.ROW_DATA_TYPE.getChildren();
+ return new CopyOnWriteInputFormat(
+ paths,
+ fieldNames.toArray(new String[0]),
+ fieldTypes.toArray(new DataType[0]),
+ IntStream.range(0, fieldNames.size()).toArray(),
+ FlinkOptions.PARTITION_DEFAULT_NAME.defaultValue(),
+ FlinkOptions.PARTITION_PATH_FIELD.defaultValue(),
+ false,
+ Collections.emptyList(),
+ limit,
+ new org.apache.hadoop.conf.Configuration(),
+ true,
+ InternalSchemaManager.DISABLED,
+ HoodieSchemaConverter.convertToSchema(
+ TestConfigurations.ROW_TYPE.copy()));
+ }
+
+ private static FileStatus fileStatus(String name) {
+ return new FileStatus(
+ 1L,
+ false,
+ 1,
+ 1L,
+ 0L,
+ new org.apache.hadoop.fs.Path("/table/" + name));
+ }
+
+ private static void setIterator(
+ CopyOnWriteInputFormat inputFormat,
+ ClosableIterator<RowData> iterator) throws ReflectiveOperationException {
+ Field field = CopyOnWriteInputFormat.class.getDeclaredField("itr");
+ field.setAccessible(true);
+ field.set(inputFormat, iterator);
+ }
+
+ @SuppressWarnings("unchecked")
+ private static ClosableIterator<RowData> mockIterator() {
+ return mock(ClosableIterator.class);
+ }
+}