[ 
https://issues.apache.org/jira/browse/FLINK-28737?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Martijn Visser updated FLINK-28737:
-----------------------------------
    Description: 
With a source that emits UPDATE_BEFORE and a sink keyed on the upsert key,
{code}
INSERT INTO sink_with_pk SELECT a, b, c FROM cdc WHERE b > UNIX_TIMESTAMP() - 
300
{code}
passes {{TRY_RESOLVE}} and EXPLAIN PLAN_ADVICE gives no advice, because the 
Calc keeps the upsert key and the requirement is already empty when the 
condition is checked. The same query into a sink without primary key is 
rejected. For upsert sources the filter lands in ChangelogNormalize, which 
retracts from state, see FLINK-40898. FLINK-22826 and FLINK-35341 are user 
reports of it.

The fix checks the condition in {{StreamNonDeterministicUpdatePlanVisitor}} 
rather than dropping the upsert key in {{FlinkRelMdUpsertKeys}}: that would 
make such queries fail the ON CONFLICT validation in the default IGNORE mode, 
and the lookup-join rule only ever granted keys under deterministic conditions.

  was:Currently FlinkRelMdUpsertKeys does consider the case when the condition 
of a calc node is non-deterministic, this should be fixed


> NDU analyzer skips the Calc condition once the upsert key covers the sink key
> -----------------------------------------------------------------------------
>
>                 Key: FLINK-28737
>                 URL: https://issues.apache.org/jira/browse/FLINK-28737
>             Project: Flink
>          Issue Type: Bug
>          Components: Table SQL / Planner
>    Affects Versions: 2.1.0
>            Reporter: lincoln lee
>            Assignee: Martijn Visser
>            Priority: Critical
>              Labels: pull-request-available
>             Fix For: 2.4.0
>
>
> With a source that emits UPDATE_BEFORE and a sink keyed on the upsert key,
> {code}
> INSERT INTO sink_with_pk SELECT a, b, c FROM cdc WHERE b > UNIX_TIMESTAMP() - 
> 300
> {code}
> passes {{TRY_RESOLVE}} and EXPLAIN PLAN_ADVICE gives no advice, because the 
> Calc keeps the upsert key and the requirement is already empty when the 
> condition is checked. The same query into a sink without primary key is 
> rejected. For upsert sources the filter lands in ChangelogNormalize, which 
> retracts from state, see FLINK-40898. FLINK-22826 and FLINK-35341 are user 
> reports of it.
> The fix checks the condition in {{StreamNonDeterministicUpdatePlanVisitor}} 
> rather than dropping the upsert key in {{FlinkRelMdUpsertKeys}}: that would 
> make such queries fail the ON CONFLICT validation in the default IGNORE mode, 
> and the lookup-join rule only ever granted keys under deterministic 
> conditions.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to