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 8c53b2cd82 [flink] Speed up lookup join tests (#8997)
8c53b2cd82 is described below

commit 8c53b2cd820880959ad9e09359e38040b63f9ee2
Author: Jingsong Lee <[email protected]>
AuthorDate: Mon Aug 3 17:44:28 2026 +0800

    [flink] Speed up lookup join tests (#8997)
---
 .../org/apache/paimon/flink/LookupJoinITCase.java  | 1072 ++------------------
 .../flink/lookup/BlobAsDescriptorRowTest.java      |   86 ++
 .../lookup/DynamicPartitionNumberLoaderTest.java   |  150 +++
 3 files changed, 326 insertions(+), 982 deletions(-)

diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/LookupJoinITCase.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/LookupJoinITCase.java
index f5b6ed27d8..385d1d1cf6 100644
--- 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/LookupJoinITCase.java
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/LookupJoinITCase.java
@@ -26,7 +26,6 @@ import 
org.apache.flink.table.planner.factories.TestValuesTableFactory;
 import org.apache.flink.types.Row;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
-import org.junit.jupiter.params.provider.EnumSource;
 import org.junit.jupiter.params.provider.ValueSource;
 
 import java.util.Arrays;
@@ -34,13 +33,18 @@ import java.util.Collections;
 import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
-import java.util.concurrent.ThreadLocalRandom;
 import java.util.concurrent.TimeUnit;
 
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
 
-/** ITCase for lookup join. */
+/**
+ * ITCase for lookup join planner and runtime wiring.
+ *
+ * <p>Lookup cache semantics such as projection, filtering, refresh, sequence 
handling, and backend
+ * compatibility belong in the direct tests under {@code 
org.apache.paimon.flink.lookup}. Keep this
+ * suite focused on behavior which requires a SQL planner or a running Flink 
job.
+ */
 public class LookupJoinITCase extends CatalogITCaseBase {
 
     @Override
@@ -83,10 +87,9 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         }
     }
 
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testLookupEmptyTable(LookupCacheMode cacheMode) throws 
Exception {
-        initTable(cacheMode);
+    @Test
+    public void testLookupEmptyTable() throws Exception {
+        initTable(LookupCacheMode.AUTO);
         String query =
                 "SELECT T.i, D.j, D.k1, D.k2 FROM T LEFT JOIN DIM for 
system_time as of T.proctime AS D ON T.i = D.i";
         BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
@@ -112,38 +115,6 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         iterator.close();
     }
 
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testLookup(LookupCacheMode cacheMode) throws Exception {
-        initTable(cacheMode);
-        sql("INSERT INTO DIM VALUES (1, 11, 111, 1111), (2, 22, 222, 2222)");
-
-        String query =
-                "SELECT T.i, D.j, D.k1, D.k2 FROM T LEFT JOIN DIM for 
system_time as of T.proctime AS D ON T.i = D.i";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (1), (2), (3)");
-        List<Row> result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, 11, 111, 1111),
-                        Row.of(2, 22, 222, 2222),
-                        Row.of(3, null, null, null));
-
-        sql("INSERT INTO DIM VALUES (2, 44, 444, 4444), (3, 33, 333, 3333)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (1), (2), (3), (4)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, 11, 111, 1111),
-                        Row.of(2, 44, 444, 4444),
-                        Row.of(3, 33, 333, 3333),
-                        Row.of(4, null, null, null));
-
-        iterator.close();
-    }
-
     @Test
     public void testLookupIgnoreScanOptions() throws Exception {
         sql(
@@ -182,275 +153,6 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         streamIter.close();
     }
 
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testLookupProjection(LookupCacheMode cacheMode) throws 
Exception {
-        initTable(cacheMode);
-        sql("INSERT INTO DIM VALUES (1, 11, 111, 1111), (2, 22, 222, 2222)");
-
-        String query =
-                "SELECT T.i, D.j, D.k1 FROM T LEFT JOIN DIM for system_time as 
of T.proctime AS D ON T.i = D.i";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (1), (2), (3)");
-        List<Row> result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, 11, 111), Row.of(2, 22, 222), Row.of(3, 
null, null));
-
-        sql("INSERT INTO DIM VALUES (2, 44, 444, 4444), (3, 33, 333, 3333)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (1), (2), (3), (4)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, 11, 111),
-                        Row.of(2, 44, 444),
-                        Row.of(3, 33, 333),
-                        Row.of(4, null, null));
-
-        iterator.close();
-    }
-
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testLookupFilterPk(LookupCacheMode cacheMode) throws Exception 
{
-        initTable(cacheMode);
-        sql("INSERT INTO DIM VALUES (1, 11, 111, 1111), (2, 22, 222, 2222)");
-
-        String query =
-                "SELECT T.i, D.j, D.k1 FROM T LEFT JOIN DIM for system_time as 
of T.proctime AS D ON T.i = D.i AND D.i > 2";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (1), (2), (3)");
-        List<Row> result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, null, null), Row.of(2, null, null), 
Row.of(3, null, null));
-
-        sql("INSERT INTO DIM VALUES (2, 44, 444, 4444), (3, 33, 333, 3333)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (1), (2), (3), (4)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, null, null),
-                        Row.of(2, null, null),
-                        Row.of(3, 33, 333),
-                        Row.of(4, null, null));
-
-        iterator.close();
-    }
-
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testLookupFilterSelect(LookupCacheMode cacheMode) throws 
Exception {
-        initTable(cacheMode);
-        sql("INSERT INTO DIM VALUES (1, 11, 111, 1111), (2, 22, 222, 2222)");
-
-        String query =
-                "SELECT T.i, D.j, D.k1 FROM T LEFT JOIN DIM for system_time as 
of T.proctime AS D ON T.i = D.i AND D.k1 > 111";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (1), (2), (3)");
-        List<Row> result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, null, null), Row.of(2, 22, 222), Row.of(3, 
null, null));
-
-        sql("INSERT INTO DIM VALUES (2, 44, 444, 4444), (3, 33, 333, 3333)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (1), (2), (3), (4)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, null, null),
-                        Row.of(2, 44, 444),
-                        Row.of(3, 33, 333),
-                        Row.of(4, null, null));
-
-        iterator.close();
-    }
-
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testLookupFilterUnSelect(LookupCacheMode cacheMode) throws 
Exception {
-        initTable(cacheMode);
-        sql("INSERT INTO DIM VALUES (1, 11, 111, 1111), (2, 22, 222, 2222)");
-
-        String query =
-                "SELECT T.i, D.j, D.k1 FROM T LEFT JOIN DIM for system_time as 
of T.proctime AS D ON T.i = D.i AND D.k2 > 1111";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (1), (2), (3)");
-        List<Row> result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, null, null), Row.of(2, 22, 222), Row.of(3, 
null, null));
-
-        sql("INSERT INTO DIM VALUES (2, 44, 444, 4444), (3, 33, 333, 3333)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (1), (2), (3), (4)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, null, null),
-                        Row.of(2, 44, 444),
-                        Row.of(3, 33, 333),
-                        Row.of(4, null, null));
-
-        iterator.close();
-    }
-
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testLookupFilterUnSelectAndUpdate(LookupCacheMode cacheMode) 
throws Exception {
-        initTable(cacheMode);
-        sql("INSERT INTO DIM VALUES (1, 11, 111, 1111), (2, 22, 222, 2222)");
-
-        String query =
-                "SELECT T.i, D.j, D.k1 FROM T LEFT JOIN DIM for system_time as 
of T.proctime AS D ON T.i = D.i AND D.k2 < 4444";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (1), (2), (3)");
-        List<Row> result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, 11, 111), Row.of(2, 22, 222), Row.of(3, 
null, null));
-
-        sql("INSERT INTO DIM VALUES (2, 44, 444, 4444), (3, 33, 333, 3333)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (1), (2), (3), (4)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, 11, 111),
-                        Row.of(2, null, null),
-                        Row.of(3, 33, 333),
-                        Row.of(4, null, null));
-
-        iterator.close();
-    }
-
-    @Test
-    public void testLookupUpdateAfterLeafPredicate0() throws Exception {
-        sql(
-                "CREATE TABLE fact (\n"
-                        + "  name string,\n"
-                        + "  k string,\n"
-                        + "  proctime as PROCTIME()\n"
-                        + ")\n"
-                        + "WITH (\n"
-                        + "    'bucket' = '1',\n"
-                        + "    'bucket-key'='name'\n"
-                        + ");");
-        sql(
-                "CREATE TABLE dim (\n"
-                        + "  id bigint,\n"
-                        + "  k string,\n"
-                        + "  v string,\n"
-                        + "  PRIMARY KEY (id) NOT ENFORCED \n"
-                        + ")\n"
-                        + "WITH (\n"
-                        + "    'bucket' = '1'\n"
-                        + ");");
-        String query =
-                "select \n"
-                        + "a.name,\n"
-                        + "a.k as ak,\n"
-                        + "b.k as bk,\n"
-                        + "b.v\n"
-                        + "from fact  /*+ 
OPTIONS('scan.mode'='latest','continuous.discovery-interval'='1s') */ a\n"
-                        + "left join dim /*+ 
OPTIONS('continuous.discovery-interval'='3s') */ FOR SYSTEM_TIME AS OF 
a.proctime AS b \n"
-                        + "on a.k = b.k and b.v<'y'";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO dim VALUES (1,'k','x')");
-        sql("INSERT INTO fact VALUES ('r1','k')");
-        iterator.collect(1);
-        sql("INSERT INTO dim VALUES (1,'k','y')");
-        sql("INSERT INTO fact VALUES ('r2','k')");
-        sql("INSERT INTO fact VALUES ('r3','k')");
-        List<Row> result = iterator.collect(2);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of("r2", "k", "k", "x"), Row.of("r3", "k", "k", 
"x"));
-
-        iterator.close();
-    }
-
-    @Test
-    public void testLookupUpdateAfterLeafPredicate1() throws Exception {
-        sql(
-                "CREATE TABLE fact (\n"
-                        + "  name string,\n"
-                        + "  k string,\n"
-                        + "  proctime as PROCTIME()\n"
-                        + ")\n"
-                        + "WITH (\n"
-                        + "    'bucket' = '1',\n"
-                        + "    'bucket-key'='name'\n"
-                        + ");");
-        sql(
-                "CREATE TABLE dim (\n"
-                        + "  id bigint,\n"
-                        + "  k string,\n"
-                        + "  v string,\n"
-                        + "  PRIMARY KEY (id) NOT ENFORCED \n"
-                        + ")\n"
-                        + "WITH (\n"
-                        + "    'bucket' = '1'\n"
-                        + ");");
-        String query =
-                "select \n"
-                        + "a.name,\n"
-                        + "a.k as ak,\n"
-                        + "b.k as bk,\n"
-                        + "b.v\n"
-                        + "from fact  /*+ 
OPTIONS('scan.mode'='latest','continuous.discovery-interval'='1s') */ a\n"
-                        + "left join dim /*+ 
OPTIONS('continuous.discovery-interval'='3s') */ FOR SYSTEM_TIME AS OF 
a.proctime AS b \n"
-                        + "on a.k = b.k and b.v<'y'";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO dim VALUES (1,'k','x')");
-        sql("INSERT INTO fact VALUES ('r1','k')");
-        Thread.sleep(5000);
-        sql("INSERT INTO dim VALUES (1,'k','y')");
-        sql("INSERT INTO fact VALUES ('r2','k')");
-        sql("INSERT INTO fact VALUES ('r3','k')");
-        List<Row> result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of("r1", "k", "k", "x"),
-                        Row.of("r2", "k", null, null),
-                        Row.of("r3", "k", null, null));
-
-        iterator.close();
-    }
-
-    @Test
-    public void testLookupUpdateAfterLeafPredicate2() throws Exception {
-        sql("CREATE TABLE fact (name STRING, i INT, `proctime` AS 
PROCTIME())");
-        sql(
-                "CREATE TABLE dim (i INT PRIMARY KEY NOT ENFORCED, j INT, k1 
INT, k2 INT) WITH"
-                        + " ('continuous.discovery-interval'='1 ms')");
-
-        String query =
-                "SELECT fact.name, fact.i, D.k1 FROM fact LEFT JOIN dim for 
system_time as of fact.proctime AS D ON fact.i = D.j AND D.k1 > 100";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO dim VALUES (1, 11, 111, 1111)");
-        sql("INSERT INTO fact VALUES ('a',11)");
-        List<Row> result = iterator.collect(1);
-        assertThat(result).containsExactlyInAnyOrder(Row.of("a", 11, 111));
-
-        sql("INSERT INTO dim VALUES (1,11,100,1111)");
-        sql("INSERT INTO fact VALUES ('b',11)");
-        result = iterator.collect(1);
-        assertThat(result).containsExactlyInAnyOrder(Row.of("b", 11, null));
-        iterator.close();
-    }
-
     @Test
     public void testLookupUpdateAfterCompoundPredicate() throws Exception {
         sql("CREATE TABLE fact (name STRING, i INT, `proctime` AS 
PROCTIME())");
@@ -507,166 +209,6 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         iterator.close();
     }
 
-    @Test
-    public void testNonPkLookupProjection() throws Exception {
-        initTable(LookupCacheMode.FULL);
-        sql("INSERT INTO DIM VALUES (1, 11, 111, 1111), (2, 22, 222, 2222), 
(3, 22, 333, 3333)");
-
-        String query =
-                "SELECT T.i, D.k1 FROM T LEFT JOIN DIM for system_time as of 
T.proctime AS D ON T.i = D.j";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (11), (22), (33)");
-        List<Row> result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(11, 111), Row.of(22, 222), Row.of(22, 333), 
Row.of(33, null));
-
-        sql("INSERT INTO DIM VALUES (2, 44, 444, 4444), (3, 33, 333, 3333)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (11), (22), (33), (44)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(11, 111), Row.of(22, null), Row.of(33, 333), 
Row.of(44, 444));
-
-        iterator.close();
-    }
-
-    @Test
-    public void testNonPkLookupFilterPk() throws Exception {
-        initTable(LookupCacheMode.FULL);
-        sql("INSERT INTO DIM VALUES (1, 11, 111, 1111), (2, 22, 222, 2222), 
(3, 22, 333, 3333)");
-
-        String query =
-                "SELECT T.i, D.k1 FROM T LEFT JOIN DIM for system_time as of 
T.proctime AS D ON T.i = D.j AND D.i > 2";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (11), (22), (33)");
-        List<Row> result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(Row.of(11, null), Row.of(22, 333), 
Row.of(33, null));
-
-        sql("INSERT INTO DIM VALUES (2, 44, 444, 4444), (3, 33, 333, 3333)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (11), (22), (33), (44)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(11, null), Row.of(22, null), Row.of(33, 333), 
Row.of(44, null));
-
-        iterator.close();
-    }
-
-    @Test
-    public void testNonPkLookupFilterSelect() throws Exception {
-        initTable(LookupCacheMode.FULL);
-        sql("INSERT INTO DIM VALUES (1, 11, 111, 1111), (2, 22, 222, 2222), 
(3, 22, 333, 3333)");
-
-        String query =
-                "SELECT T.i, D.k1 FROM T LEFT JOIN DIM for system_time as of 
T.proctime AS D ON T.i = D.j AND D.k1 > 111";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (11), (22), (33)");
-        List<Row> result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(11, null), Row.of(22, 222), Row.of(22, 333), 
Row.of(33, null));
-
-        sql("INSERT INTO DIM VALUES (2, 44, 444, 4444), (3, 33, 333, 3333)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (11), (22), (33), (44)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(11, null), Row.of(22, null), Row.of(33, 333), 
Row.of(44, 444));
-
-        iterator.close();
-    }
-
-    @Test
-    public void testNonPkLookupFilterUnSelect() throws Exception {
-        initTable(LookupCacheMode.FULL);
-        sql("INSERT INTO DIM VALUES (1, 11, 111, 1111), (2, 22, 222, 2222), 
(3, 22, 333, 3333)");
-
-        String query =
-                "SELECT T.i, D.k1 FROM T LEFT JOIN DIM for system_time as of 
T.proctime AS D ON T.i = D.j AND D.k2 > 1111";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (11), (22), (33)");
-        List<Row> result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(11, null), Row.of(22, 222), Row.of(22, 333), 
Row.of(33, null));
-
-        sql("INSERT INTO DIM VALUES (2, 44, 444, 4444), (3, 33, 333, 3333)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (11), (22), (33), (44)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(11, null), Row.of(22, null), Row.of(33, 333), 
Row.of(44, 444));
-
-        iterator.close();
-    }
-
-    @Test
-    public void testNonPkLookupFilterUnSelectAndUpdate() throws Exception {
-        initTable(LookupCacheMode.FULL);
-        sql("INSERT INTO DIM VALUES (1, 11, 111, 1111), (2, 22, 222, 2222), 
(3, 22, 333, 3333)");
-
-        String query =
-                "SELECT T.i, D.k1 FROM T LEFT JOIN DIM for system_time as of 
T.proctime AS D ON T.i = D.j AND D.k2 < 4444";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (11), (22), (33)");
-        List<Row> result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(11, 111), Row.of(22, 222), Row.of(22, 333), 
Row.of(33, null));
-
-        sql("INSERT INTO DIM VALUES (2, 44, 444, 4444), (3, 33, 333, 3333)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (11), (22), (33), (44)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(11, 111), Row.of(22, null), Row.of(33, 333), 
Row.of(44, null));
-
-        iterator.close();
-    }
-
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testRepeatRefresh(LookupCacheMode cacheMode) throws Exception {
-        initTable(cacheMode);
-        sql("INSERT INTO DIM VALUES (1, 11, 111, 1111), (2, 22, 222, 2222)");
-
-        String query =
-                "SELECT T.i, D.j, D.k1 FROM T LEFT JOIN DIM for system_time as 
of T.proctime AS D ON T.i = D.i";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (1), (2), (3)");
-        List<Row> result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, 11, 111), Row.of(2, 22, 222), Row.of(3, 
null, null));
-
-        sql("INSERT INTO DIM VALUES (2, 44, 444, 4444)");
-        sql("INSERT INTO DIM VALUES (3, 33, 333, 3333)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (1), (2), (3), (4)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, 11, 111),
-                        Row.of(2, 44, 444),
-                        Row.of(3, 33, 333),
-                        Row.of(4, null, null));
-
-        iterator.close();
-    }
-
     @Test
     public void testLookupPartialUpdateIllegal() {
         sql(
@@ -684,7 +226,6 @@ public class LookupJoinITCase extends CatalogITCaseBase {
 
     @Test
     public void testLookupPartialUpdate() throws Exception {
-        testLookupPartialUpdate("none");
         testLookupPartialUpdate("zstd");
     }
 
@@ -714,10 +255,9 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         sql("TRUNCATE TABLE T");
     }
 
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testRetryLookup(LookupCacheMode cacheMode) throws Exception {
-        initTable(cacheMode);
+    @Test
+    public void testRetryLookup() throws Exception {
+        initTable(LookupCacheMode.FULL);
         sql("INSERT INTO DIM VALUES (1, 11, 111, 1111), (2, 22, 222, 2222)");
 
         String query =
@@ -738,15 +278,14 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         iterator.close();
     }
 
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testAsyncRetryLookup(LookupCacheMode cacheMode) throws 
Exception {
-        initTable(cacheMode);
+    @Test
+    public void testAsyncRetryLookup() throws Exception {
+        initTable(LookupCacheMode.FULL);
         sql("INSERT INTO DIM VALUES (1, 11, 111, 1111), (2, 22, 222, 2222)");
 
         String query =
                 "SELECT /*+ LOOKUP('table'='D', 
'retry-predicate'='lookup_miss',"
-                        + " 'retry-strategy'='fixed_delay', 
'output-mode'='allow_unordered', 'fixed-delay'='3s','max-attempts'='30') */"
+                        + " 'retry-strategy'='fixed_delay', 
'output-mode'='allow_unordered', 'fixed-delay'='500ms','max-attempts'='180') */"
                         + " T.i, D.j, D.k1, D.k2 FROM T LEFT JOIN DIM /*+ 
OPTIONS('lookup.async'='true') */ for system_time as of T.proctime AS D ON T.i 
= D.i";
         BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
 
@@ -754,173 +293,40 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         sql("INSERT INTO T VALUES (2)");
         sql("INSERT INTO T VALUES (1)");
         assertThat(iterator.collect(2))
-                .containsExactlyInAnyOrder(Row.of(1, 11, 111, 1111), Row.of(2, 
22, 222, 2222));
-
-        sql("INSERT INTO DIM VALUES (3, 33, 333, 3333)");
-        assertThat(iterator.collect(1, 10, TimeUnit.MINUTES))
-                .containsExactlyInAnyOrder(Row.of(3, 33, 333, 3333));
-
-        iterator.close();
-    }
-
-    @Test
-    public void testLookupPartitionedTable() throws Exception {
-        initTable(LookupCacheMode.AUTO);
-        String query =
-                "SELECT T.i, D.j, D.k1, D.k2 FROM T LEFT JOIN PARTITIONED_DIM 
for system_time as of T.proctime AS D ON T.i = D.i";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (1), (2), (3)");
-
-        List<Row> result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, null, null, null),
-                        Row.of(2, null, null, null),
-                        Row.of(3, null, null, null));
-
-        sql("INSERT INTO PARTITIONED_DIM VALUES (1, 11, 111, 1111), (2, 22, 
222, 2222)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (1), (2), (4)");
-        result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, 11, 111, 1111),
-                        Row.of(2, 22, 222, 2222),
-                        Row.of(4, null, null, null));
-        iterator.close();
-    }
-
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testLookupMaxPtPartitionedTable(LookupCacheMode mode) throws 
Exception {
-        boolean testDynamicBucket = ThreadLocalRandom.current().nextBoolean();
-        String primaryKeys;
-        String bucket;
-        if (testDynamicBucket) {
-            primaryKeys = "k";
-            bucket = "-1";
-        } else {
-            primaryKeys = "pt, k";
-            bucket = "1";
-        }
-        sql(
-                "CREATE TABLE PARTITIONED_DIM (pt STRING, k INT, v INT, 
PRIMARY KEY (%s) NOT ENFORCED)"
-                        + "PARTITIONED BY (`pt`) WITH ("
-                        + "'bucket' = '%s', "
-                        + "'lookup.dynamic-partition' = 'max_pt()', "
-                        + "'lookup.dynamic-partition.refresh-interval' = '1 
ms', "
-                        + "'lookup.cache' = '%s', "
-                        + "'continuous.discovery-interval'='1 ms')",
-                primaryKeys, bucket, mode);
-        String query =
-                "SELECT T.i, D.v FROM T LEFT JOIN PARTITIONED_DIM for 
system_time as of T.proctime AS D ON T.i = D.k";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO PARTITIONED_DIM VALUES ('1', 1, 2)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (1)");
-        List<Row> result = iterator.collect(1);
-        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 2));
-
-        sql("INSERT INTO PARTITIONED_DIM VALUES ('2', 1, 3)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (1)");
-        result = iterator.collect(1);
-        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 3));
-
-        iterator.close();
-    }
-
-    @Test
-    public void testLookupNonPkAppendTable() throws Exception {
-        sql(
-                "CREATE TABLE DIM_NO_PK (i INT, j INT, k1 INT, k2 INT) "
-                        + "PARTITIONED BY (`i`) WITH 
('continuous.discovery-interval'='1 ms')");
-
-        String query =
-                "SELECT T.i, D.j, D.k1, D.k2 FROM T LEFT JOIN DIM_NO_PK for 
system_time as of T.proctime AS D ON T.i "
-                        + "= D.i";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (1), (2), (3)");
-
-        List<Row> result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, null, null, null),
-                        Row.of(2, null, null, null),
-                        Row.of(3, null, null, null));
-
-        sql(
-                "INSERT INTO DIM_NO_PK VALUES (1, 11, 111, 1111), (1, 12, 112, 
1112), (1, 11, 111, 1111)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (1), (2), (4)");
-        result = iterator.collect(5);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, 11, 111, 1111),
-                        Row.of(1, 11, 111, 1111),
-                        Row.of(1, 12, 112, 1112),
-                        Row.of(2, null, null, null),
-                        Row.of(4, null, null, null));
-        iterator.close();
-    }
-
-    @Test
-    public void testWithSequenceFieldTable() throws Exception {
-        sql(
-                "CREATE TABLE DIM_WITH_SEQUENCE (i INT PRIMARY KEY NOT 
ENFORCED, j INT, k1 INT, k2 INT) WITH"
-                        + " ('continuous.discovery-interval'='1 ms', 
'sequence.field' = 'j')");
-        sql("INSERT INTO DIM_WITH_SEQUENCE VALUES (1, 11, 111, 1111), (2, 22, 
222, 2222)");
-
-        String query =
-                "SELECT T.i, D.j, D.k1, D.k2 FROM T LEFT JOIN 
DIM_WITH_SEQUENCE for system_time as of T.proctime AS D ON T.i = D.i";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (1), (2), (3)");
-        List<Row> result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, 11, 111, 1111),
-                        Row.of(2, 22, 222, 2222),
-                        Row.of(3, null, null, null));
+                .containsExactlyInAnyOrder(Row.of(1, 11, 111, 1111), Row.of(2, 
22, 222, 2222));
 
-        sql("INSERT INTO DIM_WITH_SEQUENCE VALUES (2, 11, 444, 4444), (3, 33, 
333, 3333)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (1), (2), (3), (4)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, 11, 111, 1111),
-                        Row.of(2, 22, 222, 2222), // not change
-                        Row.of(3, 33, 333, 3333),
-                        Row.of(4, null, null, null));
+        sql("INSERT INTO DIM VALUES (3, 33, 333, 3333)");
+        assertThat(iterator.collect(1, 10, TimeUnit.MINUTES))
+                .containsExactlyInAnyOrder(Row.of(3, 33, 333, 3333));
 
         iterator.close();
     }
 
     @Test
-    public void testAsyncRetryLookupWithSequenceField() throws Exception {
+    public void testLookupMaxPtDynamicBucketTable() throws Exception {
         sql(
-                "CREATE TABLE DIM_WITH_SEQUENCE (i INT PRIMARY KEY NOT 
ENFORCED, j INT, k1 INT, k2 INT) WITH"
-                        + " ('continuous.discovery-interval'='1 ms', 
'sequence.field' = 'j')");
-        sql("INSERT INTO DIM_WITH_SEQUENCE VALUES (1, 11, 111, 1111), (2, 22, 
222, 2222)");
-
+                "CREATE TABLE PARTITIONED_DIM (pt STRING, k INT, v INT, 
PRIMARY KEY (k) NOT ENFORCED)"
+                        + "PARTITIONED BY (`pt`) WITH ("
+                        + "'bucket' = '-1', "
+                        + "'lookup.dynamic-partition' = 'max_pt()', "
+                        + "'lookup.dynamic-partition.refresh-interval' = '1 
ms', "
+                        + "'lookup.cache' = 'full', "
+                        + "'continuous.discovery-interval'='1 ms')");
         String query =
-                "SELECT /*+ LOOKUP('table'='D', 
'retry-predicate'='lookup_miss',"
-                        + " 'retry-strategy'='fixed_delay', 
'output-mode'='allow_unordered', 'fixed-delay'='3s','max-attempts'='60') */"
-                        + " T.i, D.j, D.k1, D.k2 FROM T LEFT JOIN 
DIM_WITH_SEQUENCE /*+ OPTIONS('lookup.async'='true') */ for system_time as of 
T.proctime AS D ON T.i = D.i";
+                "SELECT T.i, D.v FROM T LEFT JOIN PARTITIONED_DIM for 
system_time as of T.proctime AS D ON T.i = D.k";
         BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
 
-        sql("INSERT INTO T VALUES (3)");
-        sql("INSERT INTO T VALUES (2)");
+        sql("INSERT INTO PARTITIONED_DIM VALUES ('1', 1, 2)");
+        Thread.sleep(2000); // wait refresh
         sql("INSERT INTO T VALUES (1)");
-        assertThat(iterator.collect(2))
-                .containsExactlyInAnyOrder(Row.of(1, 11, 111, 1111), Row.of(2, 
22, 222, 2222));
+        List<Row> result = iterator.collect(1);
+        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 2));
 
-        sql("INSERT INTO DIM_WITH_SEQUENCE VALUES (3, 33, 333, 3333)");
-        assertThat(iterator.collect(1)).containsExactlyInAnyOrder(Row.of(3, 
33, 333, 3333));
+        sql("INSERT INTO PARTITIONED_DIM VALUES ('2', 1, 3)");
+        Thread.sleep(2000); // wait refresh
+        sql("INSERT INTO T VALUES (1)");
+        result = iterator.collect(1);
+        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 3));
 
         iterator.close();
     }
@@ -934,7 +340,7 @@ public class LookupJoinITCase extends CatalogITCaseBase {
 
         String query =
                 "SELECT /*+ LOOKUP('table'='D', 
'retry-predicate'='lookup_miss',"
-                        + " 'retry-strategy'='fixed_delay', 
'output-mode'='allow_unordered', 'fixed-delay'='3s','max-attempts'='60') */"
+                        + " 'retry-strategy'='fixed_delay', 
'output-mode'='allow_unordered', 'fixed-delay'='500ms','max-attempts'='360') */"
                         + " T.i, D.i, D.j, D.k2 FROM T LEFT JOIN 
DIM_WITH_SEQUENCE /*+ OPTIONS('lookup.async'='true') */ for system_time as of 
T.proctime AS D ON T.i = D.k1";
         BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
 
@@ -952,13 +358,11 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         iterator.close();
     }
 
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testPartialCacheBucketKeyOrder(LookupCacheMode mode) throws 
Exception {
+    @Test
+    public void testPartialCacheBucketKeyOrder() throws Exception {
         sql(
                 "CREATE TABLE DIM (k2 INT, k1 INT, j INT , i INT, PRIMARY 
KEY(i, j) NOT ENFORCED) WITH"
-                        + " ('continuous.discovery-interval'='1 ms', 
'lookup.cache'='%s', 'bucket' = '2', 'bucket-key' = 'j')",
-                mode);
+                        + " ('continuous.discovery-interval'='1 ms', 
'lookup.cache'='auto', 'bucket' = '2', 'bucket-key' = 'j')");
 
         sql("CREATE TABLE T2 (j INT, i INT, `proctime` AS PROCTIME())");
 
@@ -990,13 +394,11 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         iterator.close();
     }
 
-    @ParameterizedTest
-    @ValueSource(booleans = {true, false})
-    public void testOverwriteDimTable(boolean isPkTable) throws Exception {
+    @Test
+    public void testOverwriteNonPkDimTable() throws Exception {
         sql(
-                "CREATE TABLE DIM (i INT %s, v int, pt STRING) "
-                        + "PARTITIONED BY (pt) WITH 
('continuous.discovery-interval'='1 ms')",
-                isPkTable ? "PRIMARY KEY NOT ENFORCED" : "");
+                "CREATE TABLE DIM (i INT, v int, pt STRING) "
+                        + "PARTITIONED BY (pt) WITH 
('continuous.discovery-interval'='1 ms')");
 
         BlockingIterator<Row, Row> iterator =
                 streamSqlBlockIter(
@@ -1018,58 +420,15 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         iterator.close();
     }
 
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testLookupPartitionLevelMaxPt(LookupCacheMode mode) throws 
Exception {
-        sql(
-                "CREATE TABLE PARTITIONED_DIM (pt1 STRING, pt2 INT, i INT, v 
INT)"
-                        + "PARTITIONED BY (`pt1`, `pt2`) WITH ("
-                        + "'scan.partitions' = 'pt1=max_pt()', "
-                        + "'lookup.dynamic-partition.refresh-interval' = '1 
ms', "
-                        + "'lookup.cache' = '%s', "
-                        + "'continuous.discovery-interval'='1 ms')",
-                mode);
-
-        String query =
-                "SELECT D.pt1, D.pt2, T.i, D.v FROM T LEFT JOIN 
PARTITIONED_DIM for SYSTEM_TIME AS OF T.proctime AS D ON T.i = D.i";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql(
-                "INSERT INTO PARTITIONED_DIM VALUES ('202415', 14, 1, 1), 
('202415', 15, 1, 1), ('202414', 15, 1, 1)");
-        Thread.sleep(500); // wait refresh
-        sql("INSERT INTO T VALUES (1)");
-        List<Row> result = iterator.collect(2);
-        assertThat(result)
-                .containsExactlyInAnyOrder(Row.of("202415", 14, 1, 1), 
Row.of("202415", 15, 1, 1));
-
-        sql("INSERT INTO PARTITIONED_DIM VALUES ('202416', 14, 2, 2), 
('202416', 15, 2, 2)");
-        Thread.sleep(500); // wait refresh
-        sql("INSERT INTO T VALUES (2)");
-        result = iterator.collect(2);
-        assertThat(result)
-                .containsExactlyInAnyOrder(Row.of("202416", 14, 2, 2), 
Row.of("202416", 15, 2, 2));
-
-        sql("ALTER TABLE PARTITIONED_DIM DROP PARTITION (pt1 = '202416',pt2 = 
'15')");
-        Thread.sleep(500); // wait refresh
-        sql("INSERT INTO T VALUES (1), (2)");
-        result = iterator.collect(2);
-        assertThat(result)
-                .containsExactlyInAnyOrder(Row.of(null, null, 1, null), 
Row.of("202416", 14, 2, 2));
-
-        iterator.close();
-    }
-
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testLookupMultiPartitionLevelMaxPt(LookupCacheMode mode) 
throws Exception {
+    @Test
+    public void testLookupMultiPartitionLevelMaxPt() throws Exception {
         sql(
                 "CREATE TABLE PARTITIONED_DIM (pt1 STRING, pt2 INT, pt3 INT, i 
INT, v INT)"
                         + "PARTITIONED BY (`pt1`, `pt2`, `pt3`) WITH ("
                         + "'scan.partitions' = 'pt1=max_pt(),pt2=max_pt()', "
                         + "'lookup.dynamic-partition.refresh-interval' = '1 
ms', "
-                        + "'lookup.cache' = '%s', "
-                        + "'continuous.discovery-interval'='1 ms')",
-                mode);
+                        + "'lookup.cache' = 'full', "
+                        + "'continuous.discovery-interval'='1 ms')");
 
         String query =
                 "SELECT D.pt1, D.pt2, D.pt3, T.i, D.v FROM T LEFT JOIN 
PARTITIONED_DIM for SYSTEM_TIME AS OF T.proctime AS D ON T.i = D.i";
@@ -1103,55 +462,14 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         iterator.close();
     }
 
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testLookupMaxTwoPt0(LookupCacheMode mode) throws Exception {
-        sql(
-                "CREATE TABLE PARTITIONED_DIM (pt STRING, i INT, v INT)"
-                        + "PARTITIONED BY (`pt`) WITH ("
-                        + "'lookup.dynamic-partition' = 'max_two_pt()', "
-                        + "'lookup.dynamic-partition.refresh-interval' = '1 
ms', "
-                        + "'lookup.cache' = '%s', "
-                        + "'continuous.discovery-interval'='1 ms')",
-                mode);
-
-        String query =
-                "SELECT D.pt, T.i, D.v FROM T LEFT JOIN PARTITIONED_DIM for 
SYSTEM_TIME AS OF T.proctime AS D ON T.i = D.i";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO PARTITIONED_DIM VALUES ('2024-10-01', 1, 1), 
('2024-10-01', 2, 2)");
-        Thread.sleep(500); // wait refresh
-        sql("INSERT INTO T VALUES (1)");
-        List<Row> result = iterator.collect(1);
-        assertThat(result).containsExactlyInAnyOrder(Row.of("2024-10-01", 1, 
1));
-
-        sql("INSERT INTO PARTITIONED_DIM VALUES ('2024-10-02', 2, 2)");
-        Thread.sleep(500); // wait refresh
-        sql("INSERT INTO T VALUES (2)");
-        result = iterator.collect(2);
-        assertThat(result)
-                .containsExactlyInAnyOrder(Row.of("2024-10-01", 2, 2), 
Row.of("2024-10-02", 2, 2));
-
-        sql("ALTER TABLE PARTITIONED_DIM DROP PARTITION (pt = '2024-10-01')");
-        Thread.sleep(500); // wait refresh
-        sql("INSERT INTO T VALUES (1), (2)");
-        result = iterator.collect(2);
-        assertThat(result)
-                .containsExactlyInAnyOrder(Row.of(null, 1, null), 
Row.of("2024-10-02", 2, 2));
-
-        iterator.close();
-    }
-
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testLookupSpecifiedPartition(LookupCacheMode mode) throws 
Exception {
+    @Test
+    public void testLookupSpecifiedPartition() throws Exception {
         sql(
                 "CREATE TABLE PARTITIONED_DIM (pt STRING, k INT, v INT, 
PRIMARY KEY (pt, k) NOT ENFORCED) "
                         + "PARTITIONED BY (pt) WITH ( "
                         + "'bucket' = '1', "
-                        + "'lookup.cache' = '%s', "
-                        + "'continuous.discovery-interval'='1 ms')",
-                mode);
+                        + "'lookup.cache' = 'full', "
+                        + "'continuous.discovery-interval'='1 ms')");
 
         sql("INSERT INTO T VALUES (1), (2)");
         sql(
@@ -1224,41 +542,7 @@ public class LookupJoinITCase extends CatalogITCaseBase {
     }
 
     @Test
-    public void testFallbackCacheMode() throws Exception {
-        sql(
-                "CREATE TABLE DIM_WITH_SEQUENCE (i INT PRIMARY KEY NOT 
ENFORCED, j INT, k1 INT, k2 INT) WITH"
-                        + " ('continuous.discovery-interval'='1 ms', 
'sequence.field' = 'j', 'bucket' = '1')");
-        sql("INSERT INTO DIM_WITH_SEQUENCE VALUES (1, 11, 111, 1111), (2, 22, 
222, 2222)");
-
-        String query =
-                "SELECT T.i, D.j, D.k1, D.k2 FROM T LEFT JOIN 
DIM_WITH_SEQUENCE for system_time as of T.proctime AS D ON T.i = D.i";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (1), (2), (3)");
-        List<Row> result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, 11, 111, 1111),
-                        Row.of(2, 22, 222, 2222),
-                        Row.of(3, null, null, null));
-
-        sql("INSERT INTO DIM_WITH_SEQUENCE VALUES (2, 11, 444, 4444), (3, 33, 
333, 3333)");
-        Thread.sleep(2000); // wait refresh
-        sql("INSERT INTO T VALUES (1), (2), (3), (4)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of(1, 11, 111, 1111),
-                        Row.of(2, 22, 222, 2222), // not change
-                        Row.of(3, 33, 333, 3333),
-                        Row.of(4, null, null, null));
-
-        iterator.close();
-    }
-
-    @ParameterizedTest
-    @EnumSource(LookupCacheMode.class)
-    public void testAsyncPartitionRefresh(LookupCacheMode mode) throws 
Exception {
+    public void testAsyncPartitionRefresh() throws Exception {
         // This test verifies asynchronous partition refresh:
         // when max_pt() changes, the lookup table is refreshed in a 
background thread,
         // old partition data continues serving queries until the new 
partition is fully loaded.
@@ -1269,9 +553,8 @@ public class LookupJoinITCase extends CatalogITCaseBase {
                         + "'lookup.dynamic-partition' = 'max_pt()', "
                         + "'lookup.dynamic-partition.refresh-interval' = '1 
ms', "
                         + "'lookup.dynamic-partition.refresh.async' = 'true', "
-                        + "'lookup.cache' = '%s', "
-                        + "'continuous.discovery-interval'='1 ms')",
-                mode);
+                        + "'lookup.cache' = 'full', "
+                        + "'continuous.discovery-interval'='1 ms')");
 
         // insert data into partition '1'
         sql("INSERT INTO PARTITIONED_DIM VALUES ('1', 1, 100), ('1', 2, 200)");
@@ -1286,75 +569,27 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         List<Row> result = iterator.collect(2);
         assertThat(result).containsExactlyInAnyOrder(Row.of(1, 100), Row.of(2, 
200));
 
-        // insert data into a new partition '2', which will trigger async 
partition refresh
+        // The triggering lookup keeps serving the old partition until the 
background load swaps.
         sql("INSERT INTO PARTITIONED_DIM VALUES ('2', 1, 1000), ('2', 2, 
2000)");
-        Thread.sleep(500); // wait for async refresh to complete
-        // trigger a lookup to check async completion and switch to new 
partition
-        sql("INSERT INTO T VALUES (1), (2)");
-        iterator.collect(2);
-        Thread.sleep(500);
         sql("INSERT INTO T VALUES (1), (2)");
         result = iterator.collect(2);
-        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 1000), 
Row.of(2, 2000));
+        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 100), Row.of(2, 
200));
 
-        // insert another new partition '3' and verify switch again
-        sql("INSERT INTO PARTITIONED_DIM VALUES ('3', 1, 10000), ('3', 2, 
20000)");
-        Thread.sleep(500); // wait for async refresh to complete
-        sql("INSERT INTO T VALUES (1), (2)");
-        iterator.collect(2);
         Thread.sleep(500);
         sql("INSERT INTO T VALUES (1), (2)");
         result = iterator.collect(2);
-        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 10000), 
Row.of(2, 20000));
-
-        iterator.close();
-    }
-
-    @ParameterizedTest
-    @EnumSource(
-            value = LookupCacheMode.class,
-            names = {"FULL", "MEMORY"})
-    public void 
testAsyncPartitionRefreshServesOldDataDuringRefresh(LookupCacheMode mode)
-            throws Exception {
-        // Verify that during async refresh, queries still return old 
partition data
-        // until the new partition is fully loaded and switched.
-        sql(
-                "CREATE TABLE PARTITIONED_DIM (pt STRING, k INT, v INT, 
PRIMARY KEY (pt, k) NOT ENFORCED)"
-                        + "PARTITIONED BY (`pt`) WITH ("
-                        + "'bucket' = '1', "
-                        + "'lookup.dynamic-partition' = 'max_pt()', "
-                        + "'lookup.dynamic-partition.refresh-interval' = '1 
ms', "
-                        + "'lookup.dynamic-partition.refresh.async' = 'true', "
-                        + "'lookup.cache' = '%s', "
-                        + "'continuous.discovery-interval'='1 ms')",
-                mode);
-
-        sql("INSERT INTO PARTITIONED_DIM VALUES ('1', 1, 100), ('1', 2, 200)");
-
-        String query =
-                "SELECT T.i, D.v FROM T LEFT JOIN PARTITIONED_DIM "
-                        + "for system_time as of T.proctime AS D ON T.i = D.k";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (1), (2)");
-        List<Row> result = iterator.collect(2);
-        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 100), Row.of(2, 
200));
-
-        // insert new partition '2' to trigger async refresh
-        sql("INSERT INTO PARTITIONED_DIM VALUES ('2', 1, 1000), ('2', 2, 
2000)");
+        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 1000), 
Row.of(2, 2000));
 
-        // immediately query before async refresh completes — should still 
return old partition data
+        // Repeat the transition to verify more than one asynchronous swap.
+        sql("INSERT INTO PARTITIONED_DIM VALUES ('3', 1, 10000), ('3', 2, 
20000)");
         sql("INSERT INTO T VALUES (1), (2)");
         result = iterator.collect(2);
-        // old partition data (100, 200) should still be served
-        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 100), Row.of(2, 
200));
+        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 1000), 
Row.of(2, 2000));
 
-        // now wait for async refresh to complete and trigger switch
         Thread.sleep(500);
         sql("INSERT INTO T VALUES (1), (2)");
         result = iterator.collect(2);
-        // after switch, new partition data should be returned
-        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 1000), 
Row.of(2, 2000));
+        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 10000), 
Row.of(2, 20000));
 
         iterator.close();
     }
@@ -1369,9 +604,8 @@ public class LookupJoinITCase extends CatalogITCaseBase {
                         + "'scan.partitions' = 'pt1=max_pt()', "
                         + "'lookup.dynamic-partition.refresh-interval' = '1 
ms', "
                         + "'lookup.dynamic-partition.refresh.async' = 'true', "
-                        + "'lookup.cache' = '%s', "
-                        + "'continuous.discovery-interval'='1 ms')",
-                LookupCacheMode.FULL);
+                        + "'lookup.cache' = 'full', "
+                        + "'continuous.discovery-interval'='1 ms')");
 
         sql(
                 "INSERT INTO PARTITIONED_DIM VALUES "
@@ -1397,9 +631,15 @@ public class LookupJoinITCase extends CatalogITCaseBase {
                 "INSERT INTO PARTITIONED_DIM VALUES "
                         + "('2025', 1, 1, 1000), ('2025', 1, 2, 2000), "
                         + "('2025', 2, 1, 3000), ('2025', 2, 2, 4000)");
-        Thread.sleep(500);
         sql("INSERT INTO T VALUES (1), (2)");
-        iterator.collect(4);
+        result = iterator.collect(4);
+        assertThat(result)
+                .containsExactlyInAnyOrder(
+                        Row.of("2024", 1, 1, 100),
+                        Row.of("2024", 1, 2, 200),
+                        Row.of("2024", 2, 1, 300),
+                        Row.of("2024", 2, 2, 400));
+
         Thread.sleep(500);
         sql("INSERT INTO T VALUES (1), (2)");
         result = iterator.collect(4);
@@ -1414,124 +654,7 @@ public class LookupJoinITCase extends CatalogITCaseBase {
     }
 
     @Test
-    public void testAsyncPartitionRefreshWithOverwrite() throws Exception {
-        // Verify async partition refresh works correctly when a new max 
partition
-        // is created via INSERT OVERWRITE.
-        sql(
-                "CREATE TABLE PARTITIONED_DIM (pt INT, k INT, v INT, PRIMARY 
KEY (pt, k) NOT ENFORCED)"
-                        + "PARTITIONED BY (`pt`) WITH ("
-                        + "'bucket' = '1', "
-                        + "'lookup.dynamic-partition' = 'max_pt()', "
-                        + "'lookup.dynamic-partition.refresh-interval' = '1 
ms', "
-                        + "'lookup.dynamic-partition.refresh.async' = 'true', "
-                        + "'lookup.cache' = '%s', "
-                        + "'continuous.discovery-interval'='1 ms')",
-                LookupCacheMode.FULL);
-
-        sql("INSERT INTO PARTITIONED_DIM VALUES (1, 1, 100), (1, 2, 200)");
-
-        String query =
-                "SELECT T.i, D.v FROM T LEFT JOIN PARTITIONED_DIM "
-                        + "for system_time as of T.proctime AS D ON T.i = D.k";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (1), (2)");
-        List<Row> result = iterator.collect(2);
-        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 100), Row.of(2, 
200));
-
-        // overwrite current max partition with new data
-        sql("INSERT OVERWRITE PARTITIONED_DIM PARTITION (pt = 1) VALUES (1, 
150), (2, 250)");
-        Thread.sleep(500);
-        sql("INSERT INTO T VALUES (1), (2)");
-        result = iterator.collect(2);
-        assertThat(result).containsExactlyInAnyOrder(Row.of(1, 150), Row.of(2, 
250));
-
-        // overwrite to create a new max partition
-        sql(
-                "INSERT OVERWRITE PARTITIONED_DIM PARTITION (pt = 2) VALUES 
(1, 1000), (2, 2000), (3, 3000)");
-        Thread.sleep(500);
-        sql("INSERT INTO T VALUES (1), (2), (3)");
-        iterator.collect(3);
-        Thread.sleep(500);
-        sql("INSERT INTO T VALUES (1), (2), (3)");
-        result = iterator.collect(3);
-        assertThat(result)
-                .containsExactlyInAnyOrder(Row.of(1, 1000), Row.of(2, 2000), 
Row.of(3, 3000));
-
-        iterator.close();
-    }
-
-    @Test
-    public void testAsyncPartitionRefreshWithMaxTwoPt() throws Exception {
-        // Verify async partition refresh works correctly with max_two_pt() 
strategy.
-        sql(
-                "CREATE TABLE TWO_PT_DIM (pt STRING, k INT, v INT, PRIMARY KEY 
(pt, k) NOT ENFORCED)"
-                        + "PARTITIONED BY (`pt`) WITH ("
-                        + "'bucket' = '1', "
-                        + "'lookup.dynamic-partition' = 'max_two_pt()', "
-                        + "'lookup.dynamic-partition.refresh-interval' = '1 
ms', "
-                        + "'lookup.dynamic-partition.refresh.async' = 'true', "
-                        + "'lookup.cache' = '%s', "
-                        + "'continuous.discovery-interval'='1 ms')",
-                LookupCacheMode.FULL);
-
-        // insert data into partitions '1' and '2'
-        sql(
-                "INSERT INTO TWO_PT_DIM VALUES "
-                        + "('1', 1, 100), ('1', 2, 200), "
-                        + "('2', 1, 300), ('2', 2, 400)");
-
-        String query =
-                "SELECT D.pt, T.i, D.v FROM T LEFT JOIN TWO_PT_DIM "
-                        + "for system_time as of T.proctime AS D ON T.i = D.k";
-        BlockingIterator<Row, Row> iterator = 
BlockingIterator.of(sEnv.executeSql(query).collect());
-
-        sql("INSERT INTO T VALUES (1), (2)");
-        List<Row> result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of("1", 1, 100),
-                        Row.of("1", 2, 200),
-                        Row.of("2", 1, 300),
-                        Row.of("2", 2, 400));
-
-        // insert new partition '3', now max_two_pt should be '2' and '3'
-        sql("INSERT INTO TWO_PT_DIM VALUES " + "('3', 1, 1000), ('3', 2, 
2000)");
-        sql("INSERT INTO T VALUES (1), (2)");
-        iterator.collect(4);
-        Thread.sleep(500);
-        sql("INSERT INTO T VALUES (1), (2)");
-        result = iterator.collect(4);
-        // should now see data from partitions '2' and '3'
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of("2", 1, 300),
-                        Row.of("2", 2, 400),
-                        Row.of("3", 1, 1000),
-                        Row.of("3", 2, 2000));
-
-        // insert another partition '4', max_two_pt should be '3' and '4'
-        sql("INSERT INTO TWO_PT_DIM VALUES " + "('4', 1, 10000), ('4', 2, 
20000)");
-        sql("INSERT INTO T VALUES (1), (2)");
-        iterator.collect(4);
-        Thread.sleep(500);
-        sql("INSERT INTO T VALUES (1), (2)");
-        result = iterator.collect(4);
-        assertThat(result)
-                .containsExactlyInAnyOrder(
-                        Row.of("3", 1, 1000),
-                        Row.of("3", 2, 2000),
-                        Row.of("4", 1, 10000),
-                        Row.of("4", 2, 20000));
-
-        iterator.close();
-    }
-
-    @ParameterizedTest
-    @EnumSource(
-            value = LookupCacheMode.class,
-            names = {"FULL", "MEMORY"})
-    public void testAsyncPartitionRefreshWithNonPkTable(LookupCacheMode mode) 
throws Exception {
+    public void testAsyncPartitionRefreshWithNonPkTable() throws Exception {
         // Verify async partition refresh works correctly with non-primary-key 
append tables.
         sql(
                 "CREATE TABLE NON_PK_DIM (pt STRING, k INT, v INT)"
@@ -1539,9 +662,8 @@ public class LookupJoinITCase extends CatalogITCaseBase {
                         + "'lookup.dynamic-partition' = 'max_pt()', "
                         + "'lookup.dynamic-partition.refresh-interval' = '1 
ms', "
                         + "'lookup.dynamic-partition.refresh.async' = 'true', "
-                        + "'lookup.cache' = '%s', "
-                        + "'continuous.discovery-interval'='1 ms')",
-                mode);
+                        + "'lookup.cache' = 'full', "
+                        + "'continuous.discovery-interval'='1 ms')");
 
         sql("INSERT INTO NON_PK_DIM VALUES ('1', 1, 100), ('1', 1, 101), ('1', 
2, 200)");
 
@@ -1569,11 +691,8 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         iterator.close();
     }
 
-    @ParameterizedTest
-    @EnumSource(
-            value = LookupCacheMode.class,
-            names = {"FULL", "MEMORY"})
-    public void testLookupBlobAsDescriptorOnNormalBlobTable(LookupCacheMode 
mode) throws Exception {
+    @Test
+    public void testLookupBlobAsDescriptorOnNormalBlobTable() throws Exception 
{
         // Test that lookup.blob-as-descriptor works correctly even when the 
table was NOT
         // written with blob-as-descriptor=true. Previously this would fail 
with
         // "Blob data can not convert to descriptor" because BlobFormatReader 
returned BlobData
@@ -1585,9 +704,8 @@ public class LookupJoinITCase extends CatalogITCaseBase {
                         + "'data-evolution.enabled'='true', "
                         + "'blob-field'='picture', "
                         + "'lookup.blob-as-descriptor'='true', "
-                        + "'lookup.cache'='%s', "
-                        + "'continuous.discovery-interval'='1 ms')",
-                mode);
+                        + "'lookup.cache'='full', "
+                        + "'continuous.discovery-interval'='1 ms')");
 
         // Write raw blob data (NOT as descriptor) — this is the normal write 
path.
         sql("INSERT INTO BLOB_DIM VALUES (1, 'cat', X'48656C6C6F'), (2, 'dog', 
X'576F726C64')");
@@ -1626,21 +744,16 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         iterator.close();
     }
 
-    @ParameterizedTest
-    @EnumSource(
-            value = LookupCacheMode.class,
-            names = {"FULL", "MEMORY"})
-    public void 
testLookupArrayBlobAsDescriptorOnNormalBlobTable(LookupCacheMode mode)
-            throws Exception {
+    @Test
+    public void testLookupArrayBlobAsDescriptorOnNormalBlobTable() throws 
Exception {
         sql(
                 "CREATE TABLE ARRAY_BLOB_DIM (id INT, name STRING, pictures 
ARRAY<BYTES>) WITH ("
                         + "'row-tracking.enabled'='true', "
                         + "'data-evolution.enabled'='true', "
                         + "'blob-field'='pictures', "
                         + "'lookup.blob-as-descriptor'='true', "
-                        + "'lookup.cache'='%s', "
-                        + "'continuous.discovery-interval'='1 ms')",
-                mode);
+                        + "'lookup.cache'='full', "
+                        + "'continuous.discovery-interval'='1 ms')");
 
         String dataId =
                 TestValuesTableFactory.registerData(
@@ -1691,21 +804,16 @@ public class LookupJoinITCase extends CatalogITCaseBase {
         iterator.close();
     }
 
-    @ParameterizedTest
-    @EnumSource(
-            value = LookupCacheMode.class,
-            names = {"FULL", "MEMORY"})
-    public void testLookupMapBlobAsDescriptorOnNormalBlobTable(LookupCacheMode 
mode)
-            throws Exception {
+    @Test
+    public void testLookupMapBlobAsDescriptorOnNormalBlobTable() throws 
Exception {
         sql(
                 "CREATE TABLE MAP_BLOB_DIM (id INT, name STRING, pictures 
MAP<INT, BYTES>) WITH ("
                         + "'row-tracking.enabled'='true', "
                         + "'data-evolution.enabled'='true', "
                         + "'blob-field'='pictures', "
                         + "'lookup.blob-as-descriptor'='true', "
-                        + "'lookup.cache'='%s', "
-                        + "'continuous.discovery-interval'='1 ms')",
-                mode);
+                        + "'lookup.cache'='full', "
+                        + "'continuous.discovery-interval'='1 ms')");
 
         Map<Integer, byte[]> firstPictures = new LinkedHashMap<>();
         firstPictures.put(1, new byte[] {72, 101, 108, 108, 111});
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/BlobAsDescriptorRowTest.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/BlobAsDescriptorRowTest.java
new file mode 100644
index 0000000000..403e8060f2
--- /dev/null
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/BlobAsDescriptorRowTest.java
@@ -0,0 +1,86 @@
+/*
+ * 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.paimon.flink.lookup;
+
+import org.apache.paimon.data.Blob;
+import org.apache.paimon.data.BlobDescriptor;
+import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.types.DataTypeRoot;
+import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.types.RowKind;
+import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.UriReader;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+import java.util.HashSet;
+import java.util.Set;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/** Tests for {@link BlobAsDescriptorRow}. */
+class BlobAsDescriptorRowTest {
+
+    @Test
+    void testConvertBlobToDescriptor() {
+        BlobDescriptor descriptor = new 
BlobDescriptor("https://example.com/blob";, 12, 34);
+        GenericRow wrapped =
+                GenericRow.of(7, Blob.fromDescriptor(UriReader.fromHttp(), 
descriptor), null);
+        BlobAsDescriptorRow row =
+                new BlobAsDescriptorRow(wrapped, new 
HashSet<>(Arrays.asList(1, 2)));
+
+        assertThat(row.getInt(0)).isEqualTo(7);
+        
assertThat(BlobDescriptor.deserialize(row.getBinary(1))).isEqualTo(descriptor);
+        assertThat(row.getBinary(2)).isNull();
+
+        row.setRowKind(RowKind.DELETE);
+        assertThat(row.getRowKind()).isEqualTo(RowKind.DELETE);
+        assertThat(wrapped.getRowKind()).isEqualTo(RowKind.DELETE);
+    }
+
+    @Test
+    void testReplaceBlobWithVarBinary() {
+        RowType rowType =
+                RowType.builder()
+                        .field("id", DataTypes.INT())
+                        .field("picture", DataTypes.BLOB())
+                        .field("raw", DataTypes.BYTES())
+                        .build();
+
+        Set<Integer> blobPositions = 
BlobAsDescriptorRow.blobFieldPositions(rowType);
+        RowType converted = 
BlobAsDescriptorRow.replaceBlobWithVarBinary(rowType, blobPositions);
+
+        assertThat(blobPositions).containsExactlyInAnyOrder(1);
+        assertThat(converted.getTypeAt(0)).isEqualTo(DataTypes.INT());
+        
assertThat(converted.getTypeAt(1).getTypeRoot()).isEqualTo(DataTypeRoot.VARBINARY);
+        assertThat(converted.getTypeAt(2)).isEqualTo(DataTypes.BYTES());
+        assertThat(rowType.getTypeAt(1)).isEqualTo(DataTypes.BLOB());
+    }
+
+    @Test
+    void testKeepOriginalTypeWhenThereIsNoBlob() {
+        RowType rowType = RowType.of(DataTypes.INT(), DataTypes.BYTES());
+
+        assertThat(
+                        BlobAsDescriptorRow.replaceBlobWithVarBinary(
+                                rowType, 
BlobAsDescriptorRow.blobFieldPositions(rowType)))
+                .isSameAs(rowType);
+    }
+}
diff --git 
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/DynamicPartitionNumberLoaderTest.java
 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/DynamicPartitionNumberLoaderTest.java
new file mode 100644
index 0000000000..6c84cca6ef
--- /dev/null
+++ 
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/lookup/DynamicPartitionNumberLoaderTest.java
@@ -0,0 +1,150 @@
+/*
+ * 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.paimon.flink.lookup;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.data.GenericRow;
+import org.apache.paimon.flink.FlinkConnectorOptions;
+import org.apache.paimon.options.Options;
+import org.apache.paimon.schema.Schema;
+import org.apache.paimon.schema.SchemaManager;
+import org.apache.paimon.schema.TableSchema;
+import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.FileStoreTableFactory;
+import org.apache.paimon.table.sink.TableCommitImpl;
+import org.apache.paimon.table.sink.TableWriteImpl;
+import org.apache.paimon.types.DataType;
+import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.TraceableFileIO;
+
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import java.nio.file.Path;
+import java.time.Duration;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import java.util.UUID;
+import java.util.stream.Collectors;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/** Tests for {@link DynamicPartitionNumberLoader}. */
+class DynamicPartitionNumberLoaderTest {
+
+    @TempDir private Path tempDir;
+
+    private final String commitUser = UUID.randomUUID().toString();
+    private final TraceableFileIO fileIO = new TraceableFileIO();
+
+    private FileStoreTable table;
+    private TableWriteImpl<?> write;
+    private TableCommitImpl commit;
+
+    @BeforeEach
+    void before() throws Exception {
+        table = createFileStoreTable();
+        write = table.newWrite(commitUser);
+        commit = table.newCommit(commitUser);
+    }
+
+    @AfterEach
+    void after() throws Exception {
+        write.close();
+        commit.close();
+    }
+
+    @Test
+    void testMaxTwoPartitionsAndRefresh() throws Exception {
+        writePartition("2024", 1);
+        writePartition("2025", 2);
+        commit.commit(1, write.prepareCommit(true, 1));
+
+        DynamicPartitionLoader loader = createLoader("max_two_pt()");
+        assertThat(loader.checkRefresh()).isTrue();
+        assertThat(partitions(loader)).containsExactly("2025", "2024");
+        assertThat(loader.checkRefresh()).isFalse();
+
+        writePartition("2026", 3);
+        commit.commit(2, write.prepareCommit(true, 2));
+
+        assertThat(loader.checkRefresh()).isTrue();
+        assertThat(partitions(loader)).containsExactly("2026", "2025");
+    }
+
+    @Test
+    void testMaxPartition() throws Exception {
+        writePartition("2024", 1);
+        writePartition("2025", 2);
+        commit.commit(1, write.prepareCommit(true, 1));
+
+        DynamicPartitionLoader loader = createLoader("max_pt()");
+        assertThat(loader.checkRefresh()).isTrue();
+        assertThat(partitions(loader)).containsExactly("2025");
+    }
+
+    private DynamicPartitionLoader createLoader(String scanPartitions) {
+        FileStoreTable configuredTable =
+                table.copy(
+                        Collections.singletonMap(
+                                FlinkConnectorOptions.SCAN_PARTITIONS.key(), 
scanPartitions));
+        DynamicPartitionLoader loader =
+                (DynamicPartitionLoader) PartitionLoader.of(configuredTable);
+        loader.open();
+        return loader;
+    }
+
+    private void writePartition(String partition, int value) throws Exception {
+        write.write(GenericRow.of(BinaryString.fromString(partition), value, 
(long) value));
+    }
+
+    private List<String> partitions(DynamicPartitionLoader loader) {
+        return loader.partitions().stream()
+                .map(row -> row.getString(0))
+                .map(BinaryString::toString)
+                .collect(Collectors.toList());
+    }
+
+    private FileStoreTable createFileStoreTable() throws Exception {
+        org.apache.paimon.fs.Path tablePath = new 
org.apache.paimon.fs.Path(tempDir.toString());
+        SchemaManager schemaManager = new SchemaManager(fileIO, tablePath);
+        Options options = new Options();
+        options.set(CoreOptions.BUCKET, 1);
+        
options.set(FlinkConnectorOptions.LOOKUP_DYNAMIC_PARTITION_REFRESH_INTERVAL, 
Duration.ZERO);
+
+        RowType rowType =
+                RowType.of(
+                        new DataType[] {DataTypes.STRING(), DataTypes.INT(), 
DataTypes.BIGINT()},
+                        new String[] {"pt", "k", "v"});
+        Schema schema =
+                new Schema(
+                        rowType.getFields(),
+                        Collections.singletonList("pt"),
+                        Arrays.asList("pt", "k"),
+                        options.toMap(),
+                        "");
+        TableSchema tableSchema = schemaManager.createTable(schema);
+        return FileStoreTableFactory.create(fileIO, tablePath, tableSchema);
+    }
+}

Reply via email to