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

Gustavo de Morais updated FLINK-40841:
--------------------------------------
    Description: 
When an outer join's input is an upsert stream (updates arrive as UPDATE_AFTER 
without UPDATE_BEFORE) and the ON clause has a non-equi predicate, the join 
produces stale rows. The planner allows this because each input's upsert key 
contains its join key ({{{}FlinkChangelogModeInferenceProgram{}}}). It does not 
check the non-equi condition. An update then replaces the stored row without 
retracting what the old version joined with. For equi-only conditions this is 
safe, because the old and new versions always match the same rows.

Example: {{{}L LEFT JOIN R ON L.k = R.k AND R.v > L.v{}}}, both sides upsert.
 # The old version matched, the new one does not:

{code:java}
+I L(v=5)      -> +I[L, null]
+I R1(v=10)    -> -D[L, null], +I[L, R1]
+U R1(v=0)     -> nothing{code}
{{ }}
Expected {{{}[L, null]{}}}, actual {{{}[L, R1]{}}}. L's association count stays 
at 1, so L is never null-padded again.
 # The outer row is null-padded, then an update matches:

 
{code:java}
+I R(v=10)
+I L1(v=20)    -> +I[L1, null]
+U L1(v=5)     -> +I[L1', R] {code}
 

With UPDATE_BEFORE present, both cases are correct.

  was:
When an outer join's input is an upsert stream (updates arrive as UPDATE_AFTER 
without UPDATE_BEFORE) and the ON clause has a non-equi predicate, the join 
produces stale rows. The planner allows this because each input's upsert key 
contains its join key ({{{}FlinkChangelogModeInferenceProgram{}}}). It does not 
check the non-equi condition. An update then replaces the stored row without 
retracting what the old version joined with. For equi-only conditions this is 
safe, because the old and new versions always match the same rows.

Example: {{{}L LEFT JOIN R ON L.k = R.k AND R.v > L.v{}}}, both sides upsert.
 # The old version matched, the new one does not:

{code:java}
+I L(v=5) -> +I[L, null] +I R1(v=10) -> -D[L, null], +I[L, R1] +U R1(v=0) -> 
nothing{code}
{{ }}
Expected {{{}[L, null]{}}}, actual {{{}[L, R1]{}}}. L's association count stays 
at 1, so L is never null-padded again.
 # The outer row is null-padded, then an update matches:

 
{code:java}
+I R(v=10) +I L1(v=20) -> +I[L1, null] +U L1(v=5) -> +I[L1', R]
Expected [L1', R], actual [L1, null] and [L1', R]. -D[L1, null] is never sent.
{code}
 


With UPDATE_BEFORE present, both cases are correct.


> Streaming outer join with a non-equi condition produces wrong results on 
> upsert input
> -------------------------------------------------------------------------------------
>
>                 Key: FLINK-40841
>                 URL: https://issues.apache.org/jira/browse/FLINK-40841
>             Project: Flink
>          Issue Type: Sub-task
>          Components: Table SQL / API
>            Reporter: Gustavo de Morais
>            Assignee: Gustavo de Morais
>            Priority: Major
>
> When an outer join's input is an upsert stream (updates arrive as 
> UPDATE_AFTER without UPDATE_BEFORE) and the ON clause has a non-equi 
> predicate, the join produces stale rows. The planner allows this because each 
> input's upsert key contains its join key 
> ({{{}FlinkChangelogModeInferenceProgram{}}}). It does not check the non-equi 
> condition. An update then replaces the stored row without retracting what the 
> old version joined with. For equi-only conditions this is safe, because the 
> old and new versions always match the same rows.
> Example: {{{}L LEFT JOIN R ON L.k = R.k AND R.v > L.v{}}}, both sides upsert.
>  # The old version matched, the new one does not:
> {code:java}
> +I L(v=5)      -> +I[L, null]
> +I R1(v=10)    -> -D[L, null], +I[L, R1]
> +U R1(v=0)     -> nothing{code}
> {{ }}
> Expected {{{}[L, null]{}}}, actual {{{}[L, R1]{}}}. L's association count 
> stays at 1, so L is never null-padded again.
>  # The outer row is null-padded, then an update matches:
>  
> {code:java}
> +I R(v=10)
> +I L1(v=20)    -> +I[L1, null]
> +U L1(v=5)     -> +I[L1', R] {code}
>  
> With UPDATE_BEFORE present, both cases are correct.



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

Reply via email to