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

Chao Sun resolved SPARK-58378.
------------------------------
    Fix Version/s: 4.3.0
                   5.0.0
       Resolution: Fixed

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

> Fuse compatible approximate percentile sketches
> -----------------------------------------------
>
>                 Key: SPARK-58378
>                 URL: https://issues.apache.org/jira/browse/SPARK-58378
>             Project: Spark
>          Issue Type: Improvement
>          Components: SQL
>    Affects Versions: 5.0.0
>            Reporter: Chao Sun
>            Assignee: Chao Sun
>            Priority: Major
>              Labels: pull-request-available
>             Fix For: 4.3.0, 5.0.0
>
>
> h2. Problem
> Approximate percentile sketches summarize a distribution and can answer 
> multiple requested percentiles from the same state. Spark currently builds a 
> separate sketch for every scalar percentile expression, even when the input 
> and accuracy are identical.
> {code:sql}
> SELECT
>   service,
>   percentile_approx(latency_ms, 0.50, 10000) AS p50,
>   percentile_approx(latency_ms, 0.90, 10000) AS p90,
>   percentile_approx(latency_ms, 0.95, 10000) AS p95
> FROM request_metrics
> GROUP BY service;
> {code}
> This query allocates three percentile digests per group, inserts each latency 
> three times, and serializes and merges three partial digests. Spark already 
> supports the equivalent array form, but existing scalar SQL should not 
> require a manual rewrite.
> h2. Proposed change
> Add a Catalyst optimizer rule that fuses compatible scalar approximate 
> percentile aggregates into one array-valued aggregate and projects the 
> original scalar results. Preserve output names, ordering, values, and 
> expression identifiers.
> Only combine deterministic inputs with the same structural expression, 
> evaluated accuracy, filter, aggregate mode, and distinctness. Leave 
> array-valued and Structured Streaming aggregates unchanged so existing 
> checkpoint schemas remain compatible.
> h2. Validation
> Add Catalyst, end-to-end physical-plan, accuracy, ANSI, floating-point, null, 
> grouping, TIME, and streaming checkpoint recovery regressions.



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