ganeshashree opened a new pull request, #58197:
URL: https://github.com/apache/spark/pull/58197
<!--
Thanks for sending a pull request! Here are some tips for you:
1. If this is your first time, please read our contributor guidelines:
https://spark.apache.org/contributing.html
2. Ensure you have added or run the appropriate tests for your PR:
https://spark.apache.org/developer-tools.html
3. If the PR is unfinished, add '[WIP]' in your PR title, e.g.,
'[WIP][SPARK-XXXX] Your PR title ...'.
4. Be sure to keep the PR description updated to reflect all changes.
5. Please write your PR title to summarize what this PR proposes.
6. If possible, provide a concise example to reproduce the issue for a
faster review.
7. If you want to add a new configuration, please read the guideline first
for naming configurations in
'common/utils/src/main/scala/org/apache/spark/internal/config/ConfigEntry.scala'.
8. If you want to add or modify an error type or message, please read the
guideline first in
'common/utils/src/main/resources/error/README.md'.
-->
### What changes were proposed in this pull request?
<!--
Please clarify what changes you are proposing. The purpose of this section
is to outline the changes and how this PR fixes the issue.
If possible, please consider writing useful notes for better and faster
reviews in your PR. See the examples below.
1. If you refactor some codes with changing classes, showing the class
hierarchy will help reviewers.
2. If you fix some SQL features, you can provide some references of other
DBMSes.
3. If there is design documentation, please add the link.
4. If there is a discussion in the mailing list, please add the link.
-->
Aggregations whose grouping keys use non-binary (collated) strings currently
fall back to `SortAggregateExec`, because hash aggregation requires
binary-stable grouping keys (`Aggregate.supportsHashAggregate →
UnsafeRowUtils.isBinaryStable`).
Mirroring RewriteCollationJoin (SPARK-48000, which did the same for hash
joins), this PR adds an optimizer rule RewriteCollationAggregate that:
- injects CollationKey into non-binary-stable grouping keys, so grouping is
done on the collation-normalized (binary-stable) bytes; and
- preserves any projected original key by wrapping it in a First aggregate —
an arbitrary representative of each collation-equal group.
Physical planning then uses regular HashAggregateExec when the aggregation
buffer is mutable, or prefers ObjectHashAggregateExec over a sort when the
buffer is non-mutable (the projected-key First case) — guarded so it only
applies when all grouping keys are binary-stable.
The rule runs after the operators that lower DISTINCT / dropDuplicates /
INTERSECT / EXCEPT into aggregates, so those are covered too. For
multiple-distinct aggregates, RewriteDistinctAggregates replaces the grouping
expression with a fresh attribute, so CollationKey provenance is carried
forward as attribute metadata (registered in INTERNAL_METADATA_KEYS so it never
leaks into user-facing schemas); OptimizeExpand is guarded to not insert a
pre-aggregate grouping on collated keys.
Correctness never depends on the metadata/provenance signal — it only
selects object-hash vs. sort; collation-correct grouping comes from the
CollationKey normalization plus the binary-stable guard.
Gated by a new config spark.sql.collation.hashAggregation.enabled (default
true). Streaming aggregates and aggregates containing Python aggregate
functions are skipped.
Example:
-- 200M rows, UTF8_LCASE key: previously SortAggregate (sort + spill); now
hash-based, no sort
`SELECT k, COUNT(*) FROM t GROUP BY k; -- k STRING COLLATE UTF8_LCASE
`
### Why are the changes needed?
<!--
Please clarify why the changes are needed. For instance,
1. If you propose a new API, clarify the use case for a new API.
2. If you fix a bug, you can clarify why it is a bug.
-->
GROUP BY on a collated key is common, and the mandatory sort is a large,
avoidable cost — a reported query ran ~197s (with disk spill, a sort over 229M
rows) on a UTF8_LCASE key vs. ~39s when the key was cast to UTF8_BINARY. Hash
joins on collated keys were already fixed by SPARK-48000; aggregation was the
remaining gap.
### Does this PR introduce _any_ user-facing change?
<!--
Note that it means *any* user-facing change including all aspects such as
new features, bug fixes, or other behavior changes. Documentation-only updates
are not considered user-facing changes.
If yes, please clarify the previous behavior and the change this PR proposes
- provide the console output, description and/or an example to show the
behavior difference if possible.
If possible, please also clarify if this is a user-facing change compared to
the released Spark versions or within the unreleased branches such as master.
If no, write 'No'.
-->
Yes, a plan/performance change (results are unchanged): collated GROUP BY,
SELECT DISTINCT, dropDuplicates, INTERSECT/EXCEPT, and ROLLUP/CUBE now use
hash-based aggregation instead of sort-based. Compared to released versions,
this is new behavior in the unreleased branch. It can be disabled with
spark.sql.collation.hashAggregation.enabled=false. No change to query results
or output schemas.
### How was this patch tested?
<!--
If tests were added, say they were added here. Please make sure to add some
test cases that check the changes thoroughly including negative and positive
cases if possible.
If it was tested in a way different from regular unit tests, please clarify
how you tested step by step, ideally copy and paste-able, so that other
reviewers can test and check, and descendants can verify in the future.
If tests were not added, please describe why they were not added and/or why
it was difficult to add.
If benchmark tests were added, please run the benchmarks in GitHub Actions
for the consistent environment, and the instructions could accord to:
https://spark.apache.org/developer-tools.html#github-workflow-benchmarks.
-->
New tests in `CollationAggregationSuite` (21 cases): projected/unprojected
keys, struct keys, mixed unnormalizable (map) keys, single and multiple
distinct, SELECT DISTINCT, dropDuplicates, INTERSECT/EXCEPT, ROLLUP, the
OptimizeExpand guard, the streaming skip, the config kill-switch,
nullability/metadata preservation, and forced-object-hash correctness. Updated
CollationSuite, and added a PySpark test in test_pandas_udf_grouped_agg.py for
the Python-aggregate skip. Regression-checked CollationSuite,
DataFrameAggregateSuite, RewriteDistinctAggregatesSuite, OptimizeExpandSuite,
their query suites, and DataSourceV2RelationSuite - all green.
### Was this patch authored or co-authored using generative AI tooling?
<!--
If generative AI tooling has been used in the process of authoring this
patch, please include the
phrase: 'Generated-by: ' followed by the name of the tool and its version.
If no, write 'No'.
Please refer to the [ASF Generative Tooling
Guidance](https://www.apache.org/legal/generative-tooling.html) for details.
-->
Generated-by: Claude Code (Opus 4.8)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]