Fabian Hueske created FLINK-40914:
-------------------------------------

             Summary: Support partial updates on LATERAL SNAPSHOT join build 
side
                 Key: FLINK-40914
                 URL: https://issues.apache.org/jira/browse/FLINK-40914
             Project: Flink
          Issue Type: Improvement
          Components: Table SQL / Runtime
            Reporter: Fabian Hueske
            Assignee: Fabian Hueske


LATERAL SNAPSHOT join maintains the build-side of the join in state.
Currently, the state is maintained as a MapState<Row, Long>, supporting any 
kind of tables as input. One drawback of this implementation is that 
retractions require full row matches, which is not needed if the input table 
has a primary key. 
Inputs that are read from upsert connectors with partial deletes (like Kafka) 
require an expensive ChangelogNormalize operator to derive full deletion rows.

We can significantly improve the performance of such jobs by directly 
supporting partial deletes in the LateralSnapshotJoin operator.

This requires:
 * Alternative state layout for build-side input with primary key
 * Abstraction of build-side updates and lookups
 * A switch (or new operator implementation) to determine the correct 
build-side mode
 * Planner changes to configure the switch (or chose the correct operator 
implementation)
 * Adjustment of LSJ logic in `StreamNonDeterministicUpdatePlanVisitor` (make 
it less restrictive if build-side has a PK)
 * Tests for the new mode / implementation



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

Reply via email to