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

Jungtaek Lim resolved SPARK-59290.
----------------------------------
    Fix Version/s: 4.4.0
       Resolution: Fixed

Issue resolved by pull request 58562
[https://github.com/apache/spark/pull/58562]

> The operation "stream" EXCEPT "static" is allowed but does not work as 
> expected
> -------------------------------------------------------------------------------
>
>                 Key: SPARK-59290
>                 URL: https://issues.apache.org/jira/browse/SPARK-59290
>             Project: Spark
>          Issue Type: Bug
>          Components: Structured Streaming
>    Affects Versions: 4.4.0
>            Reporter: Jungtaek Lim
>            Assignee: Jungtaek Lim
>            Priority: Critical
>              Labels: pull-request-available
>             Fix For: 4.4.0
>
>
> DataFrame has several "set" operations e.g. UNION, INTERSECT, EXCEPT (MINUS) 
> which performs set operations between two DataFrames.
> But the allowance of streaming vs batch (static) DataFrame for left and right 
> is different depending on the operator.
>  * UNION: batch-batch and stream-stream is allowed, mix of stream and batch 
> isn't allowed.
>  * INTERSECT: stream isn't allowed, hence only batch-batch is allowed.
>  * EXCEPT: stream on the right side isn't allowed, so batch-batch and 
> stream-batch are allowed.
> We just found out the allowed case of EXCEPT, stream-batch, is very poorly 
> defined and it does not work as expected in most cases.
> Here is the semantic of EXCEPT DISTINCT vs EXCEPT ALL.
>   - EXCEPT DISTINCT returns values present in L but absent from R, with 
> duplicates removed.
>   - EXCEPT ALL preserves multiplicity. For each value, it returns 
> max(count(L) - count(R), 0) copies.
>   Example:
>   L = [A, A, A, B]
>   R = [A, A, C]
>   Results:
>   EXCEPT DISTINCT = [B]
>   EXCEPT ALL      = [A, B]
> Currently Spark rewrites the operator EXCEPT with AGGREGATE, which works for 
> batch query, but does not work for streaming query as following:
> > append mode
> This requires watermark to be set on the stream side. It doesn't make much 
> sense in practice since this means the data on the stream has event time 
> column which is expected to increase over time, while the data on the batch 
> side is only evaluated once (in practice the batch side can change, but if 
> this is taken into account, it is even harder to reason about the behavior).
> > update mode
> Stateful aggregate no longer asks to set the watermark on the stream side. 
> But that is even worse since it can lead to the correctness issue.
>  * EXCEPT ALL: static side is evaluated "per batch" which is problematic for 
> counting the number of the distinct row from stream - batch (batch side is 
> multiplied).
>  * EXCEPT DISTINCT: The expected behavior is deduplication among the distinct 
> row, but the actual behavior is, the operator produces the row every batch 
> when there is an update on state store. That said, deduplication across 
> microbatches does not happen.
> There "might" be a way to implement the behavior correctly (most likely to 
> have stateful implementation for EXCEPT with streaming), but even with that, 
> it is uneasy to reason about the behavior accounting the fact that batch side 
> can change over time; the output being produced in batch M cannot be 
> corrected in batch M+1 even if the batch side in batch M+1 changes it.
> Unless there is a clear definition of EXCEPT with streaming, we should 
> disallow having stream in any side in EXCEPT. Notice we already have a 
> correctness issue so it's not an option to defer the fix till we find the 
> clear definition.



--
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