Skip to content

feat: native MIN/MAX over strings and SortAggregate as a native hash aggregate - #12

Open
msafonov wants to merge 3 commits into
branch-1.1from
joom/1.1-native-string-minmax
Open

msafonov wants to merge 3 commits into
branch-1.1from
joom/1.1-native-string-minmax

Conversation

@msafonov

@msafonov msafonov commented Oct 10, 2026 •

Copy link
Copy Markdown
Collaborator

What

  • MIN/MAX accept StringType with the default UTF8_BINARY collation. The native min_udaf/max_udaf already handle Utf8 and compare UTF-8 bytes, which is Spark's binary order. Collated strings stay in Spark. max_by/min_by and window MIN/MAX are unchanged.
  • A SortAggregateExec becomes a CometHashAggregateExec built from an ObjectHashAggregateExec with the same fields:
    • Spark runs an ObjectHashAggregateExec correctly over any input order, so an aggregate reverted to Spark later (by the cost-based engine choice or the unsafe-partial-aggregate revert) stays correct.
    • A sort by exactly the grouping keys directly below the aggregate is dropped. Without such a sort, only order-insensitive aggregates (MIN, MAX, COUNT, SUM, AVG) are converted, so a sort-then-collect_list/first pattern keeps its meaning.
    • The ordering is restored the way EnsureRequirements would: when a consumer is visited, and before it is converted, each child whose required ordering is no longer met and that leads through unary operators to a converted aggregate gets a native sort on that edge. In gold_orders this happens for the sort-merge join that sits directly over the SortAggregate. A partial aggregate under a shuffle, a final aggregate read by a hash aggregate, a project or a write, and any consumer without an ordering requirement get no sort.
    • Aggregates Comet does not support keep the SortAggregate.
  • The cost-based engine choice prices the converted aggregate as an object hash aggregate (agg + aggObjectHash), which is also what it reverts to.
  • spark.comet.exec.sortAggregate.enabled (default true) turns the conversion off.

Tests

  • New tests in CometAggregateSuite:
    • MIN/MAX over strings with nulls, empty strings and non-ASCII text, including a pair whose UTF-8 and UTF-16 orders differ; grouped and ungrouped; AQE on and off.
    • The SortAggregate in the Spark plan becomes a native hash aggregate.
    • The sort-merge join over a converted aggregate still gets sorted input.
    • No sort is added when no consumer needs the ordering (a partial aggregate under a shuffle, a final aggregate read by the result).
    • An unsupported aggregate (max_by on strings) stays a SortAggregate.
    • An order-sensitive aggregate over input sorted beyond its keys stays in Spark.
    • A forced cost-based revert runs as a Spark ObjectHashAggregateExec.
  • New test in CostBasedEngineChoiceSuite: the converted aggregate is priced as Agg + AggObjectHash.
  • Expected changes in existing tests:
    • min_max.sql and first_last.sql (two queries) expected SortAggregate is not supported and now run natively; they are checked against Spark.
    • CostBasedEngineChoiceSuite and WideRowSortFallbackSuite use max(string) as their Spark SortAggregate fixture, so they set spark.comet.exec.sortAggregate.enabled=false suite-wide.
  • Results:
    • CometAggregateSuite, CostBasedEngineChoiceSuite, WideRowSortFallbackSuite, CometSqlFileTestSuite: 778 passed, 0 failed.
    • CometExecSuite, CometJoinSuite, CometWindowExecSuite, CometNativeShuffleSuite, CometShuffleSuite, DisableAQECometShuffleSuite, TPC-DS V1_4/V2_7 plan stability: 579 passed, 0 failed.
    • No plan-stability golden file changed: no approved TPC-DS plan has a SortAggregate.
    • spotless clean; test-compile OK for spark-3.4, 4.0, 4.1 and 4.2.

Real data (gold_orders incremental, data of 2026-10-09)

Run on msaf-spark-comet:joom-1.1-6 with the changed classes defined on the driver. Native code is unchanged.

  • Result vs vanilla: 20,159,446 rows. Every column is equal except the four ROW_NUMBER() columns, whose sets per (key, order time) tie group are identical.
  • Plan: SortAggregate 3 → 0, Spark Exchange 2 → 0, Spark Sort 4 → 1, Spark sort-merge joins 3 → 2.
  • Task time 4.20 → 4.15 h. The support_tickets block (0.62 h) did not get cheaper: its two joins stay in Spark because the orders side is a Spark Project, which is in Spark because of filter(discounts, x -> ...) (a higher-order function).

Synthetic aggregate

62.5M rows grouped by (date, device_id, user_id, ephemeral) with first(string), timestamp max/min and a sum, the shape of platform.load_fact_active_device, 2 repetitions:

  • Comet with SortAggregate: 243 / 215 task-s
  • this PR: 124 / 108 task-s
  • vanilla: 345 / 398 task-s

🤖 Generated with Claude Code

…aggregate

Spark plans MIN or MAX over a string as a SortAggregate, since the
aggregation buffer is not mutable, and Comet did not convert SortAggregate,
so such an aggregate and often the joins around it ran in Spark.

MIN and MAX now accept StringType with the default UTF8_BINARY collation.
The native min/max compare UTF-8 bytes, which is Spark's binary order;
collated strings stay in Spark.

A SortAggregateExec is converted to a CometHashAggregateExec built from an
ObjectHashAggregateExec with the same fields. Spark runs that operator over
any input order, so an aggregate reverted to Spark later (cost-based engine
choice, unsafe partial aggregates) stays correct, and the cost-based choice
prices it as an object hash aggregate. A sort by exactly the grouping keys
directly below the aggregate is dropped; without one only order-insensitive
aggregates (MIN, MAX, COUNT, SUM, AVG) are converted. A native sort above the
aggregate restores the output ordering a consumer such as a sort-merge join
may rely on, and is dropped again when a shuffle reads it directly.
Aggregates Comet does not support keep the SortAggregate.

spark.comet.exec.sortAggregate.enabled (default true) turns the conversion
off.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@github-actions github-actions Bot added enhancement New feature or request area:aggregation labels Oct 10, 2026
msafonov and others added 2 commits October 10, 2026 13:17
… needs it

The native sort above a SortAggregate converted to a hash aggregate was
added unconditionally and dropped only under a shuffle. It is now added the
way EnsureRequirements would: when a consumer is visited, before it is
converted, each child whose required ordering is no longer met and that
leads through unary operators to a converted SortAggregate gets a sort on
that edge. A partial aggregate under a shuffle, a final aggregate read by a
hash aggregate, a project or a write, and any consumer without an ordering
requirement get none.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Spark's sort is stable, so a SortAggregate over input sorted by its grouping
keys sees each group's rows in their input order, and FIRST, LAST, ANY_VALUE
and FIRST_VALUE return the first row of a group as read. The native hash
aggregate does not guarantee that order, so converting such a SortAggregate
broke the exact first/last results the fork promises on an ordered single
split (CometFirstLastBoolAggFuzzSuite, 240 mismatches on Spark 4.1 and 559 on
3.5 for first, any_value and first_value over strings).

A SortAggregate is now converted only when every aggregate function is
order-insensitive (MIN, MAX, COUNT, SUM, AVG), whatever its input.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:aggregation enhancement New feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant