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 d56e6747f43c feat(flink): enable LSM reader for Flink metadata table
(#19373)
d56e6747f43c is described below
commit d56e6747f43c7388d0cdd6e27ef5afab82de360e
Author: Shuo Cheng <[email protected]>
AuthorDate: Tue Jul 28 09:55:08 2026 +0800
feat(flink): enable LSM reader for Flink metadata table (#19373)
* feat(flink): enable LSM reader for Flink metadata table
---
.../metadata/HoodieBackedTableMetadataWriter.java | 12 ++
.../TestHoodieBackedTableMetadataWriter.java | 16 +++
.../table/read/lsm/HoodieLsmFileGroupReader.java | 6 +
.../hudi/metadata/HoodieBackedTableMetadata.java | 52 +++++---
.../read/lsm/TestHoodieLsmFileGroupReader.java | 134 ++++++++++++++++++++-
.../sink/TestStreamWriteOperatorCoordinator.java | 46 +++++++
.../sink/compact/ITTestHoodieFlinkCompactor.java | 2 +-
7 files changed, 251 insertions(+), 17 deletions(-)
diff --git
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java
index 336ebad1fa46..83d63b35a198 100644
---
a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java
+++
b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadataWriter.java
@@ -46,7 +46,9 @@ import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.model.HoodieRecordLocation;
import org.apache.hudi.common.model.HoodieTableType;
import org.apache.hudi.common.model.HoodieWriteStat;
+import org.apache.hudi.common.table.HoodieTableConfig;
import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.HoodieTableVersion;
import org.apache.hudi.common.table.timeline.HoodieActiveTimeline;
import org.apache.hudi.common.table.timeline.HoodieInstant;
import org.apache.hudi.common.table.timeline.HoodieInstantTimeGenerator;
@@ -562,6 +564,7 @@ public abstract class HoodieBackedTableMetadataWriter<I, O>
implements HoodieTab
.setArchiveLogFolder(getTimelineHistoryPath())
.setPayloadClassName(HoodieMetadataPayload.class.getName())
.setBaseFileFormat(HoodieFileFormat.HFILE.toString())
+ .setTableStorageLayout(getMetadataTableStorageLayout(getEngineType(),
dataWriteConfig.getWriteVersion()))
.setRecordKeyFields(RECORD_KEY_FIELD_NAME)
.setPopulateMetaFields(DEFAULT_METADATA_POPULATE_META_FIELDS)
.setKeyGeneratorClassProp(HoodieTableMetadataKeyGenerator.class.getCanonicalName())
@@ -574,6 +577,15 @@ public abstract class HoodieBackedTableMetadataWriter<I,
O> implements HoodieTab
.build();
}
+ static String getMetadataTableStorageLayout(EngineType engineType,
HoodieTableVersion tableVersion) {
+ // Return null for all other cases so TableBuilder does not persist the
property. The table config
+ // still resolves the absent property to the default layout, while
preserving the existing on-disk
+ // configuration for other engines and older table versions.
+ return engineType == EngineType.FLINK &&
tableVersion.greaterThanOrEquals(HoodieTableVersion.TEN)
+ ? HoodieTableConfig.TableStorageLayout.LSM_TREE.configValue()
+ : null;
+ }
+
/**
* Function to find hoodie partitions and list files in them in parallel.
*
diff --git
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/TestHoodieBackedTableMetadataWriter.java
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/TestHoodieBackedTableMetadataWriter.java
index e3f4ae662182..9f777c91d9bd 100644
---
a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/TestHoodieBackedTableMetadataWriter.java
+++
b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metadata/TestHoodieBackedTableMetadataWriter.java
@@ -24,12 +24,14 @@ import org.apache.hudi.client.BaseHoodieWriteClient;
import org.apache.hudi.common.config.HoodieMetadataConfig;
import org.apache.hudi.common.config.HoodieTableServiceManagerConfig;
import org.apache.hudi.common.data.HoodieData;
+import org.apache.hudi.common.engine.EngineType;
import org.apache.hudi.common.engine.HoodieEngineContext;
import org.apache.hudi.common.model.HoodieCommitMetadata;
import org.apache.hudi.common.model.HoodieFailedWritesCleaningPolicy;
import org.apache.hudi.common.model.HoodieRecord;
import org.apache.hudi.common.table.HoodieTableConfig;
import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.HoodieTableVersion;
import org.apache.hudi.common.table.timeline.HoodieActiveTimeline;
import org.apache.hudi.common.table.timeline.HoodieInstant;
import org.apache.hudi.common.table.timeline.HoodieTimeline;
@@ -67,7 +69,9 @@ import java.util.stream.Stream;
import static
org.apache.hudi.common.testutils.HoodieTestUtils.INSTANT_GENERATOR;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -85,6 +89,18 @@ import static org.mockito.Mockito.verifyNoMoreInteractions;
import static org.mockito.Mockito.when;
class TestHoodieBackedTableMetadataWriter {
+ @Test
+ void tableStorageLayoutForMetadataTable() {
+ assertEquals("lsm_tree",
HoodieBackedTableMetadataWriter.getMetadataTableStorageLayout(
+ EngineType.FLINK, HoodieTableVersion.TEN));
+ assertNull(HoodieBackedTableMetadataWriter.getMetadataTableStorageLayout(
+ EngineType.FLINK, HoodieTableVersion.NINE));
+ assertNull(HoodieBackedTableMetadataWriter.getMetadataTableStorageLayout(
+ EngineType.SPARK, HoodieTableVersion.TEN));
+ assertNull(HoodieBackedTableMetadataWriter.getMetadataTableStorageLayout(
+ EngineType.JAVA, HoodieTableVersion.TEN));
+ }
+
@Test
void completeStreamingCommitSkipsAlreadyCompletedMetadataInstant() {
String instantTime = "20260709120000000";
diff --git
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/HoodieLsmFileGroupReader.java
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/HoodieLsmFileGroupReader.java
index 600a615b5443..da4154993f18 100644
---
a/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/HoodieLsmFileGroupReader.java
+++
b/hudi-common/src/main/java/org/apache/hudi/common/table/read/lsm/HoodieLsmFileGroupReader.java
@@ -28,6 +28,7 @@ import org.apache.hudi.common.schema.HoodieSchema;
import org.apache.hudi.common.schema.internal.InternalSchema;
import org.apache.hudi.common.table.HoodieTableConfig;
import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.log.InstantRange;
import org.apache.hudi.common.table.read.BaseFileUpdateCallback;
import org.apache.hudi.common.table.read.BufferedRecord;
import org.apache.hudi.common.table.read.FileGroupReaderSchemaHandler;
@@ -141,6 +142,11 @@ public final class HoodieLsmFileGroupReader<T> implements
HoodieRecordReader<T>
.sortOutputs(false)
.inflightInstantsAllowed(allowInflightInstants)
.build();
+ // filter log files by instant range.
+ if (logFiles != null && readerContext.getInstantRange().isPresent()) {
+ InstantRange instantRange = readerContext.getInstantRange().get();
+ logFiles = logFiles.filter(logFile ->
instantRange.isInRange(logFile.getDeltaCommitTime()));
+ }
this.inputSplit = InputSplit.builder()
.baseFileOption(baseFileOption)
.logFileStream(logFiles)
diff --git
a/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadata.java
b/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadata.java
index 46cf00a5326a..3d31a2ccfed3 100644
---
a/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadata.java
+++
b/hudi-common/src/main/java/org/apache/hudi/metadata/HoodieBackedTableMetadata.java
@@ -34,6 +34,7 @@ import org.apache.hudi.common.expression.Expression;
import org.apache.hudi.common.expression.Literal;
import org.apache.hudi.common.expression.Predicate;
import org.apache.hudi.common.expression.Predicates;
+import org.apache.hudi.common.fs.FSUtils;
import org.apache.hudi.common.function.SerializableBiFunction;
import org.apache.hudi.common.function.SerializableFunction;
import org.apache.hudi.common.function.SerializableFunctionUnchecked;
@@ -52,6 +53,7 @@ import
org.apache.hudi.common.table.read.FileGroupReaderSchemaHandler;
import org.apache.hudi.common.table.read.HoodieFileGroupReader;
import org.apache.hudi.common.table.read.buffer.FileGroupRecordBufferLoader;
import
org.apache.hudi.common.table.read.buffer.ReusableFileGroupRecordBufferLoader;
+import org.apache.hudi.common.table.read.lsm.HoodieLsmFileGroupReader;
import org.apache.hudi.common.table.timeline.HoodieInstant;
import org.apache.hudi.common.table.view.HoodieTableFileSystemView;
import org.apache.hudi.common.util.ConfigUtils;
@@ -571,9 +573,10 @@ public class HoodieBackedTableMetadata extends
BaseTableMetadata {
.explicitInstants(validInstantTimestamps).build());
// If reuse is enabled and full scan is allowed for the partition, we can
reuse the file readers for base files and the reader context for the log files.
+ boolean shouldReuse = reuse &&
isFullScanAllowedForPartition(fileSlice.getPartitionPath());
+ boolean useLsmReader = !shouldReuse &&
shouldUseLsmReader(metadataMetaClient, fileSlice);
Map<StoragePath, HoodieAvroFileReader> baseFileReaders =
Collections.emptyMap();
ReusableFileGroupRecordBufferLoader<IndexedRecord> recordBufferLoader =
null;
- boolean shouldReuse = reuse &&
isFullScanAllowedForPartition(fileSlice.getPartitionPath());
TypedProperties fileGroupReaderProps =
ConfigUtils.buildFileGroupReaderProperties(metadataConfig, shouldReuse);
if (shouldReuse) {
Pair<HoodieAvroFileReader,
ReusableFileGroupRecordBufferLoader<IndexedRecord>> readers =
@@ -606,20 +609,39 @@ public class HoodieBackedTableMetadata extends
BaseTableMetadata {
baseFileReaders,
fileGroupReaderProps);
- HoodieFileGroupReader<IndexedRecord> fileGroupReader =
HoodieFileGroupReader.<IndexedRecord>builder()
- .withReaderContext(readerContext)
- .withHoodieTableMetaClient(metadataMetaClient)
- .withLatestCommitTime(latestMetadataInstantTime)
- .withBaseFileOption(fileSlice.getBaseFile())
- .withLogFiles(fileSlice.getLogFiles())
- .withPartitionPath(fileSlice.getPartitionPath())
- .withDataSchema(SCHEMA)
- .withRequestedSchema(SCHEMA)
- .withProps(fileGroupReaderProps)
- .withRecordBufferLoader(recordBufferLoader)
- .build();
-
- return fileGroupReader.getClosableIterator();
+ if (useLsmReader) {
+ return HoodieLsmFileGroupReader.<IndexedRecord>builder()
+ .withReaderContext(readerContext)
+ .withHoodieTableMetaClient(metadataMetaClient)
+ .withLatestCommitTime(latestMetadataInstantTime)
+ .withBaseFileOption(fileSlice.getBaseFile())
+ .withLogFiles(fileSlice.getLogFiles())
+ .withPartitionPath(fileSlice.getPartitionPath())
+ .withDataSchema(SCHEMA)
+ .withRequestedSchema(SCHEMA)
+ .withProps(fileGroupReaderProps)
+ .build()
+ .getClosableIterator();
+ } else {
+ return HoodieFileGroupReader.<IndexedRecord>builder()
+ .withReaderContext(readerContext)
+ .withHoodieTableMetaClient(metadataMetaClient)
+ .withLatestCommitTime(latestMetadataInstantTime)
+ .withBaseFileOption(fileSlice.getBaseFile())
+ .withLogFiles(fileSlice.getLogFiles())
+ .withPartitionPath(fileSlice.getPartitionPath())
+ .withDataSchema(SCHEMA)
+ .withRequestedSchema(SCHEMA)
+ .withProps(fileGroupReaderProps)
+ .withRecordBufferLoader(recordBufferLoader)
+ .build()
+ .getClosableIterator();
+ }
+ }
+
+ private static boolean shouldUseLsmReader(HoodieTableMetaClient metaClient,
FileSlice fileSlice) {
+ return metaClient.getTableConfig().isLSMTreeStorageLayout()
+ && fileSlice.getLogFiles().allMatch(logFile ->
FSUtils.isNativeLogFile(logFile.getFileName()));
}
private ReusableFileGroupRecordBufferLoader<IndexedRecord>
buildReusableRecordBufferLoader(FileSlice fileSlice, String
latestMetadataInstantTime,
diff --git
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestHoodieLsmFileGroupReader.java
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestHoodieLsmFileGroupReader.java
index 4bc7518a0577..67d003d36ed4 100644
---
a/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestHoodieLsmFileGroupReader.java
+++
b/hudi-common/src/test/java/org/apache/hudi/common/table/read/lsm/TestHoodieLsmFileGroupReader.java
@@ -24,6 +24,7 @@ import org.apache.hudi.common.config.RecordMergeMode;
import org.apache.hudi.common.config.TypedProperties;
import org.apache.hudi.common.engine.HoodieReaderContext;
import org.apache.hudi.common.model.HoodieAvroIndexedRecord;
+import org.apache.hudi.common.model.HoodieBaseFile;
import org.apache.hudi.common.model.HoodieEmptyRecord;
import org.apache.hudi.common.model.HoodieFileFormat;
import org.apache.hudi.common.model.HoodieKey;
@@ -35,12 +36,14 @@ import org.apache.hudi.common.schema.HoodieSchemaField;
import org.apache.hudi.common.schema.HoodieSchemaType;
import org.apache.hudi.common.table.HoodieTableConfig;
import org.apache.hudi.common.table.HoodieTableMetaClient;
+import org.apache.hudi.common.table.log.InstantRange;
import org.apache.hudi.common.table.read.BufferedRecord;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.common.util.collection.ClosableIterator;
import org.apache.hudi.storage.HoodieStorage;
import org.apache.hudi.storage.StorageConfiguration;
import org.apache.hudi.storage.StoragePath;
+import org.apache.hudi.storage.StoragePathInfo;
import org.apache.avro.generic.GenericData;
import org.apache.avro.generic.IndexedRecord;
@@ -58,8 +61,11 @@ import static org.junit.jupiter.api.Assertions.assertFalse;
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.anyLong;
+import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.spy;
import static org.mockito.Mockito.when;
class TestHoodieLsmFileGroupReader {
@@ -68,6 +74,12 @@ class TestHoodieLsmFileGroupReader {
HoodieSchemaField.of("id", HoodieSchema.create(HoodieSchemaType.STRING)),
HoodieSchemaField.of("value",
HoodieSchema.create(HoodieSchemaType.STRING)),
HoodieSchemaField.of("ts", HoodieSchema.create(HoodieSchemaType.LONG))));
+ private static final HoodieSchema SCHEMA_WITH_COMMIT_TIME =
HoodieSchema.createRecord(
+ "lsm_record_with_commit_time", null, null, Arrays.asList(
+ HoodieSchemaField.of(HoodieRecord.COMMIT_TIME_METADATA_FIELD,
HoodieSchema.create(HoodieSchemaType.STRING)),
+ HoodieSchemaField.of("id",
HoodieSchema.create(HoodieSchemaType.STRING)),
+ HoodieSchemaField.of("value",
HoodieSchema.create(HoodieSchemaType.STRING)),
+ HoodieSchemaField.of("ts",
HoodieSchema.create(HoodieSchemaType.LONG))));
private HoodieTableMetaClient metaClient;
private TypedProperties props;
@@ -165,6 +177,77 @@ class TestHoodieLsmFileGroupReader {
.build());
}
+ @Test
+ void testNativeDataAndDeleteLogsAreFilteredByInstantRange() throws
IOException {
+ List<HoodieLogFile> logFiles = Arrays.asList(
+ logFile("file1_1-0-1_002_1.log.parquet"),
+ logFile("file1_1-0-1_003_1.deletes.parquet"));
+
+ // A matching native data log remains in the split, so a non-zero base
offset is rejected.
+ assertThrows(IllegalArgumentException.class, () -> reader(
+ context(Option.of(exactMatchRange("002"))), Option.empty(), logFiles,
1L));
+
+ // Native delete logs use the same instant encoded in their file name.
+ assertThrows(IllegalArgumentException.class, () -> reader(
+ context(Option.of(exactMatchRange("003"))), Option.empty(), logFiles,
1L));
+
+ // When every log is outside the range, the split is treated as having no
logs.
+ HoodieReaderContext<IndexedRecord> filteredContext =
context(Option.of(exactMatchRange("001")));
+ try (HoodieLsmFileGroupReader<IndexedRecord> ignored =
+ reader(filteredContext, Option.empty(), logFiles, 1L)) {
+ assertFalse(filteredContext.getHasLogFiles());
+ }
+
+ // An absent instant range leaves the original log stream unchanged.
+ assertThrows(IllegalArgumentException.class, () -> reader(
+ context(Option.empty()), Option.empty(), logFiles, 1L));
+ }
+
+ @Test
+ void testAllFilteredLogsUseDataTableBaseFileOnlyPath() throws IOException {
+ HoodieReaderContext<IndexedRecord> readerContext =
spy(context(Option.of(exactMatchRange("002"))));
+ StoragePathInfo baseFilePathInfo =
pathInfo("/tmp/file1_1-0-1_001.parquet");
+ doReturn(ClosableIterator.wrap(Arrays.asList(
+ recordWithCommitTime("001", "a", "old", 1),
+ recordWithCommitTime("002", "b", "kept", 2)).iterator()))
+ .when(readerContext).getFileRecordIterator(
+ eq(baseFilePathInfo), anyLong(), anyLong(),
any(HoodieSchema.class),
+ any(HoodieSchema.class), any(HoodieStorage.class));
+
+ List<HoodieLogFile> filteredLogs = Arrays.asList(
+ logFile("file1_1-0-1_003_1.log.parquet"),
+ logFile("file1_1-0-1_004_1.deletes.parquet"));
+ try (HoodieLsmFileGroupReader<IndexedRecord> reader = reader(
+ readerContext, Option.of(new HoodieBaseFile(baseFilePathInfo)),
filteredLogs, 0L);
+ ClosableIterator<IndexedRecord> iterator =
reader.getClosableIterator()) {
+ assertFalse(readerContext.getHasLogFiles());
+ List<IndexedRecord> records = drain(iterator);
+ assertEquals(1, records.size());
+ assertEquals("b", records.get(0).get(1).toString());
+ }
+ }
+
+ @Test
+ void testMetadataTableBaseFileIsNotFilteredByInstantRange() throws
IOException {
+ when(metaClient.getBasePath()).thenReturn(
+ new StoragePath("/tmp/data-table/" +
HoodieTableMetaClient.METADATA_TABLE_FOLDER_PATH));
+ HoodieReaderContext<IndexedRecord> readerContext =
spy(context(Option.of(exactMatchRange("002"))));
+ StoragePathInfo baseFilePathInfo =
pathInfo("/tmp/file1_1-0-1_001.parquet");
+ doReturn(ClosableIterator.wrap(Collections.singletonList(
+ recordWithCommitTime("001", "a", "retained", 1)).iterator()))
+ .when(readerContext).getFileRecordIterator(
+ eq(baseFilePathInfo), anyLong(), anyLong(),
any(HoodieSchema.class),
+ any(HoodieSchema.class), any(HoodieStorage.class));
+
+ try (HoodieLsmFileGroupReader<IndexedRecord> reader = reader(
+ readerContext, Option.of(new HoodieBaseFile(baseFilePathInfo)),
Collections.emptyList(), 0L);
+ ClosableIterator<IndexedRecord> iterator =
reader.getClosableIterator()) {
+ List<IndexedRecord> records = drain(iterator);
+ assertEquals(1, records.size());
+ assertEquals("a", records.get(0).get(1).toString());
+ }
+ }
+
private HoodieLsmFileGroupReader<IndexedRecord> reader(List<HoodieRecord>
records, boolean emitDeletes) {
return HoodieLsmFileGroupReader.<IndexedRecord>builder()
.withReaderContext(context())
@@ -180,7 +263,31 @@ class TestHoodieLsmFileGroupReader {
}
private HoodieReaderContext<IndexedRecord> context() {
- return new HoodieAvroReaderContext(storageConfiguration,
metaClient.getTableConfig(), Option.empty(), Option.empty());
+ return context(Option.empty());
+ }
+
+ private HoodieReaderContext<IndexedRecord> context(Option<InstantRange>
instantRange) {
+ return new HoodieAvroReaderContext(
+ storageConfiguration, metaClient.getTableConfig(), instantRange,
Option.empty());
+ }
+
+ private HoodieLsmFileGroupReader<IndexedRecord> reader(
+ HoodieReaderContext<IndexedRecord> readerContext,
+ Option<HoodieBaseFile> baseFileOption,
+ List<HoodieLogFile> logFiles,
+ long start) {
+ return HoodieLsmFileGroupReader.<IndexedRecord>builder()
+ .withReaderContext(readerContext)
+ .withHoodieTableMetaClient(metaClient)
+ .withLatestCommitTime("004")
+ .withDataSchema(SCHEMA_WITH_COMMIT_TIME)
+ .withRequestedSchema(SCHEMA_WITH_COMMIT_TIME)
+ .withProps(props)
+ .withBaseFileOption(baseFileOption)
+ .withLogFiles(logFiles.stream())
+ .withPartitionPath("")
+ .withStart(start)
+ .build();
}
private static HoodieRecord record(String key, String value, long ts) {
@@ -195,6 +302,31 @@ class TestHoodieLsmFileGroupReader {
return new HoodieEmptyRecord<>(new HoodieKey(key, ""),
HoodieOperation.DELETE, ts, HoodieRecord.HoodieRecordType.AVRO);
}
+ private static IndexedRecord recordWithCommitTime(
+ String commitTime, String key, String value, long ts) {
+ GenericData.Record record = new
GenericData.Record(SCHEMA_WITH_COMMIT_TIME.toAvroSchema());
+ record.put(HoodieRecord.COMMIT_TIME_METADATA_FIELD, commitTime);
+ record.put("id", key);
+ record.put("value", value);
+ record.put("ts", ts);
+ return record;
+ }
+
+ private static InstantRange exactMatchRange(String instant) {
+ return InstantRange.builder()
+ .rangeType(InstantRange.RangeType.EXACT_MATCH)
+ .explicitInstants(Collections.singleton(instant))
+ .build();
+ }
+
+ private static HoodieLogFile logFile(String fileName) {
+ return new HoodieLogFile(new StoragePath("/tmp/" + fileName), 10);
+ }
+
+ private static StoragePathInfo pathInfo(String path) {
+ return new StoragePathInfo(new StoragePath(path), 10, false, (short) 1,
1024, 0);
+ }
+
private static <T> List<T> drain(ClosableIterator<T> iterator) {
List<T> values = new ArrayList<>();
iterator.forEachRemaining(values::add);
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
index 8f39399f1953..4ac60852998c 100644
---
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/TestStreamWriteOperatorCoordinator.java
@@ -19,6 +19,7 @@
package org.apache.hudi.sink;
import org.apache.hudi.client.WriteStatus;
+import org.apache.hudi.client.common.HoodieFlinkEngineContext;
import org.apache.hudi.client.heartbeat.HoodieHeartbeatClient;
import org.apache.hudi.common.fs.FSUtils;
import org.apache.hudi.common.model.HoodieFailedWritesCleaningPolicy;
@@ -40,6 +41,7 @@ import org.apache.hudi.configuration.HadoopConfigurations;
import org.apache.hudi.exception.HoodieException;
import org.apache.hudi.hadoop.fs.HadoopFSUtils;
import org.apache.hudi.metadata.HoodieTableMetadata;
+import org.apache.hudi.metadata.MetadataPartitionType;
import org.apache.hudi.sink.event.Correspondent;
import org.apache.hudi.sink.event.WriteMetadataEvent;
import org.apache.hudi.sink.muttley.AthenaIngestionGateway;
@@ -48,6 +50,8 @@ import org.apache.hudi.sink.utils.EventBuffers;
import org.apache.hudi.sink.utils.MockCoordinatorExecutor;
import org.apache.hudi.sink.utils.NonThrownExecutor;
import org.apache.hudi.storage.HoodieStorage;
+import org.apache.hudi.storage.StoragePath;
+import org.apache.hudi.storage.StoragePathInfo;
import org.apache.hudi.storage.hadoop.HadoopStorageConfiguration;
import org.apache.hudi.util.StreamerUtil;
import org.apache.hudi.utils.TestConfigurations;
@@ -77,6 +81,7 @@ import java.lang.reflect.Field;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap;
+import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
@@ -431,6 +436,11 @@ public class TestStreamWriteOperatorCoordinator {
final String metadataTableBasePath =
HoodieTableMetadata.getMetadataTableBasePath(tempFile.getAbsolutePath());
HoodieTableMetaClient metadataTableMetaClient =
HoodieTestUtils.createMetaClient(new
HadoopStorageConfiguration(HadoopConfigurations.getHadoopConf(conf)),
metadataTableBasePath);
+
assertTrue(metadataTableMetaClient.getTableConfig().isLSMTreeStorageLayout());
+ assertTrue(FSUtils.getAllDataFilesInPartition(
+ metadataTableMetaClient.getStorage(),
+ new StoragePath(metadataTableBasePath,
MetadataPartitionType.FILES.getPartitionPath())).stream()
+ .anyMatch(pathInfo ->
FSUtils.isNativeLogFile(pathInfo.getPath().getName())));
HoodieTimeline completedTimeline =
metadataTableMetaClient.getActiveTimeline().filterCompletedInstants();
HoodieTableMetaClient dataTableMetaClient =
HoodieTestUtils.createMetaClient(new
HadoopStorageConfiguration(HadoopConfigurations.getHadoopConf(conf)), new
Path(metadataTableBasePath).getParent().getParent().toString());
@@ -449,6 +459,8 @@ public class TestStreamWriteOperatorCoordinator {
assertThat("One instant need to sync to metadata table",
completedTimeline.countInstants(), is(numCommits + 1));
assertThat(completedTimeline.lastInstant().get().requestedTime(),
is(instant));
}
+ assertLsmMetadataTableReads(conf, dataTableMetaClient);
+
// the 5th commit triggers the compaction
mockWriteWithMetadata(ckp++);
metadataTableMetaClient.reloadActiveTimeline();
@@ -758,6 +770,40 @@ public class TestStreamWriteOperatorCoordinator {
return instant;
}
+ private static void assertLsmMetadataTableReads(Configuration conf,
HoodieTableMetaClient dataTableMetaClient)
+ throws Exception {
+ String dataTableBasePath = dataTableMetaClient.getBasePath().toString();
+ String partitionPath = new Path(dataTableBasePath, "par1").toString();
+ List<String> partitionPaths = Arrays.asList(
+ partitionPath,
+ new Path(dataTableBasePath, "missing-partition").toString());
+
+ List<String> actualPartitions;
+ Map<String, List<StoragePathInfo>> actualFiles;
+ try (HoodieTableMetadata metadataTable =
dataTableMetaClient.getTableFormat().getMetadataFactory().create(
+ HoodieFlinkEngineContext.DEFAULT,
+ dataTableMetaClient.getStorage(),
+ StreamerUtil.metadataConfig(conf),
+ dataTableBasePath,
+ false)) {
+ actualPartitions = metadataTable.getAllPartitionPaths();
+ actualFiles = metadataTable.getAllFilesInPartitions(partitionPaths);
+ }
+
+ try (HoodieTableMetadata metadataTable =
dataTableMetaClient.getTableFormat().getMetadataFactory().create(
+ HoodieFlinkEngineContext.DEFAULT,
+ dataTableMetaClient.getStorage(),
+ StreamerUtil.metadataConfig(conf),
+ dataTableBasePath,
+ true)) {
+ assertEquals(actualPartitions, metadataTable.getAllPartitionPaths());
+ assertEquals(actualFiles,
metadataTable.getAllFilesInPartitions(partitionPaths));
+ }
+
+ assertEquals(Collections.singletonList("par1"), actualPartitions);
+ assertFalse(actualFiles.get(partitionPath).isEmpty());
+ }
+
private static WriteMetadataEvent createBootstrapEvent(int taskId, long
checkpointId, String instant, String partitionPath) {
return createOperatorEvent(taskId, checkpointId, instant, partitionPath,
true, false, 0.1);
}
diff --git
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/compact/ITTestHoodieFlinkCompactor.java
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/compact/ITTestHoodieFlinkCompactor.java
index 218435050250..188f868cd88e 100644
---
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/compact/ITTestHoodieFlinkCompactor.java
+++
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/compact/ITTestHoodieFlinkCompactor.java
@@ -271,7 +271,7 @@ public class ITTestHoodieFlinkCompactor {
// Mark instant as compaction inflight
metaClient.getActiveTimeline().transitionCompactionRequestedToInflight(instant);
- conf.set(FlinkOptions.WRITE_TABLE_VERSION, upgrade ?
HoodieTableVersion.EIGHT.versionCode() : HoodieTableVersion.SIX.versionCode());
+ conf.set(FlinkOptions.WRITE_TABLE_VERSION, upgrade ?
HoodieTableVersion.current().versionCode() :
HoodieTableVersion.SIX.versionCode());
env.addSource(new
CompactionPlanSourceFunction(Collections.singletonList(Pair.of(compactionInstantTime,
compactionPlan)), conf))
.name("compaction_source")
.uid("uid_compaction_source")