Gustavo de Morais created FLINK-40939:
-----------------------------------------
Summary: TRY_RESOLVE accepts non-deterministic join keys and join
conditions on updating inputs with a unique key
Key: FLINK-40939
URL: https://issues.apache.org/jira/browse/FLINK-40939
Project: Flink
Issue Type: Bug
Components: Table SQL / Planner
Affects Versions: 2.2.1, 2.3.0, 1.20.5
Reporter: Gustavo de Morais
With {{table.optimizer.non-deterministic-update.strategy}} set to
{{TRY_RESOLVE}}, {{StreamNonDeterministicUpdatePlanVisitor}} does not require
the join key or the other columns of the join condition to be deterministic
when an updating join input has a unique key ({{HasUniqueKey}} or
{{JoinKeyContainsUniqueKey}}). Only the columns required by downstream are
passed to the input. The same applies to {{StreamPhysicalMultiJoin}}.
The join input is hash-partitioned by the join key and the join state is keyed
by it. A retraction also evaluates the join condition on the incoming row to
find the partners to retract. If an update recomputes a join key
non-deterministically, the -U is shuffled to a different key and, with
parallelism > 1, usually to a different subtask than the row it should retract.
If it recomputes a non-equi condition column, it matches different partners. In
both cases the stored row is never retracted and the result keeps stale rows.
The planner accepts such queries silently, also when the non-deterministic call
is written directly in the equi condition (it is pushed down into a Calc).
{code:sql}
SET 'table.optimizer.non-deterministic-update.strategy' = 'TRY_RESOLVE';
CREATE TABLE orders (order_id INT, status STRING, PRIMARY KEY (order_id) NOT
ENFORCED)
WITH ('connector' = 'values', 'changelog-mode' = 'I,UA,D');
CREATE TABLE promos (promo_day STRING, discount INT) WITH ('connector' =
'values');
INSERT INTO snk
SELECT o.order_id, o.status, p.discount
FROM (SELECT order_id, status, DATE_FORMAT(NOW(), 'yyyy-MM-dd') AS order_day
FROM orders) o
JOIN promos p ON o.order_day = p.promo_day;
{code}
If +I[1, CREATED] is processed on 2026-10-05 and +U[1, SHIPPED] on 2026-10-06,
the -U is routed to key 2026-10-06 and never retracts the row stored under
2026-10-05. The sink ends with [1, CREATED, 10] and [1, SHIPPED, 20] although
orders has a single row. Expected: TRY_RESOLVE rejects the query with "can not
satisfy the determinism requirement".
We should require all columns referenced by the join condition of regular joins
and multi-joins to be deterministic on updating inputs. Temporal joins are
apparently not affected since the versioned side is never retracted by value.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)