Repository navigation
Conversation
…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>
… 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>
This was referenced Oct 10, 2026
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What
MIN/MAXacceptStringTypewith the default UTF8_BINARY collation. The nativemin_udaf/max_udafalready handle Utf8 and compare UTF-8 bytes, which is Spark's binary order. Collated strings stay in Spark.max_by/min_byand window MIN/MAX are unchanged.SortAggregateExecbecomes aCometHashAggregateExecbuilt from anObjectHashAggregateExecwith the same fields:ObjectHashAggregateExeccorrectly 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.agg+aggObjectHash), which is also what it reverts to.spark.comet.exec.sortAggregate.enabled(default true) turns the conversion off.Tests
CometAggregateSuite:max_byon strings) stays a SortAggregate.ObjectHashAggregateExec.CostBasedEngineChoiceSuite: the converted aggregate is priced asAgg + AggObjectHash.min_max.sqlandfirst_last.sql(two queries) expectedSortAggregate is not supportedand now run natively; they are checked against Spark.CostBasedEngineChoiceSuiteandWideRowSortFallbackSuiteusemax(string)as their Spark SortAggregate fixture, so they setspark.comet.exec.sortAggregate.enabled=falsesuite-wide.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.
ROW_NUMBER()columns, whose sets per (key, order time) tie group are identical.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 ofplatform.load_fact_active_device, 2 repetitions:🤖 Generated with Claude Code