From 62480be5056a9c31daf22c8f9a161a79fb732a6d Mon Sep 17 00:00:00 2001 From: Fabian Hueske Date: Fri, 2 Oct 2026 12:11:38 +0200 Subject: [PATCH 1/3] [FLINK-40885][table] Support LATERAL SNAPSHOT join in the Table API Enable Table.joinLateral / leftOuterJoinLateral with a call("SNAPSHOT", ...) build side, producing the same plan as the equivalent SQL. Release note: The LATERAL SNAPSHOT join can now be expressed through the programmatic Table API, not just SQL, via Table.joinLateral / Table.leftOuterJoinLateral with a call("SNAPSHOT", ...) build side. Generated-by: Claude Code (Claude Opus 4.8) --- .../utils/CorrelatedFunctionTableFactory.java | 15 +- .../utils/JoinOperationFactory.java | 20 +- .../planner/plan/QueryOperationConverter.java | 156 ++++-- ...teralSnapshotJoinSemanticTestPrograms.java | 18 +- .../LateralSnapshotJoinSemanticTests.java | 6 +- ...pshotJoinTableApiSemanticTestPrograms.java | 197 +++++++ .../join/LateralSnapshotJoinTableApiTest.java | 499 ++++++++++++++++++ .../sql/join/LateralSnapshotJoinTest.java | 26 + .../sql/join/LateralSnapshotJoinTest.xml | 32 ++ 9 files changed, 902 insertions(+), 67 deletions(-) create mode 100644 flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/LateralSnapshotJoinTableApiSemanticTestPrograms.java create mode 100644 flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/join/LateralSnapshotJoinTableApiTest.java diff --git a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/CorrelatedFunctionTableFactory.java b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/CorrelatedFunctionTableFactory.java index 0a477a3d0d8e1..53c5227148eb8 100644 --- a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/CorrelatedFunctionTableFactory.java +++ b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/CorrelatedFunctionTableFactory.java @@ -27,6 +27,7 @@ import org.apache.flink.table.expressions.ExpressionUtils; import org.apache.flink.table.expressions.ResolvedExpression; import org.apache.flink.table.expressions.utils.ResolvedExpressionDefaultVisitor; +import org.apache.flink.table.functions.BuiltInFunctionDefinitions; import org.apache.flink.table.functions.FunctionDefinition; import org.apache.flink.table.functions.FunctionKind; import org.apache.flink.table.operations.CorrelatedFunctionQueryOperation; @@ -58,6 +59,14 @@ QueryOperation create(ResolvedExpression callExpr, List leftTableFieldNa return callExpr.accept(calculatedTableCreator); } + /** + * Whether the given function is {@code SNAPSHOT}, currently the only {@link + * FunctionKind#PROCESS_TABLE} function that is valid as the build side of a lateral join. + */ + static boolean isSnapshot(FunctionDefinition definition) { + return definition.equals(BuiltInFunctionDefinitions.SNAPSHOT); + } + private static class FunctionTableCallVisitor extends ResolvedExpressionDefaultVisitor { private final List leftTableFieldNames; @@ -91,7 +100,11 @@ private CorrelatedFunctionQueryOperation unwrapFromAlias(CallExpression call) { + alias))) .collect(toList()); - if (!isFunctionOfKind(children.get(0), FunctionKind.TABLE)) { + // SNAPSHOT is PROCESS_TABLE rather than TABLE but is valid in a LATERAL join. + Expression firstChild = children.get(0); + if (!isFunctionOfKind(children.get(0), FunctionKind.TABLE) + && !(firstChild instanceof CallExpression + && isSnapshot(((CallExpression) firstChild).getFunctionDefinition()))) { throw fail(); } diff --git a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/JoinOperationFactory.java b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/JoinOperationFactory.java index 5b5cd81d0c41b..c01123f7ff114 100644 --- a/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/JoinOperationFactory.java +++ b/flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/operations/utils/JoinOperationFactory.java @@ -80,15 +80,29 @@ private void validateCondition( JoinType joinType, ResolvedExpression condition, boolean correlated) { - boolean alwaysTrue = ExpressionUtils.extractValue(condition, Boolean.class).orElse(false); - - if (alwaysTrue) { + final boolean alwaysTrue = + ExpressionUtils.extractValue(condition, Boolean.class).orElse(false); + + // A lateral SNAPSHOT join is rewritten into a dedicated join that supports an ON predicate + // but requires at least one equi-join predicate (for both INNER and LEFT OUTER). An + // always-true/empty condition is therefore not allowed; fall through to the equi-join check + // below so the user gets a clear error instead of a later planner failure. + final boolean isLateralSnapshot = + correlated + && right instanceof CorrelatedFunctionQueryOperation + && CorrelatedFunctionTableFactory.isSnapshot( + ((CorrelatedFunctionQueryOperation) right) + .getResolvedFunction() + .getDefinition()); + + if (alwaysTrue && !isLateralSnapshot) { return; } Boolean equiJoinExists = condition.accept(equiJoinExistsChecker); if (correlated && right instanceof CorrelatedFunctionQueryOperation + && !isLateralSnapshot && joinType != JoinType.INNER) { throw new ValidationException( "Predicate for lateral left outer join with table function can only be empty or literal true."); diff --git a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/QueryOperationConverter.java b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/QueryOperationConverter.java index 468430dbc7b31..298104659699e 100644 --- a/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/QueryOperationConverter.java +++ b/flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/QueryOperationConverter.java @@ -284,7 +284,17 @@ private AggCall getAggCall(Expression aggregateExpression) { @Override public RelNode visit(JoinQueryOperation join) { final Set corSet; - if (join.isCorrelated()) { + // Check if this is a LATERAL SNAPSHOT join. The SNAPSHOT function never references the + // outer row. Emit a plain LogicalJoin instead of a Correlate, so that + // LogicalJoinToLateralSnapshotJoinRule matches. + final QueryOperation right = join.getChildren().get(1); + final boolean isLateralSnapshotJoin = + right instanceof CorrelatedFunctionQueryOperation + && isSnapshot( + ((CorrelatedFunctionQueryOperation) right) + .getResolvedFunction() + .getDefinition()); + if (join.isCorrelated() && !isLateralSnapshotJoin) { corSet = Collections.singleton(relBuilder.peek().getCluster().createCorrel()); } else { corSet = Collections.emptySet(); @@ -352,44 +362,11 @@ public RelNode visit(FunctionQueryOperation functionTable) { .map( resolvedArg -> { if (resolvedArg instanceof TableReferenceExpression) { - final TableReferenceExpression tableRef = - (TableReferenceExpression) resolvedArg; - final LogicalType tableArgType = - tableRef.getOutputDataType().getLogicalType(); - final RelDataType rowType = - typeFactory.buildRelNodeRowType( - (RowType) tableArgType); - final int[] partitionKeys; - final int[] orderKeys; - final SortOrder[] sortOrders; - if (tableRef.getQueryOperation() - instanceof PartitionQueryOperation) { - final PartitionQueryOperation partitionOperation = - (PartitionQueryOperation) - tableRef.getQueryOperation(); - partitionKeys = - partitionOperation.getPartitionKeys(); - orderKeys = partitionOperation.getOrderKeys(); - final SortDirection[] directions = - partitionOperation.getOrderDirections(); - sortOrders = new SortOrder[directions.length]; - for (int i = 0; i < directions.length; i++) { - sortOrders[i] = - SortOrder.fromSortDirection( - directions[i]); - } - } else { - partitionKeys = new int[0]; - orderKeys = new int[0]; - sortOrders = new SortOrder[0]; - } final RexTableArgCall tableArgCall = - new RexTableArgCall( - rowType, + buildTableArgCall( + (TableReferenceExpression) resolvedArg, inputStack.size(), - partitionKeys, - orderKeys, - sortOrders); + typeFactory); inputStack.add(relBuilder.build()); return tableArgCall; } @@ -400,22 +377,9 @@ public RelNode visit(FunctionQueryOperation functionTable) { // relBuilder.build() works in LIFO fashion, this restores the original input order Collections.reverse(inputStack); - final BridgingSqlFunction sqlFunction = - BridgingSqlFunction.of(relBuilder.getCluster(), contextFunction); - - final RexNode call = - relBuilder - .getRexBuilder() - .makeCall(outputRelDataType, sqlFunction, rexNodeArgs); - final RelNode functionScan = - LogicalTableFunctionScan.create( - relBuilder.getCluster(), - inputStack, - call, - null, - outputRelDataType, - Collections.emptySet()); - relBuilder.push(functionScan); + relBuilder.push( + createTableFunctionScan( + contextFunction, rexNodeArgs, inputStack, outputRelDataType)); return relBuilder.build(); } @@ -430,6 +394,11 @@ public RelNode visit(PartitionQueryOperation partition) { public RelNode visit(CorrelatedFunctionQueryOperation correlatedFunction) { final ContextResolvedFunction contextFunction = correlatedFunction.getResolvedFunction(); + // SNAPSHOT is currently the only PTF valid in a lateral join and hence the only + // function accepting table args. + if (isSnapshot(contextFunction.getDefinition())) { + return convertCorrelatedFunctionWithTableArgs(correlatedFunction); + } final List parameters = convertToRexNodes(correlatedFunction.getArguments()); final FunctionDefinition functionDefinition = contextFunction.getDefinition(); @@ -455,6 +424,87 @@ public RelNode visit(CorrelatedFunctionQueryOperation correlatedFunction) { return relBuilder.build(); } + private RelNode convertCorrelatedFunctionWithTableArgs( + CorrelatedFunctionQueryOperation correlatedFunction) { + final FlinkTypeFactory typeFactory = ShortcutUtils.unwrapTypeFactory(relBuilder); + final List inputs = new ArrayList<>(); + final List args = new ArrayList<>(); + for (ResolvedExpression arg : correlatedFunction.getArguments()) { + if (arg instanceof TableReferenceExpression) { + final TableReferenceExpression tableRef = (TableReferenceExpression) arg; + args.add(buildTableArgCall(tableRef, inputs.size(), typeFactory)); + inputs.add(tableRef.getQueryOperation().accept(QueryOperationConverter.this)); + } else { + args.add(convertExprToRexNode(arg)); + } + } + // Derive the output row type like the FROM-clause PTF path + // (FunctionQueryOperation#getOutputDataType). + final RelDataType outputType = + typeFactory.buildRelNodeRowType( + (RowType) + DataTypeUtils.fromResolvedSchemaPreservingTimeAttributes( + correlatedFunction.getResolvedSchema()) + .getLogicalType()); + return createTableFunctionScan( + correlatedFunction.getResolvedFunction(), args, inputs, outputType); + } + + private boolean isSnapshot(FunctionDefinition definition) { + // SNAPSHOT is currently the only PTF valid as the build side of a lateral join; the + // lateral-join-specific handling in this converter keys off this single predicate. + return definition.equals(BuiltInFunctionDefinitions.SNAPSHOT); + } + + /** + * Converts a {@link TableReferenceExpression} table argument into a {@link + * RexTableArgCall}, carrying PARTITION BY / ORDER BY keys when the referenced operation is + * a {@link PartitionQueryOperation}. + */ + private RexTableArgCall buildTableArgCall( + TableReferenceExpression tableRef, int inputIndex, FlinkTypeFactory typeFactory) { + final RelDataType rowType = + typeFactory.buildRelNodeRowType( + (RowType) tableRef.getOutputDataType().getLogicalType()); + final int[] partitionKeys; + final int[] orderKeys; + final SortOrder[] sortOrders; + if (tableRef.getQueryOperation() instanceof PartitionQueryOperation) { + final PartitionQueryOperation partitionOperation = + (PartitionQueryOperation) tableRef.getQueryOperation(); + partitionKeys = partitionOperation.getPartitionKeys(); + orderKeys = partitionOperation.getOrderKeys(); + final SortDirection[] directions = partitionOperation.getOrderDirections(); + sortOrders = new SortOrder[directions.length]; + for (int i = 0; i < directions.length; i++) { + sortOrders[i] = SortOrder.fromSortDirection(directions[i]); + } + } else { + partitionKeys = new int[0]; + orderKeys = new int[0]; + sortOrders = new SortOrder[0]; + } + return new RexTableArgCall(rowType, inputIndex, partitionKeys, orderKeys, sortOrders); + } + + /** Builds a {@link LogicalTableFunctionScan} for a function call with table arguments. */ + private RelNode createTableFunctionScan( + ContextResolvedFunction function, + List args, + List inputs, + RelDataType outputType) { + final BridgingSqlFunction sqlFunction = + BridgingSqlFunction.of(relBuilder.getCluster(), function); + final RexNode call = relBuilder.getRexBuilder().makeCall(outputType, sqlFunction, args); + return LogicalTableFunctionScan.create( + relBuilder.getCluster(), + inputs, + call, + null, + outputType, + Collections.emptySet()); + } + private RelNode convertLegacyTableFunction( CorrelatedFunctionQueryOperation correlatedFunction, TableFunctionDefinition functionDefinition, diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/LateralSnapshotJoinSemanticTestPrograms.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/LateralSnapshotJoinSemanticTestPrograms.java index 275323610ec5d..0d0582668e4fe 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/LateralSnapshotJoinSemanticTestPrograms.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/LateralSnapshotJoinSemanticTestPrograms.java @@ -57,7 +57,7 @@ public class LateralSnapshotJoinSemanticTestPrograms { "load_completed_time => CAST(TIMESTAMP '2100-01-01 00:00:00' AS TIMESTAMP_LTZ(3))"; /** Event time of the flip-trigger row; equal to the {@link #MID_FLIP} timestamp. */ - private static final String FLIP_TRIGGER_TS = "00:00:10"; + static final String FLIP_TRIGGER_TS = "00:00:10"; /** A build-side key that never matches any probe row. */ private static final String FLIP_TRIGGER_KEY = "__flip_trigger__"; @@ -360,14 +360,14 @@ private static String join(String joinType, String projection, String flip, Stri + condition; } - private static List defaultProbe() { + static List defaultProbe() { return Arrays.asList( Row.of("a", 100, ts("00:01:00")), Row.of("b", 200, ts("00:01:01")), Row.of("c", 300, ts("00:01:02"))); } - private static List defaultBuild() { + static List defaultBuild() { return Arrays.asList( Row.of("a", 10, ts("00:00:01")), Row.of("b", 20, ts("00:00:02")), @@ -384,20 +384,20 @@ private static List manyProbes(int count) { * Appends a non-matching build row at the {@link #FLIP_TRIGGER_TS} timestamp; its watermark * flips the operator to the JOIN phase mid-stream (all real build rows are earlier). */ - private static List withFlipTrigger(List data) { + static List withFlipTrigger(List data) { final List withTrigger = new ArrayList<>(data); withTrigger.add(Row.of(FLIP_TRIGGER_KEY, 0, ts(FLIP_TRIGGER_TS))); return withTrigger; } - private static SourceTestStep probe(List data) { + static SourceTestStep probe(List data) { return SourceTestStep.newBuilder("probe") .addSchema(PROBE_SCHEMA) .producedValues(data.toArray(new Row[0])) .build(); } - private static SourceTestStep throttledProbe(List data, long sleepMillis) { + static SourceTestStep throttledProbe(List data, long sleepMillis) { return SourceTestStep.newBuilder("probe") .addSchema(PROBE_SCHEMA) .addOptions(throttleOptions(sleepMillis)) @@ -405,7 +405,7 @@ private static SourceTestStep throttledProbe(List data, long sleepMillis) { .build(); } - private static SourceTestStep appendBuild(List data) { + static SourceTestStep appendBuild(List data) { return SourceTestStep.newBuilder("b") .addSchema(BUILD_SCHEMA) .producedValues(data.toArray(new Row[0])) @@ -426,7 +426,7 @@ private static SourceTestStep throttledUpsertBuild(List data, long sleepMil .build(); } - private static SinkTestStep.Builder keyValueSink() { + static SinkTestStep.Builder keyValueSink() { return SinkTestStep.newBuilder("sink") .addSchema("pk STRING", "pv INT", "bk STRING", "bv INT") .testMaterializedData(); @@ -439,7 +439,7 @@ private static Map throttleOptions(long sleepMillis) { return options; } - private static LocalDateTime ts(String time) { + static LocalDateTime ts(String time) { return LocalDateTime.parse("2020-01-01T" + time); } } diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/LateralSnapshotJoinSemanticTests.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/LateralSnapshotJoinSemanticTests.java index 6536f34a59149..c167bec341a5c 100644 --- a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/LateralSnapshotJoinSemanticTests.java +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/LateralSnapshotJoinSemanticTests.java @@ -38,6 +38,10 @@ public List programs() { LateralSnapshotJoinSemanticTestPrograms.FLIP_AT_END, LateralSnapshotJoinSemanticTestPrograms.DEFAULT_COMPILE_TIME, LateralSnapshotJoinSemanticTestPrograms.LIVE_JOIN, - LateralSnapshotJoinSemanticTestPrograms.BUFFERED_THEN_DRAINED); + LateralSnapshotJoinSemanticTestPrograms.BUFFERED_THEN_DRAINED, + LateralSnapshotJoinTableApiSemanticTestPrograms.INNER_JOIN_TABLE_API, + LateralSnapshotJoinTableApiSemanticTestPrograms.LEFT_JOIN_TABLE_API, + LateralSnapshotJoinTableApiSemanticTestPrograms.TABLE_API_BUILD_SIDE, + LateralSnapshotJoinTableApiSemanticTestPrograms.DOWNSTREAM_WINDOW_TABLE_API); } } diff --git a/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/LateralSnapshotJoinTableApiSemanticTestPrograms.java b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/LateralSnapshotJoinTableApiSemanticTestPrograms.java new file mode 100644 index 0000000000000..8f606b3313916 --- /dev/null +++ b/flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/LateralSnapshotJoinTableApiSemanticTestPrograms.java @@ -0,0 +1,197 @@ +/* + * 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.nodes.exec.stream; + +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.Tumble; +import org.apache.flink.table.test.program.SinkTestStep; +import org.apache.flink.table.test.program.TableTestProgram; +import org.apache.flink.types.Row; + +import java.util.Arrays; + +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.planner.plan.nodes.exec.stream.LateralSnapshotJoinSemanticTestPrograms.FLIP_TRIGGER_TS; +import static org.apache.flink.table.planner.plan.nodes.exec.stream.LateralSnapshotJoinSemanticTestPrograms.appendBuild; +import static org.apache.flink.table.planner.plan.nodes.exec.stream.LateralSnapshotJoinSemanticTestPrograms.defaultBuild; +import static org.apache.flink.table.planner.plan.nodes.exec.stream.LateralSnapshotJoinSemanticTestPrograms.defaultProbe; +import static org.apache.flink.table.planner.plan.nodes.exec.stream.LateralSnapshotJoinSemanticTestPrograms.keyValueSink; +import static org.apache.flink.table.planner.plan.nodes.exec.stream.LateralSnapshotJoinSemanticTestPrograms.throttledProbe; +import static org.apache.flink.table.planner.plan.nodes.exec.stream.LateralSnapshotJoinSemanticTestPrograms.ts; +import static org.apache.flink.table.planner.plan.nodes.exec.stream.LateralSnapshotJoinSemanticTestPrograms.withFlipTrigger; + +/** + * Table API counterparts of {@link LateralSnapshotJoinSemanticTestPrograms}, reusing the same + * sources, sinks, and flip-trigger determinism tricks but expressing the {@code LATERAL SNAPSHOT} + * join through {@link Table#joinLateral} / {@link Table#leftOuterJoinLateral} instead of SQL. + */ +public class LateralSnapshotJoinTableApiSemanticTestPrograms { + + public static final TableTestProgram INNER_JOIN_TABLE_API = + TableTestProgram.of( + "lateral-snapshot-inner-join-table-api", + "LATERAL SNAPSHOT inner join expressed through the Table API") + .setupTableSource(throttledProbe(defaultProbe(), 40L)) + .setupTableSource(appendBuild(withFlipTrigger(defaultBuild()))) + .setupTableSink( + keyValueSink() + .consumedValues( + "+I[a, 100, a, 10]", + "+I[a, 100, a, 11]", + "+I[b, 200, b, 20]") + .build()) + .runTableApi( + env -> + env.from("probe") + .joinLateral( + call( + "SNAPSHOT", + env.from("b").asArgument("input"), + descriptor("bts").asArgument("on_time"), + loadCompletedTime(FLIP_TRIGGER_TS) + .asArgument( + "load_completed_time")), + $("pk").isEqual($("bk"))) + .select($("pk"), $("pv"), $("bk"), $("bv")), + "sink") + .build(); + + public static final TableTestProgram LEFT_JOIN_TABLE_API = + TableTestProgram.of( + "lateral-snapshot-left-join-table-api", + "LATERAL SNAPSHOT left join expressed through the Table API pads " + + "unmatched probe rows with null") + .setupTableSource(throttledProbe(defaultProbe(), 40L)) + .setupTableSource(appendBuild(withFlipTrigger(defaultBuild()))) + .setupTableSink( + keyValueSink() + .consumedValues( + "+I[a, 100, a, 10]", + "+I[a, 100, a, 11]", + "+I[b, 200, b, 20]", + "+I[c, 300, null, null]") + .build()) + .runTableApi( + env -> + env.from("probe") + .leftOuterJoinLateral( + call( + "SNAPSHOT", + env.from("b").asArgument("input"), + descriptor("bts").asArgument("on_time"), + loadCompletedTime(FLIP_TRIGGER_TS) + .asArgument( + "load_completed_time")), + $("pk").isEqual($("bk"))) + .select($("pk"), $("pv"), $("bk"), $("bv")), + "sink") + .build(); + + public static final TableTestProgram TABLE_API_BUILD_SIDE = + TableTestProgram.of( + "lateral-snapshot-table-api-build-side", + "the SNAPSHOT 'input' argument is itself a Table API expression " + + "(a filtered table), not a plain scan") + .setupTableSource(throttledProbe(defaultProbe(), 40L)) + .setupTableSource(appendBuild(withFlipTrigger(defaultBuild()))) + .setupTableSink( + keyValueSink() + .consumedValues("+I[a, 100, a, 11]", "+I[b, 200, b, 20]") + .build()) + .runTableApi( + env -> { + final Table filteredBuild = + env.from("b").filter($("bv").isGreater(10)); + return env.from("probe") + .joinLateral( + call( + "SNAPSHOT", + filteredBuild.asArgument("input"), + descriptor("bts").asArgument("on_time"), + loadCompletedTime(FLIP_TRIGGER_TS) + .asArgument("load_completed_time")), + $("pk").isEqual($("bk"))) + .select($("pk"), $("pv"), $("bk"), $("bv")); + }, + "sink") + .build(); + + public static final TableTestProgram DOWNSTREAM_WINDOW_TABLE_API = + TableTestProgram.of( + "lateral-snapshot-downstream-window-table-api", + "a TUMBLE window downstream of the Table API join groups by the " + + "probe-side rowtime, which the join preserves") + .setupTableSource( + throttledProbe( + Arrays.asList( + Row.of("a", 1, ts("00:01:00")), + Row.of("b", 2, ts("00:01:01")), + Row.of("c", 3, ts("00:01:02"))), + 40L)) + .setupTableSource( + appendBuild( + withFlipTrigger( + Arrays.asList( + Row.of("a", 10, ts("00:00:01")), + Row.of("b", 20, ts("00:00:02")), + Row.of("c", 30, ts("00:00:03")))))) + .setupTableSink( + SinkTestStep.newBuilder("sink") + .addSchema("wStart TIMESTAMP(3)", "pvSum INT", "bvSum INT") + .testMaterializedData() + .consumedValues("+I[2020-01-01T00:01, 6, 60]") + .build()) + .runTableApi( + env -> + env.from("probe") + .joinLateral( + call( + "SNAPSHOT", + env.from("b").asArgument("input"), + descriptor("bts").asArgument("on_time"), + loadCompletedTime(FLIP_TRIGGER_TS) + .asArgument( + "load_completed_time")), + $("pk").isEqual($("bk"))) + .select($("pts"), $("pv"), $("bv")) + .window( + Tumble.over(lit(1).minutes()) + .on($("pts")) + .as("w")) + .groupBy($("w")) + .select($("w").start(), $("pv").sum(), $("bv").sum()), + "sink") + .build(); + + /** + * Builds the {@code load_completed_time} argument as a Table API expression equivalent to the + * SQL {@code CAST(TIMESTAMP '2020-01-01