From cbed23cbceef23c47f3c8cbbf529cc704f29765b Mon Sep 17 00:00:00 2001 From: zhangjunfan Date: Wed, 23 Sep 2026 16:23:55 +0800 Subject: [PATCH 1/2] [lake/paimon] Explictly snapshotId setting to improve scan-based lookup performance --- .../org/apache/fluss/config/TableConfig.java | 2 + .../lake/lakestorage/LakeTableLookuper.java | 23 +++++- .../lookup/PaimonScanBasedTableLookuper.java | 14 +++- .../lookup/PaimonLakeTableLookuperTest.java | 75 +++++++++++++++++-- .../HistoricalLakeLookupManager.java | 3 +- 5 files changed, 109 insertions(+), 8 deletions(-) diff --git a/fluss-common/src/main/java/org/apache/fluss/config/TableConfig.java b/fluss-common/src/main/java/org/apache/fluss/config/TableConfig.java index ed2c111c335..0c87dfc752e 100644 --- a/fluss-common/src/main/java/org/apache/fluss/config/TableConfig.java +++ b/fluss-common/src/main/java/org/apache/fluss/config/TableConfig.java @@ -17,6 +17,7 @@ package org.apache.fluss.config; +import org.apache.fluss.annotation.Internal; import org.apache.fluss.annotation.PublicEvolving; import org.apache.fluss.compression.ArrowCompressionInfo; import org.apache.fluss.metadata.ChangelogImage; @@ -136,6 +137,7 @@ public boolean isHistoricalPartitionEnabled() { } /** Gets the lookup mode for historical partitions of the table. */ + @Internal public LakeLookupMode getHistoricalLookupMode() { return config.get(ConfigOptions.TABLE_DATALAKE_HISTORICAL_PARTITION_LOOKUP_MODE); } diff --git a/fluss-common/src/main/java/org/apache/fluss/lake/lakestorage/LakeTableLookuper.java b/fluss-common/src/main/java/org/apache/fluss/lake/lakestorage/LakeTableLookuper.java index 6c93d9226e4..3b684d0feda 100644 --- a/fluss-common/src/main/java/org/apache/fluss/lake/lakestorage/LakeTableLookuper.java +++ b/fluss-common/src/main/java/org/apache/fluss/lake/lakestorage/LakeTableLookuper.java @@ -18,6 +18,7 @@ package org.apache.fluss.lake.lakestorage; import org.apache.fluss.annotation.PublicEvolving; +import org.apache.fluss.annotation.VisibleForTesting; import org.apache.fluss.metadata.ResolvedPartitionSpec; import org.apache.fluss.types.RowType; @@ -78,6 +79,18 @@ final class LookupContext { private final short schemaId; private final RowType valueRowType; private final LookupMetricRecorder lookupMetricRecorder; + private final @Nullable Long lakeSnapshotId; + + /** Creates a lookup context when the lake snapshot is unknown. */ + @VisibleForTesting + public LookupContext( + ResolvedPartitionSpec partitionSpec, + @Nullable Integer bucketId, + short schemaId, + RowType valueRowType, + LookupMetricRecorder lookupMetricRecorder) { + this(partitionSpec, bucketId, schemaId, valueRowType, lookupMetricRecorder, null); + } /** * Creates a lookup context. @@ -88,19 +101,22 @@ final class LookupContext { * @param schemaId schema id to encode the returned Fluss value with * @param valueRowType row type to encode the returned Fluss value with * @param lookupMetricRecorder recorder for lake table point lookup metrics + * @param lakeSnapshotId known lake snapshot ID, or null if unknown */ public LookupContext( ResolvedPartitionSpec partitionSpec, @Nullable Integer bucketId, short schemaId, RowType valueRowType, - LookupMetricRecorder lookupMetricRecorder) { + LookupMetricRecorder lookupMetricRecorder, + @Nullable Long lakeSnapshotId) { this.partitionSpec = checkNotNull(partitionSpec, "partitionSpec must not be null."); this.bucketId = bucketId; this.schemaId = schemaId; this.valueRowType = checkNotNull(valueRowType, "valueRowType must not be null."); this.lookupMetricRecorder = checkNotNull(lookupMetricRecorder, "lookupMetricRecorder must not be null."); + this.lakeSnapshotId = lakeSnapshotId; } /** Returns the resolved Fluss partition spec for the lookup. */ @@ -130,5 +146,10 @@ public RowType valueRowType() { public LookupMetricRecorder lookupMetricRecorder() { return lookupMetricRecorder; } + + /** Returns the known lake snapshot ID, or null if unknown. */ + public @Nullable Long lakeSnapshotId() { + return lakeSnapshotId; + } } } diff --git a/fluss-lake/fluss-lake-paimon/src/main/java/org/apache/fluss/lake/paimon/lookup/PaimonScanBasedTableLookuper.java b/fluss-lake/fluss-lake-paimon/src/main/java/org/apache/fluss/lake/paimon/lookup/PaimonScanBasedTableLookuper.java index 4b6aeee134c..184b7014a8b 100644 --- a/fluss-lake/fluss-lake-paimon/src/main/java/org/apache/fluss/lake/paimon/lookup/PaimonScanBasedTableLookuper.java +++ b/fluss-lake/fluss-lake-paimon/src/main/java/org/apache/fluss/lake/paimon/lookup/PaimonScanBasedTableLookuper.java @@ -28,6 +28,7 @@ import org.apache.fluss.utils.ExceptionUtils; import org.apache.fluss.utils.IOUtils; +import org.apache.paimon.CoreOptions; import org.apache.paimon.catalog.Catalog; import org.apache.paimon.catalog.CatalogContext; import org.apache.paimon.catalog.CatalogFactory; @@ -161,8 +162,19 @@ private FileStoreTable table() throws Exception { private @Nullable byte[] scanLookup(FileStoreTable table, byte[] key, LookupContext context) throws Exception { + FileStoreTable scanTable = table; + Long lakeSnapshotId = context.lakeSnapshotId(); + if (lakeSnapshotId != null) { + // Paimon propagates the table's snapshot and manifest caches to this copy. + scanTable = + scanTable.copy( + Collections.singletonMap( + CoreOptions.SCAN_SNAPSHOT_ID.key(), + String.valueOf(lakeSnapshotId))); + } ReadBuilder readBuilder = - table.newReadBuilder() + scanTable + .newReadBuilder() .withFilter(createKeyPredicates(table, key, context)) .withPartitionFilter(createPartitionPredicate(table, context)) .withReadType( diff --git a/fluss-lake/fluss-lake-paimon/src/test/java/org/apache/fluss/lake/paimon/lookup/PaimonLakeTableLookuperTest.java b/fluss-lake/fluss-lake-paimon/src/test/java/org/apache/fluss/lake/paimon/lookup/PaimonLakeTableLookuperTest.java index 9dfb4092489..8732aec4706 100644 --- a/fluss-lake/fluss-lake-paimon/src/test/java/org/apache/fluss/lake/paimon/lookup/PaimonLakeTableLookuperTest.java +++ b/fluss-lake/fluss-lake-paimon/src/test/java/org/apache/fluss/lake/paimon/lookup/PaimonLakeTableLookuperTest.java @@ -178,6 +178,67 @@ void testLookupPartitionedPrimaryKeyTable(LakeLookupMode lookupMode) throws Exce } } + @Test + void testScanLookupUsesRequestedSnapshot() throws Exception { + TablePath tablePath = TablePath.of(DB, "scan_snapshot"); + Schema schema = pkSchema(); + FileStoreTable table = createPaimonTable(tablePath, partitionedPkDescriptor(schema)); + long firstSnapshotId = + writeAndCommitData( + table, + Collections.singletonMap( + 0, Collections.singletonList(paimonRow(1, "20240101", "Alice")))); + ResolvedPartitionSpec partitionSpec = + ResolvedPartitionSpec.fromPartitionName( + Collections.singletonList("dt"), "20240101"); + LakeTableLookuper.LookupContext firstContext = + new LakeTableLookuper.LookupContext( + partitionSpec, + 0, + SCHEMA_ID, + schema.getRowType(), + NO_OP_LOOKUP_METRIC_RECORDER, + firstSnapshotId); + byte[] key = paimonKey(schema, 1, "20240101"); + + try (LakeTableLookuper lookuper = + createLookuper(LakeLookupMode.SCAN, tablePath, KvFormat.COMPACTED)) { + assertRow( + decodeValue(lookuper.lookup(key, firstContext), SCHEMA_ID, schema).row, + 1, + "20240101", + "Alice"); + + long secondSnapshotId = + writeAndCommitData( + table, + Collections.singletonMap( + 0, + Collections.singletonList( + paimonRow(1, "20240101", "Updated Alice")))); + // A committed newer snapshot does not change the snapshot pinned by this context. + assertRow( + decodeValue(lookuper.lookup(key, firstContext), SCHEMA_ID, schema).row, + 1, + "20240101", + "Alice"); + + LakeTableLookuper.LookupContext secondContext = + new LakeTableLookuper.LookupContext( + partitionSpec, + 0, + SCHEMA_ID, + schema.getRowType(), + NO_OP_LOOKUP_METRIC_RECORDER, + secondSnapshotId); + assertRow( + decodeValue(lookuper.lookup(key, secondContext), SCHEMA_ID, schema).row, + 1, + "20240101", + "Updated Alice"); + } + } + @ParameterizedTest(name = "lookupMode={0}") @EnumSource(LakeLookupMode.class) void testLookupKeysInComputedBuckets(LakeLookupMode lookupMode) throws Exception { @@ -836,9 +897,11 @@ void testLookupWithNonStringPartitionKey(LakeLookupMode lookupMode) throws Excep .distributedBy(2, "id") .build(); FileStoreTable table = createPaimonTable(tablePath, tableDescriptor); - writeAndCommitData( - table, - Collections.singletonMap(0, Collections.singletonList(paimonRow(1, 7, "Alice")))); + long snapshotId = + writeAndCommitData( + table, + Collections.singletonMap( + 0, Collections.singletonList(paimonRow(1, 7, "Alice")))); try (LakeTableLookuper lookuper = createLookuper(lookupMode, tablePath, KvFormat.COMPACTED)) { @@ -849,7 +912,8 @@ void testLookupWithNonStringPartitionKey(LakeLookupMode lookupMode) throws Excep 0, SCHEMA_ID, schema.getRowType(), - NO_OP_LOOKUP_METRIC_RECORDER); + NO_OP_LOOKUP_METRIC_RECORDER, + snapshotId); BinaryValue decodedValue = decodeValue( @@ -882,7 +946,8 @@ void testRejectAppendOnlyTableAndLookupAfterClose(LakeLookupMode lookupMode) thr 0, SCHEMA_ID, schema.getRowType(), - NO_OP_LOOKUP_METRIC_RECORDER); + NO_OP_LOOKUP_METRIC_RECORDER, + null); assertThatThrownBy(() -> lookuper.lookup(new byte[0], context)) .isInstanceOf(UnsupportedOperationException.class) diff --git a/fluss-server/src/main/java/org/apache/fluss/server/replica/historical/HistoricalLakeLookupManager.java b/fluss-server/src/main/java/org/apache/fluss/server/replica/historical/HistoricalLakeLookupManager.java index 8210d9ce8e5..9c3abe5d7f9 100644 --- a/fluss-server/src/main/java/org/apache/fluss/server/replica/historical/HistoricalLakeLookupManager.java +++ b/fluss-server/src/main/java/org/apache/fluss/server/replica/historical/HistoricalLakeLookupManager.java @@ -373,7 +373,8 @@ private LookupContext createLookupContext( lakeBucketId, (short) schemaInfo.getSchemaId(), schemaInfo.getSchema().getRowType(), - lookupMetricRecorder); + lookupMetricRecorder, + requiredLakeSnapshotIds.get(tableInfo.getTableId())); return new LookupContext( tableInfo.getTableId(), schemaInfo.getSchemaId(), tablePath, lookupContext); } From a50e5d06b7f5ab573bffe5765863caac8755ae30 Mon Sep 17 00:00:00 2001 From: zhangjunfan Date: Thu, 24 Sep 2026 17:46:06 +0800 Subject: [PATCH 2/2] comment --- .../lake/paimon/lookup/PaimonScanBasedTableLookuper.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/fluss-lake/fluss-lake-paimon/src/main/java/org/apache/fluss/lake/paimon/lookup/PaimonScanBasedTableLookuper.java b/fluss-lake/fluss-lake-paimon/src/main/java/org/apache/fluss/lake/paimon/lookup/PaimonScanBasedTableLookuper.java index 184b7014a8b..945181d50f9 100644 --- a/fluss-lake/fluss-lake-paimon/src/main/java/org/apache/fluss/lake/paimon/lookup/PaimonScanBasedTableLookuper.java +++ b/fluss-lake/fluss-lake-paimon/src/main/java/org/apache/fluss/lake/paimon/lookup/PaimonScanBasedTableLookuper.java @@ -62,7 +62,8 @@ import static org.apache.fluss.utils.concurrent.LockUtils.inWriteLock; /** - * Looks up a primary key by scanning the latest Paimon snapshot with a limit of one row. + * Looks up a primary key by scanning the requested Paimon snapshot, or the latest snapshot when no + * snapshot ID is provided, with a limit of one row. * *

Each scan is restricted to the requested partition, bucket, and complete primary key. It does * not create local lookup files. Lookups use independent readers and encoders and may run in @@ -120,7 +121,7 @@ public PaimonScanBasedTableLookuper( @Override public void requestRefresh() { - // Each lookup already plans a fresh scan of the latest snapshot. + // Each lookup plans a fresh scan, so there are no cached data files to refresh. } @Override