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)