Gustavo de Morais created FLINK-40875:
-----------------------------------------
Summary: Optimization to let streaming joins fill key-only deletes
instead of adding a ChangelogNormalize
Key: FLINK-40875
URL: https://issues.apache.org/jira/browse/FLINK-40875
Project: Flink
Issue Type: Improvement
Components: Table SQL / API
Reporter: Gustavo de Morais
When an upsert source sends key-only deletes and a downstream consumer requires
full deletes, the planner adds a ChangelogNormalize. Its only job is to keep
the full row per key in state and fill in the deletes. If the source feeds a
streaming join, the join already stores the same full rows in its own state. It
can look up the full record of a key-only delete, which it already does for
non-equi conditions (FLINK-40858).
Proposal:
* In the delete-kind inference, let a join input with key-only deletes stay
key-only even when the join's parent requires full deletes, and mark the join's
output as full deletes.
* Persist the join's output changelog mode in the compiled plan. The operator
looks up the full record for every key-only delete when its output must contain
full deletes.
* Lookups stay limited to where they are needed: one point read per key-only
delete, and none otherwise.
Benefit: no ChangelogNormalize between the source and the join, which saves one
operator and a full copy of the source state per key.
To consider:
* Records missing from state (expired by TTL) cannot be filled. Define the
behavior and align it with ChangelogNormalize.
* It only applies when downstream does not require UPDATE_BEFORE, which still
needs ChangelogNormalize.
* Plans of new statements change, which affects upgrades that re-plan SQL from
a savepoint. Consider a config option.
* Mini-batch, async-state, semi/anti and multi-way joins must support it or be
excluded.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)