fhueske commented on code in PR #29368: URL: https://github.com/apache/flink/pull/29368#discussion_r4217692197
########## flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/join/LateralSnapshotJoinTableApiTest.java: ########## @@ -0,0 +1,407 @@ +/* + * 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.flink.table.planner.plan.stream.sql.join; + +import org.apache.flink.table.api.ApiExpression; +import org.apache.flink.table.api.DataTypes; +import org.apache.flink.table.api.Table; +import org.apache.flink.table.api.TableConfig; +import org.apache.flink.table.api.ValidationException; +import org.apache.flink.table.api.config.TableConfigOptions.ColumnExpansionStrategy; +import org.apache.flink.table.planner.utils.TableTestBase; +import org.apache.flink.table.planner.utils.TableTestUtil; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.MethodSource; + +import java.time.Duration; +import java.time.LocalDateTime; +import java.time.ZoneId; +import java.util.Collections; +import java.util.List; +import java.util.stream.Collectors; +import java.util.stream.Stream; + +import static org.apache.flink.table.api.Expressions.$; +import static org.apache.flink.table.api.Expressions.call; +import static org.apache.flink.table.api.Expressions.descriptor; +import static org.apache.flink.table.api.Expressions.lit; +import static org.apache.flink.table.api.config.TableConfigOptions.ColumnExpansionStrategy.EXCLUDE_DEFAULT_VIRTUAL_METADATA_COLUMNS; +import static org.apache.flink.table.api.config.TableConfigOptions.TABLE_COLUMN_EXPANSION_STRATEGY; +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/** + * Verifies that the Table API surface of the {@code LATERAL SNAPSHOT} join produces the exact same + * optimized physical plan as the equivalent SQL. + * + * <p>This is the primary coverage for the Table API surface: it converges on the same {@code + * LateralSnapshotJoin} {@code RelNode} / {@code StreamExecLateralSnapshotJoin} operator as SQL, + * whose execution semantics are already fully covered by the SQL semantic tests ({@code + * LateralSnapshotJoinSemanticTests}) and {@code LateralSnapshotJoinITCase}. {@code + * LateralSnapshotJoinSemanticTestPrograms#INNER_JOIN_TABLE_API} keeps a single execution smoke on + * top of this. + */ +public class LateralSnapshotJoinTableApiTest extends TableTestBase { + + private static final String LOAD_COMPLETED_TIME_SQL = + "CAST(TIMESTAMP '2026-07-01 00:00:00' AS TIMESTAMP_LTZ(3))"; + + private static final ApiExpression LOAD_COMPLETED_TIME_TABLE = + lit(LocalDateTime.parse("2026-07-01T00:00:00")).cast(DataTypes.TIMESTAMP_LTZ(3)); + + private TableTestUtil util; + + @BeforeEach + void setup() { + final TableConfig config = TableConfig.getDefault(); + config.setLocalTimeZone(ZoneId.of("UTC")); + util = streamTestUtil(config); + + util.tableEnv() + .executeSql( + "CREATE TABLE probe (" + + " pk STRING," + + " pv INT," + + " pts TIMESTAMP(3)," + + " WATERMARK FOR pts AS pts" + + ") WITH ('connector' = 'values', 'bounded' = 'false')"); + + util.tableEnv() + .executeSql( + "CREATE TABLE b (" + + " bk STRING," + + " bv INT," + + " bts TIMESTAMP(3)," + + " WATERMARK FOR bts AS bts" + + ") WITH (" + + " 'connector' = 'values'," + + " 'bounded' = 'false'," + + " 'changelog-mode' = 'I,UB,UA,D'" + + ")"); + } + + // ------------------------------------------------------------------------------------------ + // Plan equivalence: the Table API surface should produce the same optimized plan as SQL. + // ------------------------------------------------------------------------------------------ + + @Test + void testInnerJoinPlanParity() { + final String sql = + "SELECT probe.pk, probe.pv, s.bk, s.bv FROM probe JOIN LATERAL SNAPSHOT(" + + "input => TABLE b, on_time => DESCRIPTOR(bts), " + + "load_completed_time => " + + LOAD_COMPLETED_TIME_SQL + + ") AS s " + + "ON probe.pk = s.bk"; + + final Table apiResult = + util.tableEnv() + .from("probe") + .joinLateral( + call( + "SNAPSHOT", + util.tableEnv().from("b").asArgument("input"), + descriptor("bts").asArgument("on_time"), + LOAD_COMPLETED_TIME_TABLE.asArgument( + "load_completed_time")), + $("pk").isEqual($("bk"))) + .select($("pk"), $("pv"), $("bk"), $("bv")); + + assertSameOptimizedPhysicalPlan(sql, apiResult); + } + + @Test + void testLeftOuterJoinPlanParity() { + final String sql = + "SELECT probe.pk, probe.pv, s.bk, s.bv FROM probe LEFT JOIN LATERAL SNAPSHOT(" + + "input => TABLE b, on_time => DESCRIPTOR(bts), " + + "load_completed_time => " + + LOAD_COMPLETED_TIME_SQL + + ") AS s " + + "ON probe.pk = s.bk"; + + final Table apiResult = + util.tableEnv() + .from("probe") + .leftOuterJoinLateral( + call( + "SNAPSHOT", + util.tableEnv().from("b").asArgument("input"), + descriptor("bts").asArgument("on_time"), + LOAD_COMPLETED_TIME_TABLE.asArgument( + "load_completed_time")), + $("pk").isEqual($("bk"))) + .select($("pk"), $("pv"), $("bk"), $("bv")); + + assertSameOptimizedPhysicalPlan(sql, apiResult); + } + + @Test + void testTableApiBuildSidePlanParity() { + // SQL build side goes through a named view (CTE) to mirror how SQL can only reference a + // TABLE argument by name; the Table API passes the transformed Table directly. + final String sql = + "WITH cte AS (SELECT * FROM b WHERE bv > 10) " + + "SELECT probe.pk, probe.pv, s.bk, s.bv FROM probe JOIN LATERAL SNAPSHOT(" + + "input => TABLE cte, on_time => DESCRIPTOR(bts), " + + "load_completed_time => " + + LOAD_COMPLETED_TIME_SQL + + ") AS s " + + "ON probe.pk = s.bk"; + + final Table filteredBuild = util.tableEnv().from("b").filter($("bv").isGreater(10)); + final Table apiResult = + util.tableEnv() + .from("probe") + .joinLateral( + call( + "SNAPSHOT", + filteredBuild.asArgument("input"), + descriptor("bts").asArgument("on_time"), + LOAD_COMPLETED_TIME_TABLE.asArgument( + "load_completed_time")), + $("pk").isEqual($("bk"))) + .select($("pk"), $("pv"), $("bk"), $("bv")); + + assertSameOptimizedPhysicalPlan(sql, apiResult); + } + + @Test + void testAliasedColumnsPlanParity() { + // No narrowing projection on top: SELECT * keeps the aliased build columns, so this + // asserts both plan parity and that the Table API and SQL agree on the (aliased) schema. + final String sql = + "SELECT * FROM probe JOIN LATERAL SNAPSHOT(" + + "input => TABLE b, on_time => DESCRIPTOR(bts), " + + "load_completed_time => " + + LOAD_COMPLETED_TIME_SQL + + ") AS s(k, v, t) " + + "ON probe.pk = s.k"; + + final Table apiResult = + util.tableEnv() + .from("probe") + .joinLateral( + call( + "SNAPSHOT", + util.tableEnv().from("b").asArgument("input"), + descriptor("bts").asArgument("on_time"), + LOAD_COMPLETED_TIME_TABLE.asArgument( + "load_completed_time")) + .as("k", "v", "t"), + $("pk").isEqual($("k"))); + + assertSameOptimizedPhysicalPlan(sql, apiResult); + assertThat(apiResult.getResolvedSchema()) + .isEqualTo(util.tableEnv().sqlQuery(sql).getResolvedSchema()); + } + + @Test + void testOptionalArgsPlanParity() { + // The optional scalar SNAPSHOT arguments (load_completed_idle_timeout, state_ttl) must be + // passed through the Table API the same way as named SQL arguments and produce the same + // plan. The order of Table API args is shuffled to assert correct ordering of named args. + final String sql = + "SELECT probe.pk, probe.pv, s.bk, s.bv FROM probe JOIN LATERAL SNAPSHOT(" + + "input => TABLE b, on_time => DESCRIPTOR(bts), " + + "load_completed_time => " + + LOAD_COMPLETED_TIME_SQL + + ", load_completed_idle_timeout => INTERVAL '10' SECOND" + + ", state_ttl => INTERVAL '1' DAY" + + ") AS s " + + "ON probe.pk = s.bk"; + + final Table apiResult = + util.tableEnv() + .from("probe") + .joinLateral( + call( + "SNAPSHOT", + LOAD_COMPLETED_TIME_TABLE.asArgument("load_completed_time"), + lit(Duration.ofDays(1)).asArgument("state_ttl"), + descriptor("bts").asArgument("on_time"), + lit(Duration.ofSeconds(10)) + .asArgument("load_completed_idle_timeout"), + util.tableEnv().from("b").asArgument("input")), + $("pk").isEqual($("bk"))) + .select($("pk"), $("pv"), $("bk"), $("bv")); + + assertSameOptimizedPhysicalPlan(sql, apiResult); + } + + @Test + void testDefaultLoadCompletedTimePlanParity() { Review Comment: You're right, nothing Table API specific here — defaulting an omitted optional arg is generic PTF handling. Removed this test (and the optional-args and column-expansion ones for the same reason). -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
