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)

Reply via email to