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)