[
https://issues.apache.org/jira/browse/SPARK-59788?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18119460#comment-18119460
]
Jonathan Gao commented on SPARK-59788:
--------------------------------------
Upstream PR: https://github.com/apache/spark/pull/59048
> Initialize the global aggregation once per empty batch in the streamline plan
> -----------------------------------------------------------------------------
>
> Key: SPARK-59788
> URL: https://issues.apache.org/jira/browse/SPARK-59788
> Project: Spark
> Issue Type: Bug
> Components: SQL, Structured Streaming
> Affects Versions: 5.0.0
> Reporter: Jonathan Gao
> Priority: Major
> Labels: pull-request-available
>
> A global streaming aggregation on the streamline (Real-Time Mode) plan drops
> the initialized result for an executed empty batch:
> * a source batch whose RDD has no partitions: the buffer is never
> initialized, no state is created, and no row is emitted;
> * a batch with multiple empty partitions: one seeded row per empty partition
> (duplicate rows);
> * a non-empty batch that has an empty sibling partition: an extra seeded row.
> The ordinary micro-batch plan initializes the global buffer exactly once and
> emits the initialized result (count 0, sum NULL) for an executed empty batch,
> so the streamline plan diverges from it.
> Fix: plan the non-final ProjectAggregationBufferExec with a single-partition
> (AllTuples) requirement for a global aggregation, so its empty-input guard
> fires exactly once per executed empty batch and never for a batch that has
> input. The projection propagates the single partition, so the plan keeps
> exactly one exchange; grouped aggregations are unchanged.
> Follow-up to SPARK-58635.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]