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

Ganesha S updated SPARK-58611:
------------------------------
        Parent: SPARK-32883
    Issue Type: Sub-task  (was: Improvement)

> Left anti stream-stream join
> ----------------------------
>
>                 Key: SPARK-58611
>                 URL: https://issues.apache.org/jira/browse/SPARK-58611
>             Project: Spark
>          Issue Type: Sub-task
>          Components: Structured Streaming
>    Affects Versions: 4.3.0
>            Reporter: Ganesha S
>            Priority: Major
>              Labels: pull-request-available
>
> Stream-stream join currently supports Inner, LeftOuter, RightOuter, FullOuter 
> and LeftSemi. LeftAnti is rejected at analysis time:
> {noformat}
> LeftAnti joins with a streaming DataFrame/Dataset on the right are not 
> supported
> {noformat}
> This leaves the common "find left rows with no match on the right" pattern 
> without a native streaming implementation: impressions without clicks, orders 
> without shipments, sessions without conversion, alerts without 
> acknowledgement. Users work around it with NOT IN / NOT EXISTS rewrites or 
> hand-written transformWithState logic, both more expensive and easy to get 
> subtly wrong.
> LeftSemi was added under this umbrella in SPARK-32862; its complement was 
> never done. Note stream-static LEFT ANTI already works when only the left 
> side is streaming -- only a streaming right side is rejected -- so the gap is 
> specifically stream-stream.
> h2. Approach
> Unlike left semi, left anti cannot emit while joining. A semi match is 
> positively determined: the moment a left row matches, it can be emitted. "No 
> match exists" is only decidable once the watermark guarantees no future right 
> row can match. Left anti is therefore implemented on the eviction path, 
> reusing the existing left outer plumbing in 
> \{{StreamingSymmetricHashJoinExec}}:
> * nothing is emitted when a left row matches;
> * at left-side state eviction, rows whose \{{matched}} flag is false are 
> emitted as bare left rows (left outer emits them joined with nulls instead);
> * a matched left row stays in state carrying \{{matched = true}} so that it 
> is suppressed at eviction time, rather than being dropped from state early 
> the way left semi does;
> * a left row that fails the pre-join filter can never match, so it is emitted 
> immediately without being added to state.
> h2. Requirements
> * *Watermark on the right side + time constraints are mandatory*, same 
> contract as left outer. Without them the engine can never conclude "no match".
> * *Append output mode only.* Update mode would have to emit rows before the 
> watermark can rule out a future match, and such a row could be invalidated by 
> a later batch. This is the opposite of SPARK-56384, which allowed Update for 
> Inner/LeftSemi precisely because those have no early-firing concern.
> * *No new state format version.* The \{{matched}} flag already persisted by 
> state format v2/v3/v4 is exactly the signal left anti needs, so existing 
> checkpoints need no migration. Format v1 is already rejected for non-inner 
> joins.
> * *RightAnti is out of scope* -- it is not among the join types planned onto 
> this operator.
> h2. Trade-off
> Left anti buffers every left row until eviction, whereas left semi drops 
> matched left rows from state eagerly. This is inherent: a matched row must be 
> retained so that it can be _suppressed_ at eviction rather than emitted. 
> State size for left anti is therefore comparable to left outer, not to left 
> semi.



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

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to