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 6543aec2553e feat(lance): enable lance as native log format (#19283)
6543aec2553e is described below
commit 6543aec2553ecbc0dec84d5a88aa7b09292ffc9b
Author: Shuo Cheng <[email protected]>
AuthorDate: Wed Jul 15 18:42:43 2026 +0800
feat(lance): enable lance as native log format (#19283)
* feat(lance): enable lance as native log format
---
.../org/apache/hudi/io/AppendHandleFactory.java | 4 +--
.../FileGroupReaderBasedInlineLogAppendHandle.java | 4 +++
.../hudi/io/FileGroupReaderBasedMergeHandle.java | 2 +-
.../FileGroupReaderBasedNativeLogAppendHandle.java | 4 +++
.../hudi/io/HoodieNativeLogAppendHandle.java | 2 +-
.../hudi/io/cdc/HoodieCDCLogWriterFactory.java | 2 +-
.../java/org/apache/hudi/table/HoodieTable.java | 8 +++++
.../hudi/table/action/compact/HoodieCompactor.java | 28 +++++++++--------
.../HoodieLogCompactionPlanGenerator.java | 3 +-
.../org/apache/hudi/util/CommonClientUtils.java | 16 +++-------
.../TestScheduleCompactionActionExecutor.java | 2 +-
.../apache/hudi/utils/TestCommonClientUtils.java | 35 +++++-----------------
.../apache/hudi/io/FlinkWriteHandleFactory.java | 4 +--
.../hudi/table/HoodieFlinkMergeOnReadTable.java | 2 +-
.../hudi/table/HoodieSparkMergeOnReadTable.java | 2 +-
.../org/apache/hudi/table/HoodieSparkTable.java | 12 ++++++++
.../SparkFileFormatInternalRowReaderContext.scala | 15 +++++-----
.../apache/hudi/table/TestHoodieSparkTable.java | 32 ++++++++++++++++++++
.../hudi/common/avro/HoodieAvroReaderContext.java | 11 ++++++-
.../common/engine/AvroReaderContextFactory.java | 14 +++++++--
.../java/org/apache/hudi/common/fs/FSUtils.java | 7 +++--
.../org/apache/hudi/common/fs/FileNameParser.java | 24 +++++++++++++--
.../hudi/common/model/TestHoodieLogFile.java | 11 +++++++
.../handler/MetadataTableCompactHandler.java | 7 +++--
.../table/format/FlinkRowDataReaderContext.java | 2 +-
.../hudi/functional/TestLanceDataSource.scala | 28 ++++++++++++++++-
26 files changed, 201 insertions(+), 80 deletions(-)
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/AppendHandleFactory.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/AppendHandleFactory.java
index 25d49890366d..d8f3f37ac4a9 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/AppendHandleFactory.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/AppendHandleFactory.java
@@ -34,7 +34,7 @@ public class AppendHandleFactory<T, I, K, O> extends
WriteHandleFactory<T, I, K,
final String fileIdPrefix,
final TaskContextSupplier sparkTaskContextSupplier) {
String fileId = getNextFileId(fileIdPrefix);
- if (CommonClientUtils.shouldWriteNativeLogs(hoodieConfig,
hoodieTable.getMetaClient().getTableConfig())) {
+ if (CommonClientUtils.shouldWriteNativeLogs(hoodieConfig)) {
return new HoodieNativeLogAppendHandle<>(hoodieConfig, commitTime,
hoodieTable, partitionPath,
fileId, sparkTaskContextSupplier);
}
@@ -46,7 +46,7 @@ public class AppendHandleFactory<T, I, K, O> extends
WriteHandleFactory<T, I, K,
final HoodieTable<T, I, K, O>
hoodieTable, final String partitionPath,
final String fileId, final
Iterator<HoodieRecord<T>> recordItr,
final TaskContextSupplier
sparkTaskContextSupplier) {
- if (CommonClientUtils.shouldWriteNativeLogs(hoodieConfig,
hoodieTable.getMetaClient().getTableConfig())) {
+ if (CommonClientUtils.shouldWriteNativeLogs(hoodieConfig)) {
return new HoodieNativeLogAppendHandle<>(hoodieConfig, commitTime,
hoodieTable, partitionPath,
fileId, recordItr, sparkTaskContextSupplier);
}
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedInlineLogAppendHandle.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedInlineLogAppendHandle.java
index d2c0d1e0a002..3d3fdb840c74 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedInlineLogAppendHandle.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedInlineLogAppendHandle.java
@@ -69,6 +69,10 @@ public class FileGroupReaderBasedInlineLogAppendHandle<T, I,
K, O> extends Hoodi
super(config, instantTime, hoodieTable, operation.getPartitionPath(),
operation.getFileId(), taskContextSupplier);
this.operation = operation;
this.readerContext = readerContext;
+ // File-group reader output already conforms to the writer schema.
Preserve its metadata fields
+ // instead of prepending another metadata overlay and shifting the
data-field ordinals.
+ this.isLogCompaction = true;
+ this.useWriterSchema = true;
}
@Override
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
index 42cf14cc9c9d..beba9067c7e6 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedMergeHandle.java
@@ -182,7 +182,7 @@ public class FileGroupReaderBasedMergeHandle<T, I, K, O>
extends HoodieWriteMerg
private HoodieCDCLogWriter<?> createCDCLogWriter() {
HoodieTableConfig tableConfig =
hoodieTable.getMetaClient().getTableConfig();
- if (CommonClientUtils.shouldWriteNativeLogs(config, tableConfig)) {
+ if (CommonClientUtils.shouldWriteNativeLogs(config)) {
return new HoodieNativeCDCLogger(
instantTime,
config,
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedNativeLogAppendHandle.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedNativeLogAppendHandle.java
index 98cedd037a83..e58c12dfdbc5 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedNativeLogAppendHandle.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/FileGroupReaderBasedNativeLogAppendHandle.java
@@ -65,6 +65,10 @@ public class FileGroupReaderBasedNativeLogAppendHandle<T, I,
K, O> extends Hoodi
super(config, instantTime, hoodieTable, operation.getPartitionPath(),
operation.getFileId(), taskContextSupplier);
this.operation = operation;
this.readerContext = readerContext;
+ // File-group reader output is materialized with the writer schema,
including Hudi meta fields.
+ // Treating it as a regular append would prepend another meta-field
overlay and shift all data ordinals.
+ this.isLogCompaction = true;
+ this.useWriterSchema = true;
}
@Override
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogAppendHandle.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogAppendHandle.java
index 365a9487529d..860f6055d5e2 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogAppendHandle.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/HoodieNativeLogAppendHandle.java
@@ -107,7 +107,7 @@ public class HoodieNativeLogAppendHandle<T, I, K, O>
extends HoodieAppendHandle<
hoodieTable.getBaseFileFormat(),
writeSchemaWithMetaFields,
taskContextSupplier,
-
hoodieTable.getReaderContextFactoryForWrite().getContext().getRecordContext(),
+ hoodieTable.getRecordContextForWrite(),
orderingFields,
baseFileInstantTimeOfPositions);
} catch (IOException e) {
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/cdc/HoodieCDCLogWriterFactory.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/cdc/HoodieCDCLogWriterFactory.java
index 73c4341f5ea3..8df1de8759bc 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/cdc/HoodieCDCLogWriterFactory.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/io/cdc/HoodieCDCLogWriterFactory.java
@@ -55,7 +55,7 @@ public final class HoodieCDCLogWriterFactory {
TaskContextSupplier taskContextSupplier,
Supplier<HoodieLogFormat.Writer> logWriterSupplier) {
HoodieTableConfig tableConfig =
hoodieTable.getMetaClient().getTableConfig();
- if (CommonClientUtils.shouldWriteNativeLogs(config, tableConfig)) {
+ if (CommonClientUtils.shouldWriteNativeLogs(config)) {
return new HoodieAvroNativeCDCLogger(
instantTime,
config,
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
index dd2b9eadd17b..76cb3b291e2f 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/HoodieTable.java
@@ -35,6 +35,7 @@ import org.apache.hudi.common.config.HoodieMetadataConfig;
import org.apache.hudi.common.engine.HoodieEngineContext;
import org.apache.hudi.common.engine.HoodieLocalEngineContext;
import org.apache.hudi.common.engine.ReaderContextFactory;
+import org.apache.hudi.common.engine.RecordContext;
import org.apache.hudi.common.engine.TaskContextSupplier;
import org.apache.hudi.common.fs.ConsistencyGuard;
import org.apache.hudi.common.fs.ConsistencyGuard.FileVisibility;
@@ -1281,4 +1282,11 @@ public abstract class HoodieTable<T, I, K, O> implements
Serializable {
return (ReaderContextFactory<T>)
getContext().getReaderContextFactoryForWrite(metaClient,
config.getRecordMerger().getRecordType(),
config.getProps());
}
+
+ /**
+ * Returns the record context used by the write path.
+ */
+ public RecordContext<?> getRecordContextForWrite() {
+ return getReaderContextFactoryForWrite().getContext().getRecordContext();
+ }
}
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/HoodieCompactor.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/HoodieCompactor.java
index 208f64af3cd0..88cdf689288e 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/HoodieCompactor.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/HoodieCompactor.java
@@ -20,8 +20,8 @@ package org.apache.hudi.table.action.compact;
import org.apache.hudi.avro.model.HoodieCompactionPlan;
import org.apache.hudi.client.WriteStatus;
-import org.apache.hudi.common.avro.HoodieAvroReaderContext;
import org.apache.hudi.common.data.HoodieData;
+import org.apache.hudi.common.engine.AvroReaderContextFactory;
import org.apache.hudi.common.engine.HoodieEngineContext;
import org.apache.hudi.common.engine.HoodieReaderContext;
import org.apache.hudi.common.engine.ReaderContextFactory;
@@ -50,7 +50,6 @@ import org.apache.hudi.table.HoodieTable;
import org.apache.hudi.util.CommonClientUtils;
import lombok.extern.slf4j.Slf4j;
-import org.apache.avro.generic.IndexedRecord;
import java.io.IOException;
import java.io.Serializable;
@@ -133,8 +132,16 @@ public abstract class HoodieCompactor<T, I, K, O>
implements Serializable {
Option<InstantRange> instantRange =
CompactHelpers.getInstance().getInstantRange(metaClient);
if (operationType == WriteOperationType.LOG_COMPACT) {
+ ReaderContextFactory<?> readerContextFactory;
+ if (metaClient.isMetadataTable()) {
+ readerContextFactory = new AvroReaderContextFactory(metaClient,
metaClient.getTableConfig().getPayloadClass(), instantRange, config.getProps());
+ } else {
+ readerContextFactory = context.getReaderContextFactory(metaClient);
+ }
+
return context.parallelize(operations).map(
- operation -> logCompact(config, operation,
compactionInstantTime, instantRange, table, taskContextSupplier))
+ operation -> logCompact(config, operation,
compactionInstantTime, table, taskContextSupplier,
+ readerContextFactory.getContext()))
.flatMap(List::iterator);
} else {
ReaderContextFactory<T> readerContextFactory;
@@ -167,15 +174,12 @@ public abstract class HoodieCompactor<T, I, K, O>
implements Serializable {
}
public List<WriteStatus> logCompact(HoodieWriteConfig writeConfig,
- CompactionOperation operation,
- String instantTime,
- Option<InstantRange> instantRange,
- HoodieTable table,
- TaskContextSupplier taskContextSupplier)
throws IOException {
- HoodieReaderContext<IndexedRecord> readerContext = new
HoodieAvroReaderContext(
- table.getStorageConf(), table.getMetaClient().getTableConfig(),
instantRange, Option.empty(), writeConfig.getProps());
- HoodieAppendHandle<IndexedRecord, ?, ?, ?> appendHandle =
CommonClientUtils.shouldWriteNativeLogs(
- writeConfig, table.getMetaClient().getTableConfig())
+ CompactionOperation operation,
+ String instantTime,
+ HoodieTable table,
+ TaskContextSupplier
taskContextSupplier,
+ HoodieReaderContext readerContext)
throws IOException {
+ HoodieAppendHandle<T, ?, ?, ?> appendHandle =
CommonClientUtils.shouldWriteNativeLogs(writeConfig)
? new FileGroupReaderBasedNativeLogAppendHandle<>(writeConfig,
instantTime, table, operation, taskContextSupplier, readerContext)
: new FileGroupReaderBasedInlineLogAppendHandle<>(writeConfig,
instantTime, table, operation, taskContextSupplier, readerContext);
appendHandle.doAppend();
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/plan/generators/HoodieLogCompactionPlanGenerator.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/plan/generators/HoodieLogCompactionPlanGenerator.java
index e0a696b24565..3b9dc0997e60 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/plan/generators/HoodieLogCompactionPlanGenerator.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/table/action/compact/plan/generators/HoodieLogCompactionPlanGenerator.java
@@ -75,7 +75,8 @@ public class HoodieLogCompactionPlanGenerator<T extends
HoodieRecordPayload, I,
@Override
protected boolean filterFileSlice(FileSlice fileSlice, String
lastCompletedInstantTime,
Set<HoodieFileGroupId>
pendingFileGroupIds, Option<InstantRange> instantRange) {
- return super.filterFileSlice(fileSlice, lastCompletedInstantTime,
pendingFileGroupIds, instantRange) &&
isFileSliceEligibleForLogCompaction(fileSlice, lastCompletedInstantTime,
instantRange);
+ return super.filterFileSlice(fileSlice, lastCompletedInstantTime,
pendingFileGroupIds, instantRange)
+ && isFileSliceEligibleForLogCompaction(fileSlice,
lastCompletedInstantTime, instantRange);
}
@Override
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/util/CommonClientUtils.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/util/CommonClientUtils.java
index 9c7d3463453d..0c8d17809ab4 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/util/CommonClientUtils.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/util/CommonClientUtils.java
@@ -133,22 +133,14 @@ public class CommonClientUtils {
/**
* Whether log blocks should be written in the native (v2) log format
(standalone native
* files written via {@code HoodieNativeLogFormatWriter}) instead of the
legacy inline
- * log format. The native format is the default for write version >=
{@link HoodieTableVersion#TEN},
- * except for Lance base files.
+ * log format. The native format is the default for write version >=
{@link HoodieTableVersion#TEN}.
*
- * <p>This decision is keyed on the effective write version (i.e. {@code
HoodieWriteConfig#getWriteVersion()})
- * and the effective base file format, consistent with how the inline log
block layout is
- * selected, so that the on-disk format follows what the writer is targeting
during
- * upgrade/downgrade windows. Lance remains on the legacy inline log format
until native Lance log
- * support is complete.
+ * <p>This decision is keyed on the effective write version (i.e. {@code
HoodieWriteConfig#getWriteVersion()}),
+ * so that the on-disk format follows what the writer is targeting during
upgrade/downgrade windows.
*
* @param writeConfig the writer configuration.
- * @param tableConfig the persisted table configuration.
*/
- public static boolean shouldWriteNativeLogs(HoodieWriteConfig writeConfig,
HoodieTableConfig tableConfig) {
- if (getBaseFileFormat(writeConfig, tableConfig) == HoodieFileFormat.LANCE)
{
- return false;
- }
+ public static boolean shouldWriteNativeLogs(HoodieWriteConfig writeConfig) {
return
writeConfig.getWriteVersion().greaterThanOrEquals(HoodieTableVersion.TEN);
}
diff --git
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/action/compact/TestScheduleCompactionActionExecutor.java
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/action/compact/TestScheduleCompactionActionExecutor.java
index 49dbd27048d3..4d9ba937a4eb 100644
---
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/action/compact/TestScheduleCompactionActionExecutor.java
+++
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/table/action/compact/TestScheduleCompactionActionExecutor.java
@@ -77,4 +77,4 @@ class TestScheduleCompactionActionExecutor extends
HoodieCommonTestHarness {
return count;
}
}
-}
\ No newline at end of file
+}
diff --git
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/utils/TestCommonClientUtils.java
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/utils/TestCommonClientUtils.java
index 2a70fc1592ff..573ccf6523be 100644
---
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/utils/TestCommonClientUtils.java
+++
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/utils/TestCommonClientUtils.java
@@ -21,7 +21,6 @@
package org.apache.hudi.utils;
import org.apache.hudi.common.engine.TaskContextSupplier;
-import org.apache.hudi.common.model.HoodieFileFormat;
import org.apache.hudi.common.table.HoodieTableConfig;
import org.apache.hudi.common.table.HoodieTableVersion;
import org.apache.hudi.config.HoodieWriteConfig;
@@ -37,7 +36,6 @@ import java.util.stream.Stream;
import static
org.apache.hudi.util.CommonClientUtils.areTableVersionsCompatible;
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.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
@@ -104,42 +102,25 @@ class TestCommonClientUtils {
assertEquals("0-0-0",
CommonClientUtils.generateWriteToken(taskContextSupplier));
}
- @ParameterizedTest(name = "Write version {0} with base file format {1}
should write native log format: {2}")
+ @ParameterizedTest(name = "Write version {0} should write native log format:
{1}")
@MethodSource("provideWriteVersionNativeLogExpectations")
- void testShouldWriteNativeLogs(HoodieTableVersion writeVersion,
HoodieFileFormat baseFileFormat, boolean expected) {
+ void testShouldWriteNativeLogs(HoodieTableVersion writeVersion, boolean
expected) {
HoodieWriteConfig writeConfig = mock(HoodieWriteConfig.class);
- HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
when(writeConfig.getWriteVersion()).thenReturn(writeVersion);
- when(tableConfig.getBaseFileFormat()).thenReturn(baseFileFormat);
- assertEquals(expected,
CommonClientUtils.shouldWriteNativeLogs(writeConfig, tableConfig));
+ assertEquals(expected,
CommonClientUtils.shouldWriteNativeLogs(writeConfig));
}
private static Stream<Arguments> provideWriteVersionNativeLogExpectations() {
- // Native log format is the default for write version >= TEN, except for
Lance base files.
+ // Native log format is the default for write version >= TEN.
return Stream.of(
- Arguments.of(HoodieTableVersion.SIX, HoodieFileFormat.PARQUET, false),
- Arguments.of(HoodieTableVersion.EIGHT, HoodieFileFormat.PARQUET,
false),
- Arguments.of(HoodieTableVersion.NINE, HoodieFileFormat.PARQUET, false),
- Arguments.of(HoodieTableVersion.TEN, HoodieFileFormat.PARQUET, true),
- Arguments.of(HoodieTableVersion.TEN, HoodieFileFormat.ORC, true),
- Arguments.of(HoodieTableVersion.TEN, HoodieFileFormat.LANCE, false)
+ Arguments.of(HoodieTableVersion.SIX, false),
+ Arguments.of(HoodieTableVersion.EIGHT, false),
+ Arguments.of(HoodieTableVersion.NINE, false),
+ Arguments.of(HoodieTableVersion.TEN, true)
);
}
- @Test
- void testShouldWriteInlineLogFormatForMultiFormatLanceWrites() {
- HoodieWriteConfig writeConfig = mock(HoodieWriteConfig.class);
- HoodieTableConfig tableConfig = mock(HoodieTableConfig.class);
- when(writeConfig.getWriteVersion()).thenReturn(HoodieTableVersion.TEN);
-
when(writeConfig.contains(HoodieWriteConfig.BASE_FILE_FORMAT)).thenReturn(true);
- when(writeConfig.getBaseFileFormat()).thenReturn(HoodieFileFormat.LANCE);
- when(tableConfig.isMultipleBaseFileFormatsEnabled()).thenReturn(true);
- when(tableConfig.getBaseFileFormat()).thenReturn(HoodieFileFormat.PARQUET);
-
- assertFalse(CommonClientUtils.shouldWriteNativeLogs(writeConfig,
tableConfig));
- }
-
@ParameterizedTest(name = "Table version {0} with write version {1} should
be valid: {2}")
@MethodSource("provideValidTableVersionWriteVersionPairs")
void testValidTableVersionWriteVersionPairs(
diff --git
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkWriteHandleFactory.java
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkWriteHandleFactory.java
index 483a8e8687ea..f78266949fce 100644
---
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkWriteHandleFactory.java
+++
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/FlinkWriteHandleFactory.java
@@ -328,7 +328,7 @@ public class FlinkWriteHandleFactory {
final String fileID = bucketInfo.getFileIdPrefix();
final String partitionPath = bucketInfo.getPartitionPath();
final TaskContextSupplier contextSupplier =
table.getTaskContextSupplier();
- if (CommonClientUtils.shouldWriteNativeLogs(config,
table.getMetaClient().getTableConfig())) {
+ if (CommonClientUtils.shouldWriteNativeLogs(config)) {
if (table.requireSortedRecords()) {
recordItr = HoodieRecordUtils.sortRecordsByRecordKey(recordItr);
}
@@ -359,7 +359,7 @@ public class FlinkWriteHandleFactory {
String instantTime,
HoodieTable<T, I, K, O> table,
Iterator<HoodieRecord<T>> recordIterator) {
- if (CommonClientUtils.shouldWriteNativeLogs(config,
table.getMetaClient().getTableConfig())) {
+ if (CommonClientUtils.shouldWriteNativeLogs(config)) {
return new RowDataNativeLogWriteHandle<>(
config,
instantTime,
diff --git
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/HoodieFlinkMergeOnReadTable.java
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/HoodieFlinkMergeOnReadTable.java
index df73654db741..1dfa7c5e3c33 100644
---
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/HoodieFlinkMergeOnReadTable.java
+++
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/table/HoodieFlinkMergeOnReadTable.java
@@ -142,7 +142,7 @@ public class HoodieFlinkMergeOnReadTable<T>
public Iterator<List<WriteStatus>> handleInsertsForLogCompaction(String
instantTime, String partitionPath, String fileId,
Map<String,
HoodieRecord<?>> recordMap,
Map<HoodieLogBlock.HeaderMetadataType, String> header) {
- HoodieWriteHandle appendHandle =
CommonClientUtils.shouldWriteNativeLogs(config,
getMetaClient().getTableConfig())
+ HoodieWriteHandle appendHandle =
CommonClientUtils.shouldWriteNativeLogs(config)
? new HoodieNativeLogAppendHandle(config, instantTime, this,
partitionPath, fileId, recordMap.values().iterator(),
taskContextSupplier, header)
: new HoodieInlineLogAppendHandle(config, instantTime, this,
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkMergeOnReadTable.java
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkMergeOnReadTable.java
index f6aa9c8a53dd..0a49c4b69570 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkMergeOnReadTable.java
+++
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkMergeOnReadTable.java
@@ -203,7 +203,7 @@ public class HoodieSparkMergeOnReadTable<T> extends
HoodieSparkCopyOnWriteTable<
public Iterator<List<WriteStatus>> handleInsertsForLogCompaction(String
instantTime, String partitionPath, String fileId,
Map<String,
HoodieRecord<?>> recordMap,
Map<HoodieLogBlock.HeaderMetadataType, String> header) {
- HoodieWriteHandle appendHandle =
CommonClientUtils.shouldWriteNativeLogs(config,
getMetaClient().getTableConfig())
+ HoodieWriteHandle appendHandle =
CommonClientUtils.shouldWriteNativeLogs(config)
? new HoodieNativeLogAppendHandle(config, instantTime, this,
partitionPath, fileId, recordMap.values().iterator(),
taskContextSupplier, header)
: new HoodieInlineLogAppendHandle(config, instantTime, this,
diff --git
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkTable.java
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkTable.java
index 95e76b863659..4a87e8e96d9b 100644
---
a/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkTable.java
+++
b/hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/table/HoodieSparkTable.java
@@ -18,9 +18,11 @@
package org.apache.hudi.table;
+import org.apache.hudi.SparkFileFormatInternalRecordContext;
import org.apache.hudi.client.WriteStatus;
import org.apache.hudi.common.data.HoodieData;
import org.apache.hudi.common.engine.HoodieEngineContext;
+import org.apache.hudi.common.engine.RecordContext;
import org.apache.hudi.common.model.HoodieFailedWritesCleaningPolicy;
import org.apache.hudi.common.model.HoodieKey;
import org.apache.hudi.common.model.HoodieRecord;
@@ -83,6 +85,16 @@ public abstract class HoodieSparkTable<T>
return SparkHoodieIndexFactory.createIndex(config);
}
+ @Override
+ public RecordContext<?> getRecordContextForWrite() {
+ if (config.getRecordMerger().getRecordType() ==
HoodieRecord.HoodieRecordType.SPARK) {
+ // The engine context is transient and a serialized table falls back to
HoodieLocalEngineContext on executors.
+ // Select the Spark record context from the configured record type
instead of the table's runtime context.
+ return SparkFileFormatInternalRecordContext.getFieldAccessorInstance();
+ }
+ return super.getRecordContextForWrite();
+ }
+
/**
* Fetch instance of {@link HoodieTableMetadataWriter}.
*
diff --git
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
index 95156d417909..7671c64bc219 100644
---
a/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
+++
b/hudi-client/hudi-spark-client/src/main/scala/org/apache/hudi/SparkFileFormatInternalRowReaderContext.scala
@@ -143,11 +143,12 @@ class
SparkFileFormatInternalRowReaderContext(baseFileReader: SparkColumnarFileR
// Parquet stores VECTOR as FIXED_LEN_BYTE_ARRAY, so the reader needs
BinaryType
// and we decode back to ArrayType below. Lance returns ArrayType
natively, so skip
- // the rewrite only for Lance base files; log files always go through the
rewrite path.
- val isLogFile = FSUtils.isLogFile(filePath)
- val isLanceBaseFile = !isLogFile && FSUtils.isBaseFile(filePath) &&
- tableConfig.getBaseFileFormat == HoodieFileFormat.LANCE
- val vectorColumnInfo: Map[Int, HoodieSchema.Vector] = if (isLanceBaseFile)
{
+ // the rewrite only for Lance files including native log files; inline log
files
+ // always go through the rewrite path.
+ val isInlineLog = FSUtils.isInlineLogFile(filePath)
+ val isLanceFile = !isInlineLog &&
+ HoodieFileFormat.fromFileExtension(filePath.getFileExtension) ==
HoodieFileFormat.LANCE
+ val vectorColumnInfo: Map[Int, HoodieSchema.Vector] = if (isLanceFile) {
Map.empty
} else {
SparkFileFormatInternalRowReaderContext.detectVectorColumns(requiredSchema)
@@ -159,7 +160,7 @@ class
SparkFileFormatInternalRowReaderContext(baseFileReader: SparkColumnarFileR
}
val (readSchema, readFilters) =
getSchemaAndFiltersForRead(parquetReadStructType, hasRowIndexField)
- if (isLogFile) {
+ if (FSUtils.isLogFile(filePath)) {
// NOTE: now only primary key based filtering is supported for log files
// Position-based merging pairs log records with the RECORD_POSITIONS
bitmap by index (see
// PositionBasedFileGroupRecordBuffer), which requires the record stream
to contain every
@@ -169,7 +170,7 @@ class
SparkFileFormatInternalRowReaderContext(baseFileReader: SparkColumnarFileR
// Variant alignment happens later via getLogBlockRecordProjection in
the merge buffer.
// Log files reach this method either as an inline parquet data block or
as a native log file.
// Inline paths do not carry the parquet extension, so resolve them by
the log block contract.
- val fileFormat = if (FSUtils.isInlineLogFile(filePath.getName)) {
+ val fileFormat = if (isInlineLog) {
HoodieFileFormat.PARQUET
} else {
HoodieFileFormat.fromFileExtension(filePath.getFileExtension)
diff --git
a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/TestHoodieSparkTable.java
b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/TestHoodieSparkTable.java
index fdadfe2c3445..0820dd08796f 100644
---
a/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/TestHoodieSparkTable.java
+++
b/hudi-client/hudi-spark-client/src/test/java/org/apache/hudi/table/TestHoodieSparkTable.java
@@ -18,11 +18,15 @@
package org.apache.hudi.table;
+import org.apache.hudi.DefaultSparkRecordMerger;
+import org.apache.hudi.common.model.HoodieAvroRecordMerger;
+import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.model.HoodieTableType;
import org.apache.hudi.common.table.HoodieTableConfig;
import org.apache.hudi.common.table.HoodieTableMetaClient;
import org.apache.hudi.common.table.marker.MarkerType;
import org.apache.hudi.common.testutils.HoodieCommonTestHarness;
+import org.apache.hudi.common.testutils.HoodieTestUtils;
import org.apache.hudi.common.util.HoodieStorageUtils;
import org.apache.hudi.config.HoodieWriteConfig;
import org.apache.hudi.exception.HoodieException;
@@ -31,6 +35,7 @@ import org.apache.hudi.storage.StorageConfiguration;
import org.apache.hudi.storage.StoragePath;
import org.apache.hudi.table.marker.WriteMarkers;
+import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.EnumSource;
@@ -40,9 +45,11 @@ import java.util.ArrayList;
import java.util.Arrays;
import java.util.HashSet;
import java.util.List;
+import java.util.Properties;
import java.util.Set;
import static
org.apache.hudi.common.testutils.HoodieTestUtils.getDefaultStorageConf;
+import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
@@ -54,6 +61,31 @@ public class TestHoodieSparkTable extends
HoodieCommonTestHarness {
private static final StorageConfiguration<?> CONF = getDefaultStorageConf();
+ @Test
+ void testRecordContextForWriteUsesConfiguredRecordType() throws IOException {
+ initPath();
+ Properties tableProps = new Properties();
+ tableProps.setProperty(HoodieTableConfig.NAME.key(),
"record-context-test");
+ HoodieTableMetaClient metaClient = HoodieTestUtils.init(
+ CONF, basePath, HoodieTableType.COPY_ON_WRITE, tableProps);
+
+ HoodieWriteConfig sparkRecordConfig = HoodieWriteConfig.newBuilder()
+ .withPath(basePath)
+ .withRecordMergeImplClasses(DefaultSparkRecordMerger.class.getName())
+ .build();
+ HoodieSparkTable sparkRecordTable =
HoodieSparkTable.create(sparkRecordConfig, getEngineContext(), metaClient);
+ assertEquals(HoodieRecord.HoodieRecordType.SPARK,
+ sparkRecordTable.getRecordContextForWrite().getEngineRecordType());
+
+ HoodieWriteConfig avroRecordConfig = HoodieWriteConfig.newBuilder()
+ .withPath(basePath)
+ .withRecordMergeImplClasses(HoodieAvroRecordMerger.class.getName())
+ .build();
+ HoodieSparkTable avroRecordTable =
HoodieSparkTable.create(avroRecordConfig, getEngineContext(), metaClient);
+ assertEquals(HoodieRecord.HoodieRecordType.AVRO,
+ avroRecordTable.getRecordContextForWrite().getEngineRecordType());
+ }
+
@ParameterizedTest
@EnumSource(DeleteFailureType.class)
public void testDeleteFailureDuringMarkerReconciliation(DeleteFailureType
failureType) throws IOException {
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/avro/HoodieAvroReaderContext.java
b/hudi-common/src/main/java/org/apache/hudi/common/avro/HoodieAvroReaderContext.java
index 9735e688d9a2..788810b7ba3f 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/avro/HoodieAvroReaderContext.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/avro/HoodieAvroReaderContext.java
@@ -133,7 +133,16 @@ public class HoodieAvroReaderContext extends
HoodieReaderContext<IndexedRecord>
HoodieTableConfig tableConfig,
String payloadClassName,
TypedProperties props) {
- this(storageConfiguration, tableConfig, Option.empty(), Option.empty(),
Collections.emptyMap(), payloadClassName, new HoodieConfig(props));
+ this(storageConfiguration, tableConfig, Option.empty(), payloadClassName,
props);
+ }
+
+ public HoodieAvroReaderContext(
+ StorageConfiguration<?> storageConfiguration,
+ HoodieTableConfig tableConfig,
+ Option<InstantRange> instantRangeOpt,
+ String payloadClassName,
+ TypedProperties props) {
+ this(storageConfiguration, tableConfig, instantRangeOpt, Option.empty(),
Collections.emptyMap(), payloadClassName, new HoodieConfig(props));
}
private HoodieAvroReaderContext(
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/engine/AvroReaderContextFactory.java
b/hudi-common/src/main/java/org/apache/hudi/common/engine/AvroReaderContextFactory.java
index 0d6b643387dd..db1c72cb5874 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/engine/AvroReaderContextFactory.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/engine/AvroReaderContextFactory.java
@@ -21,6 +21,8 @@ package org.apache.hudi.common.engine;
import org.apache.hudi.common.avro.HoodieAvroReaderContext;
import org.apache.hudi.common.config.TypedProperties;
import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.log.InstantRange;
+import org.apache.hudi.common.util.Option;
import org.apache.avro.generic.IndexedRecord;
@@ -30,20 +32,28 @@ import org.apache.avro.generic.IndexedRecord;
public class AvroReaderContextFactory implements
ReaderContextFactory<IndexedRecord> {
private final HoodieTableMetaClient metaClient;
private final String payloadClassName;
+ private final Option<InstantRange> instantRange;
private final TypedProperties props;
public AvroReaderContextFactory(HoodieTableMetaClient metaClient,
TypedProperties props) {
- this(metaClient, metaClient.getTableConfig().getPayloadClass(), props);
+ this(metaClient, metaClient.getTableConfig().getPayloadClass(),
Option.empty(), props);
}
public AvroReaderContextFactory(HoodieTableMetaClient metaClient, String
payloadClassName, TypedProperties props) {
+ this(metaClient, payloadClassName, Option.empty(), props);
+ }
+
+ public AvroReaderContextFactory(HoodieTableMetaClient metaClient, String
payloadClassName,
+ Option<InstantRange> instantRange,
TypedProperties props) {
this.metaClient = metaClient;
this.payloadClassName = payloadClassName;
+ this.instantRange = instantRange;
this.props = props;
}
@Override
public HoodieReaderContext<IndexedRecord> getContext() {
- return new HoodieAvroReaderContext(metaClient.getStorageConf(),
metaClient.getTableConfig(), payloadClassName, props);
+ return new HoodieAvroReaderContext(
+ metaClient.getStorageConf(), metaClient.getTableConfig(),
instantRange, payloadClassName, props);
}
}
diff --git a/hudi-common/src/main/java/org/apache/hudi/common/fs/FSUtils.java
b/hudi-common/src/main/java/org/apache/hudi/common/fs/FSUtils.java
index 86846d51db79..75b6eada4f55 100644
--- a/hudi-common/src/main/java/org/apache/hudi/common/fs/FSUtils.java
+++ b/hudi-common/src/main/java/org/apache/hudi/common/fs/FSUtils.java
@@ -456,8 +456,11 @@ public class FSUtils {
.orElse(false);
}
- public static boolean isInlineLogFile(String fileName) {
- return !isNativeLogFile(fileName);
+ public static boolean isInlineLogFile(StoragePath filePath) {
+ String scheme = filePath.toUri().getScheme();
+ String fileName = InLineFSUtils.SCHEME.equals(scheme)
+ ? InLineFSUtils.getOuterFilePathFromInlinePath(filePath).getName() :
filePath.getName();
+ return FileNameParser.parseInlineLogFile(fileName).isPresent();
}
public static boolean isNativeLogFile(String fileName) {
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/fs/FileNameParser.java
b/hudi-common/src/main/java/org/apache/hudi/common/fs/FileNameParser.java
index 16ef14bfcdd3..a89e774eaee0 100644
--- a/hudi-common/src/main/java/org/apache/hudi/common/fs/FileNameParser.java
+++ b/hudi-common/src/main/java/org/apache/hudi/common/fs/FileNameParser.java
@@ -152,8 +152,23 @@ public final class FileNameParser {
return nativeLogFileName;
}
- Matcher matcher = INLINE_LOG_FILE_PATTERN.matcher(actualFileName);
- return matcher.matches() ?
Option.of(LogFileName.fromInlineLogFile(matcher)) : Option.empty();
+ return parseInlineLogFileFromActualFileName(actualFileName);
+ }
+
+ /**
+ * Parses only an inline log file name.
+ *
+ * <p>Use {@link #parseLogFile(String)} when both inline and native log
formats are acceptable.</p>
+ *
+ * @param fileName file name or full path
+ * @return decoded inline log file name when the input matches the inline
log-file format
+ */
+ public static Option<LogFileName> parseInlineLogFile(String fileName) {
+ if (StringUtils.isNullOrEmpty(fileName)) {
+ return Option.empty();
+ }
+
+ return parseInlineLogFileFromActualFileName(getActualFileName(fileName));
}
/**
@@ -192,6 +207,11 @@ public final class FileNameParser {
return
matchNativeLogFileFromActualFileName(actualFileName).map(LogFileName::fromNativeLogFile);
}
+ private static Option<LogFileName>
parseInlineLogFileFromActualFileName(String actualFileName) {
+ Matcher matcher = INLINE_LOG_FILE_PATTERN.matcher(actualFileName);
+ return matcher.matches() ?
Option.of(LogFileName.fromInlineLogFile(matcher)) : Option.empty();
+ }
+
private static Option<Matcher> matchNativeLogFileFromActualFileName(String
actualFileName) {
Matcher matcher = NATIVE_LOG_FILE_PATTERN.matcher(actualFileName);
return matcher.matches() ? Option.of(matcher) : Option.empty();
diff --git
a/hudi-common/src/test/java/org/apache/hudi/common/model/TestHoodieLogFile.java
b/hudi-common/src/test/java/org/apache/hudi/common/model/TestHoodieLogFile.java
index d041b6a850ec..9bf56ab21861 100644
---
a/hudi-common/src/test/java/org/apache/hudi/common/model/TestHoodieLogFile.java
+++
b/hudi-common/src/test/java/org/apache/hudi/common/model/TestHoodieLogFile.java
@@ -19,6 +19,7 @@
package org.apache.hudi.common.model;
import org.apache.hudi.common.fs.FSUtils;
+import org.apache.hudi.common.fs.FileNameParser;
import org.apache.hudi.storage.StoragePath;
import org.apache.hudi.storage.StoragePathInfo;
@@ -109,6 +110,16 @@ public class TestHoodieLogFile {
assertEquals(1, FSUtils.getTaskAttemptIdFromLogPath(nativeLogPath));
}
+ @Test
+ void inlineLogFileDetectionOnlyMatchesInlineLogs() {
+ assertTrue(FileNameParser.parseInlineLogFile(pathStr).isPresent());
+
assertFalse(FileNameParser.parseInlineLogFile(nativeLogPath(LogExtensions.DATA_LOG_EXTENSION,
2)).isPresent());
+ assertFalse(FileNameParser.parseInlineLogFile(fileId +
"_1-0-1_100.parquet").isPresent());
+ assertTrue(FSUtils.isInlineLogFile(new StoragePath(pathStr)));
+ assertFalse(FSUtils.isInlineLogFile(new
StoragePath(nativeLogPath(LogExtensions.DATA_LOG_EXTENSION, 2))));
+ assertFalse(FSUtils.isInlineLogFile(new StoragePath(fileId +
"_1-0-1_100.parquet")));
+ }
+
@Test
void createFromNativeDeleteParquetLogFile() {
String nativeDeleteLogPathStr = "file:///tmp/hoodie/2021/01/01/"
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/compact/handler/MetadataTableCompactHandler.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/compact/handler/MetadataTableCompactHandler.java
index db75e7efe721..58f01a1981ab 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/compact/handler/MetadataTableCompactHandler.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/compact/handler/MetadataTableCompactHandler.java
@@ -105,13 +105,16 @@ public class MetadataTableCompactHandler extends
DataTableCompactHandler {
metaClient.reload();
}
Option<InstantRange> instantRange =
CompactHelpers.getInstance().getInstantRange(metaClient);
+ String payloadClass =
ConfigUtils.getPayloadClass(writeClient.getConfig().getProps());
+ AvroReaderContextFactory readerContextFactory = new
AvroReaderContextFactory(
+ metaClient, payloadClass, instantRange,
writeClient.getConfig().getProps());
List<WriteStatus> writeStatuses = compactor.logCompact(
writeClient.getConfig(),
event.getOperation(),
event.getCompactionInstantTime(),
- instantRange,
table,
- table.getTaskContextSupplier());
+ table.getTaskContextSupplier(),
+ readerContextFactory.getContext());
compactionMetrics.endCompaction();
collector.collect(createCommitEvent(event, writeStatuses));
}
diff --git
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FlinkRowDataReaderContext.java
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FlinkRowDataReaderContext.java
index bd5b05e1d163..46251866e660 100644
---
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FlinkRowDataReaderContext.java
+++
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/table/format/FlinkRowDataReaderContext.java
@@ -107,7 +107,7 @@ public class FlinkRowDataReaderContext extends
HoodieReaderContext<RowData> {
// Log files only reach this method for parquet data blocks; base files
are resolved by their extension.
// Format-specific handling lives in the readers themselves, so this
method stays format-agnostic.
- boolean isInlineLogFile = isLogFile &&
FSUtils.isInlineLogFile(filePath.getName());
+ boolean isInlineLogFile = FSUtils.isInlineLogFile(filePath);
HoodieFileFormat format = isInlineLogFile ? HoodieFileFormat.PARQUET :
HoodieFileFormat.fromFileExtension(filePath.getFileExtension());
HoodieRowDataFileReader reader = (HoodieRowDataFileReader)
HoodieIOFactory.getIOFactory(storage)
.getReaderFactory(HoodieRecord.HoodieRecordType.FLINK)
diff --git
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestLanceDataSource.scala
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestLanceDataSource.scala
index b281b83bf6c5..532968ec1fc9 100644
---
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestLanceDataSource.scala
+++
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestLanceDataSource.scala
@@ -17,8 +17,8 @@
package org.apache.hudi.functional
+import org.apache.hudi.{DefaultSparkRecordMerger, HoodieCLIUtils}
import org.apache.hudi.DataSourceWriteOptions._
-import org.apache.hudi.DefaultSparkRecordMerger
import org.apache.hudi.blob.BlobTestHelpers
import org.apache.hudi.common.config.{HoodieCommonConfig, HoodieMetadataConfig}
import org.apache.hudi.common.engine.HoodieLocalEngineContext
@@ -27,6 +27,7 @@ import org.apache.hudi.common.schema.HoodieSchema
import org.apache.hudi.common.table.{HoodieTableConfig, HoodieTableMetaClient}
import org.apache.hudi.common.table.view.{FileSystemViewManager,
FileSystemViewStorageConfig}
import org.apache.hudi.common.testutils.HoodieTestUtils
+import org.apache.hudi.common.util.{Option => HOption}
import org.apache.hudi.config.HoodieWriteConfig
import org.apache.hudi.io.storage.HoodieSparkLanceReader
import org.apache.hudi.metadata.MetadataPartitionType
@@ -2005,6 +2006,31 @@ class TestLanceDataSource extends
HoodieSparkClientTestBase {
if (tableType == HoodieTableType.MERGE_ON_READ) {
assertCompactionCommitPresence(tablePath, expectPresent = false,
"No compaction commit should be present before max.delta.commits=6
threshold is reached")
+
+ // The update and delete above create native Lance log files. Execute
log compaction explicitly to
+ // verify that its executor-side reader context can read those logs and
its writer keeps Spark records.
+ spark.sql(s"alter table $tableName set tblproperties (" +
+ "'hoodie.log.compaction.enable' = 'true', " +
+ "'hoodie.log.compaction.blocks.threshold' = '1')")
+ val client = HoodieCLIUtils.createHoodieWriteClient(spark, tablePath,
Map.empty, Option(tableName))
+ val logCompactionInstant = client.scheduleLogCompaction(HOption.empty())
+ assertTrue(logCompactionInstant.isPresent, "Native Lance log files
should be scheduled for log compaction")
+ client.logCompact(logCompactionInstant.get(), true)
+
+ val tableFiles = Files.walk(Paths.get(tablePath))
+ try {
+ assertTrue(tableFiles.anyMatch(path =>
path.toString.contains(logCompactionInstant.get())
+ && path.toString.endsWith(".log.lance")),
+ "Log compaction should produce a native Lance log file")
+ } finally {
+ tableFiles.close()
+ }
+
+ checkAnswer(s"select id, name, age, score, dt from $tableName order by
id")(
+ Seq(1, "Alice", 31, 99.9, "2025-01-01"),
+ Seq(2, "Bob", 25, 87.3, "2025-01-02"),
+ Seq(4, "Diana", 40, null, "2025-01-01")
+ )
}
// Test 6: INSERT with static partition (only for partitioned tables)