gustavodemorais commented on code in PR #29380:
URL: https://github.com/apache/flink/pull/29380#discussion_r4183708391
##########
flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/optimize/StreamNonDeterministicUpdatePlanVisitor.java:
##########
@@ -630,6 +631,10 @@ private StreamPhysicalRel visitJoin(
throwNonDeterministicConditionError(
Review Comment:
nit: LSJ still goes through the generic non-deterministic condition check.
That rejects a non-deterministic join condition whenever the build side updates
or the join is a LEFT join, even with insert-only inputs. For LSJ that's
stricter than needed: the output is append-only and the condition is only
evaluated when a probe row arrives, never when build rows are retracted. Fine
to keep it as is for now, but maybe worth a follow-up. A dedicated
`visitLateralSnapshotJoin` could relax this, and would also avoid the flag in
`visitJoinChild`.
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/join/LateralSnapshotJoinTest.java:
##########
@@ -259,6 +263,71 @@ void
testInnerJoinWithUpsertBuildSourceMaterializesRetractions() {
assertThat(plan).doesNotContain("DropUpdateBefore");
}
+ @Test
+ void testRejectNonDeterministicBuildColumnWithPrimaryKey() {
+ // The LSJ operator keys its state by the whole build row, so
retractions are matched by
+ // exact row equality. A non-deterministic build column therefore
breaks retraction even
+ // when the build side has a unique key (FLINK-40890). TRY_RESOLVE
must reject it.
+ enableTryResolve();
+ createUpdatingPkBuildSource();
+ util.tableEnv()
+ .executeSql(
+ "CREATE VIEW b_view AS "
+ + "SELECT bk, CAST(NOW() AS STRING) AS bn, bts
FROM b_pk");
+
+ final String sql =
+ "SELECT probe.pk, s.bk, s.bn FROM probe JOIN LATERAL SNAPSHOT("
+ + "input => TABLE b_view, on_time => DESCRIPTOR(bts), "
+ + "load_completed_time => CAST(TIMESTAMP '2026-07-01
00:00:00' AS TIMESTAMP_LTZ(3))"
+ + ") AS s ON probe.pk = s.bk";
+
+ assertThatThrownBy(() -> util.tableEnv().explainSql(sql))
+ .isInstanceOf(TableException.class)
+ .hasMessageContaining("can not satisfy the determinism
requirement");
+ }
+
+ @Test
+ void testAcceptDeterministicBuildColumnWithPrimaryKey() {
+ // Counterpart to the rejection test: a deterministic build column
over the same
+ // primary-keyed, updating source must still pass TRY_RESOLVE (no
over-rejection).
+ enableTryResolve();
+ createUpdatingPkBuildSource();
+ util.tableEnv()
+ .executeSql("CREATE VIEW b_view AS SELECT bk, UPPER(bk) AS bn,
bts FROM b_pk");
+
+ final String sql =
+ "SELECT probe.pk, s.bk, s.bn FROM probe JOIN LATERAL SNAPSHOT("
+ + "input => TABLE b_view, on_time => DESCRIPTOR(bts), "
+ + "load_completed_time => CAST(TIMESTAMP '2026-07-01
00:00:00' AS TIMESTAMP_LTZ(3))"
+ + ") AS s ON probe.pk = s.bk";
+
+ assertThatCode(() ->
util.tableEnv().explainSql(sql)).doesNotThrowAnyException();
+ }
+
+ private void enableTryResolve() {
Review Comment:
nit: could we move `enableTryResolve()` and `createUpdatingPkBuildSource()`
down to the other private helpers at the bottom of the file?
##########
flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/stream/sql/join/LateralSnapshotJoinTest.java:
##########
@@ -259,6 +263,71 @@ void
testInnerJoinWithUpsertBuildSourceMaterializesRetractions() {
assertThat(plan).doesNotContain("DropUpdateBefore");
}
+ @Test
+ void testRejectNonDeterministicBuildColumnWithPrimaryKey() {
Review Comment:
nit: what's being tested is whether the build side is updating, not the
primary key itself. Something like
`testRejectNonDeterministicBuildColumnWithUpdatingBuildSide` (same for the
accept test at line 290) might read clearer.
--
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]