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)

Reply via email to