This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new d27d85ad05 [flink] Fix lookup scan mode after auto cache fallback
(#8885)
d27d85ad05 is described below
commit d27d85ad055cefff32210297c3ce545398513032
Author: zhoulii <[email protected]>
AuthorDate: Tue Jul 28 21:09:34 2026 +0800
[flink] Fix lookup scan mode after auto cache fallback (#8885)
---
.../flink/lookup/FileStoreLookupFunction.java | 10 ++++++++-
.../paimon/flink/lookup/LookupFileStoreTable.java | 4 +++-
.../flink/lookup/FileStoreLookupFunctionTest.java | 25 ++++++++++++++++++++++
3 files changed, 37 insertions(+), 2 deletions(-)
diff --git
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/FileStoreLookupFunction.java
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/FileStoreLookupFunction.java
index 2e815af0aa..1545a092c4 100644
---
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/FileStoreLookupFunction.java
+++
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/FileStoreLookupFunction.java
@@ -249,9 +249,17 @@ public class FileStoreLookupFunction implements
Serializable, Closeable {
}
if (lookupTable == null) {
+ FileStoreTable fullCacheTable = table;
+ // Resolve fallback AUTO to FULL for scan mode selection, but
preserve explicit MEMORY.
+ if (options.get(LOOKUP_CACHE_MODE) == LookupCacheMode.AUTO) {
+ fullCacheTable =
+ table.copy(
+ Collections.singletonMap(
+ LOOKUP_CACHE_MODE.key(),
LookupCacheMode.FULL.toString()));
+ }
FullCacheLookupTable.Context context =
new FullCacheLookupTable.Context(
- table,
+ fullCacheTable,
projection,
predicate,
createProjectedPredicate(projection),
diff --git
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/LookupFileStoreTable.java
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/LookupFileStoreTable.java
index bde9aa6c02..dcf36a397b 100644
---
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/LookupFileStoreTable.java
+++
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/LookupFileStoreTable.java
@@ -20,6 +20,7 @@ package org.apache.paimon.flink.lookup;
import org.apache.paimon.CoreOptions;
import org.apache.paimon.KeyValueFileStore;
+import org.apache.paimon.annotation.VisibleForTesting;
import org.apache.paimon.flink.FlinkConnectorOptions;
import org.apache.paimon.flink.utils.TableScanUtils;
import org.apache.paimon.options.Options;
@@ -127,7 +128,8 @@ public class LookupFileStoreTable extends
DelegatedFileStoreTable {
return this;
}
- private LookupStreamScanMode lookupStreamScanMode(FileStoreTable table,
List<String> joinKeys) {
+ @VisibleForTesting
+ LookupStreamScanMode lookupStreamScanMode(FileStoreTable table,
List<String> joinKeys) {
Options options = Options.fromMap(table.options());
if (options.get(LOOKUP_CACHE_MODE) ==
FlinkConnectorOptions.LookupCacheMode.AUTO
&& new HashSet<>(table.primaryKeys()).equals(new
HashSet<>(joinKeys))) {
diff --git
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/FileStoreLookupFunctionTest.java
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/FileStoreLookupFunctionTest.java
index eb89ca1081..2db4d2f3be 100644
---
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/FileStoreLookupFunctionTest.java
+++
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/FileStoreLookupFunctionTest.java
@@ -77,7 +77,10 @@ import java.util.UUID;
import java.util.concurrent.CopyOnWriteArrayList;
import static org.apache.paimon.data.BinaryRow.EMPTY_ROW;
+import static org.apache.paimon.flink.FlinkConnectorOptions.LOOKUP_CACHE_MODE;
import static
org.apache.paimon.flink.FlinkConnectorOptions.LOOKUP_REFRESH_TIME_PERIODS_BLACKLIST;
+import static
org.apache.paimon.flink.FlinkConnectorOptions.LookupCacheMode.FULL;
+import static
org.apache.paimon.flink.lookup.LookupFileStoreTable.LookupStreamScanMode.CHANGELOG;
import static org.apache.paimon.service.ServiceManager.PRIMARY_KEY_LOOKUP;
import static
org.apache.paimon.testutils.assertj.PaimonAssertions.anyCauseMatches;
import static org.assertj.core.api.Assertions.assertThat;
@@ -218,6 +221,28 @@ public class FileStoreLookupFunctionTest {
assertThat(queryExecutor).isInstanceOf(RemoteQueryExecutor.class);
}
+ @Test
+ public void testFallbackUpdatesCacheModeToFull() throws Exception {
+ table =
+ createFileStoreTable(false, false, false, null)
+
.copy(Collections.singletonMap(CoreOptions.SEQUENCE_FIELD.key(), "v"));
+ lookupFunction = createLookupFunction(table, true);
+ lookupFunction.open(tempDir.toString());
+
+
assertThat(lookupFunction.lookupTable()).isInstanceOf(FullCacheLookupTable.class);
+ FullCacheLookupTable fullCacheLookupTable =
+ (FullCacheLookupTable) lookupFunction.lookupTable();
+ assertThat(
+
Options.fromMap(fullCacheLookupTable.context.table.options())
+ .get(LOOKUP_CACHE_MODE))
+ .isEqualTo(FULL);
+ assertThat(
+
fullCacheLookupTable.context.table.lookupStreamScanMode(
+ fullCacheLookupTable.context.table.wrapped(),
+ fullCacheLookupTable.context.joinKey))
+ .isEqualTo(CHANGELOG);
+ }
+
@Test
public void testLookupScanLeak() throws Exception {
createLookupFunction(false);