[ 
https://issues.apache.org/jira/browse/FLINK-40681?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated FLINK-40681:
-----------------------------------
    Labels: pull-request-available  (was: )

> Streaming outer join drops the null-padded row after an inner-side update 
> followed by delete when UPDATE_BEFORE is not required
> -------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: FLINK-40681
>                 URL: https://issues.apache.org/jira/browse/FLINK-40681
>             Project: Flink
>          Issue Type: Bug
>          Components: Table SQL / Runtime
>    Affects Versions: 2.3.0, 2.2.1, 1.20.5, 2.4.0
>            Reporter: Martijn Visser
>            Priority: Critical
>              Labels: pull-request-available
>
> When the planner does not require UPDATE_BEFORE on the inputs of a regular 
> outer join (both inputs have an upsert key that contains the join key and the 
> sink's primary key is satisfied by the query's upsert key), the inner side 
> sends a bare UPDATE_AFTER for a changed row, and no ChangelogNormalize sits 
> in front of the join. {{StreamingJoinOperator}} treats every accumulate 
> message alike: it overwrites the row in the {{JoinKeyContainsUniqueKey}} 
> state but still increments {{numOfAssociations}} of every matching outer row 
> (line 274 on master). The later DELETE of that inner row finds the counter at 
> 2 instead of 1 (line 321) and never emits the null-padded outer row again, so 
> the outer row disappears from an upsert sink. A repeated INSERT for the same 
> key has the same effect, so the defect is "an accumulate that replaces a row 
> on a unique-key side is counted as a new association", not the UPDATE_AFTER 
> kind specifically. The planner's assumption that the operator can handle 
> upsert input once it can identify the record is not honored by the outer-join 
> bookkeeping.
> Reproducer with two upsert tables and an upsert sink, all keyed on {{k}}, 
> parallelism 1:
> {code:sql}
> CREATE TABLE left_t (k INT NOT NULL, v STRING, PRIMARY KEY (k) NOT ENFORCED) 
> WITH ('connector' = 'values', 'changelog-mode' = 'I,UA,D');
> CREATE TABLE right_t (k INT NOT NULL, w STRING, PRIMARY KEY (k) NOT ENFORCED) 
> WITH ('connector' = 'values', 'changelog-mode' = 'I,UA,D');
> CREATE TABLE sink (k INT NOT NULL, v STRING, w STRING, PRIMARY KEY (k) NOT 
> ENFORCED) WITH ('connector' = 'values', 'sink-insert-only' = 'false');
> INSERT INTO sink SELECT l.k, l.v, r.w FROM left_t l LEFT JOIN right_t r ON 
> l.k = r.k;
> {code}
> left_t receives {{+I(1, a)}}; right_t receives {{+I(1, x)}}, {{+U(1, y)}}, 
> {{-D(1, y)}}. The {{+U}} must arrive after the left row has been counted (if 
> the left row arrives after the {{+U}} it is stored with the right count and 
> the delete behaves); the attached ITCase enforces that by gating the right 
> source on the sink contents. The plan is
> {code}
> Join(joinType=[LeftOuterJoin], where=[=(k, k0)], 
> leftInputSpec=[JoinKeyContainsUniqueKey], 
> rightInputSpec=[JoinKeyContainsUniqueKey], changelogMode=[I,UA,D])
> :- Exchange(distribution=[hash[k]], changelogMode=[I,UA,D])
> :  +- TableSourceScan(table=[[default_catalog, default_database, left_t]], 
> fields=[k, v], changelogMode=[I,UA,D])
> +- Exchange(distribution=[hash[k]], changelogMode=[I,UA,D])
>    +- TableSourceScan(table=[[default_catalog, default_database, right_t]], 
> fields=[k, w], changelogMode=[I,UA,D])
> {code}
> The raw changelog the sink receives (parallelism 1, HEAP and ROCKSDB alike):
> {code}
> +I[1, a, null]
> -D[1, a, null]
> +I[1, a, x]
> +I[1, a, y]
> -D[1, a, y]
> {code}
> Expected is the same sequence followed by {{+I[1, a, null]}}. The 
> materialized sink content is therefore empty instead of {{[1, a, null]}}, and 
> a later delete of the left row reaches the sink as a DELETE for a key it no 
> longer holds.
> As a control, not a workaround: with the sink primary key changed to {{(k, 
> v)}} the planner requires UPDATE_BEFORE, adds a ChangelogNormalize in front 
> of each input, the join sees an explicit {{-U}} that decrements the counter, 
> and the result is correct.
> The same arithmetic is in {{JoinHelper}} (async state) and in 
> {{StreamingSemiAntiJoinOperator}} (an anti join never re-emits the left row 
> after {{+I, +U, -D}} on the right side). {{MiniBatchStreamingJoinOperator}} 
> has a second problem: {{JoinKeyContainsUniqueKeyBundle}} folds an 
> UPDATE_AFTER followed by a DELETE within one bundle into nothing while the 
> previous version of the row stays in state, so nothing is emitted at all.
> Three existing tests assert the wrong output: 
> {{StreamingJoinOperatorTest#testLeftOuterJoinWithStateRetentionDisabled}}, 
> {{StreamingSemiAntiJoinOperatorTest#testLeftSemiJoinWithDifferentStateRetentionTime}}
>  and {{#testLeftSemiJoinWithStateRetentionDisabled}}.
> I'll open a draft PR with failing tests for all variants (operator harness 
> for LEFT/RIGHT/FULL with sync and async state, duplicate INSERT, mini-batch, 
> anti join; ITCases asserting the raw changelog; plan tests pinning the 
> changelog modes)



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

Reply via email to