Gustavo de Morais created FLINK-40842:
-----------------------------------------

             Summary: Streaming semi/anti join over-counts matches for 
right-side updates without UPDATE_BEFORE
                 Key: FLINK-40842
                 URL: https://issues.apache.org/jira/browse/FLINK-40842
             Project: Flink
          Issue Type: Sub-task
          Components: Table SQL / API
            Reporter: Gustavo de Morais
            Assignee: Gustavo de Morais


{{StreamingSemiAntiJoinOperator}} keeps, for every left row, the number of 
matching right rows. An accumulate message on the right side always increments 
this count. When the right input is an upsert stream (UPDATE_AFTER without 
UPDATE_BEFORE), an update that replaces a right row with the same unique key is 
counted as a second match. A later delete then only decrements back to 1. This 
is the same bug as FLINK-40809, in the semi/anti join operator.

Example: {{SELECT * FROM L WHERE L.k IN (SELECT R.k FROM R)}} (semi) or {{NOT 
IN}} / {{NOT EXISTS}} (anti), where R is an upsert stream with PK {{{}(id, 
k){}}}.


{code:java}
+I L(k=1)
+I R(id=1, k=1, v=a)    count 1   semi: +I[L]    anti: -D[L]
+U R(id=1, k=1, v=b)    count 2   <- should stay 1
-D R(id=1, k=1)         count 1   semi: L is not retracted
                                  anti: L is not emitted again{code}
{{ }}
Expected: the semi join emits {{-D[L]}} and the anti join emits {{+I[L]}} after 
the delete. Actual: nothing is emitted, and L stays wrong.



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

Reply via email to