[
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)