feat: support kurtosis aggregate - #4818
Conversation
Adds a Comet-owned `Kurtosis` UDAF and `CometKurtosis` serde so that Spark's excess-kurtosis aggregate runs natively. The native accumulator stores `[n, avg, m2, m3, m4]` `Float64` state to mirror Spark's `CentralMomentAgg` (momentOrder = 4) buffer, and the update/merge kernels are a direct port of Spark's `updateExpressionsDef` / `mergeExpressions` expressions. `nullOnDivideByZero` is threaded through so `spark.sql.legacy.statisticalAggregate` both branches behave the same as Spark (NULL vs NaN when `m2 == 0`). Window use falls back to Spark today: the window path doesn't wire the Comet aggregate for kurtosis. Captured as an `expect_fallback` case in the SQL test. Scaffolding produced by the `implement-comet-expression` skill; audit performed by the `audit-comet-expression` skill.
mbutrovich
left a comment
There was a problem hiding this comment.
First pass, thanks @andygrove!
| ]) | ||
| } | ||
|
|
||
| fn default_value(&self, _data_type: &DataType) -> Result<ScalarValue> { |
There was a problem hiding this comment.
default_value returns ScalarValue::Float64(None), which is what the AggregateUDFImpl trait default already produces for a Float64 return type via ScalarValue::try_from(data_type) (datafusion/expr/src/udaf.rs:828-830). This override can be deleted.
This is the same redundancy already flagged on mode.rs:116 in #4782. Worth fixing in both so the pattern does not propagate to the next aggregate.
| /// Online update for the first four central moments. Direct port of Spark's | ||
| /// `CentralMomentAgg.updateExpressionsDef` for `momentOrder = 4`. | ||
| #[inline] | ||
| fn kurtosis_update( |
There was a problem hiding this comment.
There is no GroupsAccumulator here, so grouped kurtosis runs through DataFusion's generic GroupsAccumulatorAdapter, one boxed Accumulator and a ScalarValue round trip per group per batch.
That is a defensible starting point, but it is worth stating explicitly why, because the neighbouring central-moment aggregates in this same crate do have one: VarianceGroupsAccumulator (variance.rs:268) keeps flat Vec<f64> state and StddevGroupsAccumulator (stddev.rs:206) reuses it. 4817 in this same batch also ships a vectorized grouped accumulator and benchmarks it. Please either add one or record in the PR description that grouped kurtosis is adapter-backed, so the gap is a decision rather than an oversight.
| } | ||
| } | ||
|
|
||
| /// Online update for the first four central moments. Direct port of Spark's |
There was a problem hiding this comment.
kurtosis_update and kurtosis_merge are free functions in this file, but welford.rs is exactly the module this crate uses to share moment recurrences across aggregates: it already holds variance_update, variance_merge, covariance_update, covariance_merge and finalize_moments (welford.rs:27-141), consumed by both variance.rs and covariance.rs.
The order-4 recurrence subsumes the order-2 one, and skewness is the obvious next sibling and needs [n, avg, m2, m3] from the same recurrence. Putting these two functions in welford.rs alongside the existing ones keeps all the moment math in one place and means skewness is a momentOrder = 3 evaluate on top of shared state rather than a third copy of the same algebra.
| -- becomes a plain query.) | ||
| -- ============================================================ | ||
|
|
||
| query expect_fallback(unsupported Spark aggregate function: skewness) |
There was a problem hiding this comment.
expect_fallback(unsupported Spark aggregate function: skewness) pins Comet's current lack of skewness support as expected behavior in the kurtosis fixture.
That is the wrong place for it. It couples an unrelated expression's status to this file, and when skewness lands, this assertion fails in a fixture whose name gives no hint why. The comment above at :161 already anticipates this. Either drop the query or move it to a fixture named for skewness.
Which issue does this PR close?
Closes #.
Rationale for this change
kurtosisis a standard SQL statistical aggregate and one of the last remainingCentralMomentAggsiblings that Comet didn't run natively (variance and stddev are already supported). Adding it lets queries using kurtosis stay in Comet's native path instead of falling back to Spark for the whole aggregate.What changes are included in this PR?
KurtosisUDAF (native/spark-expr/src/agg_funcs/kurtosis.rs) with a per-rowKurtosisAccumulator. Intermediate state is[n, avg, m2, m3, m4]Float64, mirroring Spark'sCentralMomentAgg(momentOrder = 4) wire format so a Spark-produced Partial and a Comet-produced Final can share bytes without conversion. Update and merge kernels are direct ports of Spark'supdateExpressionsDefandmergeExpressions(Meng 2015 recurrence).nullOnDivideByZerothrough the proto sospark.sql.legacy.statisticalAggregatebehaves the same on both engines: the default returnsNULLwhenm2 == 0, legacy mode returnsNaN.Kurtosisprotobuf message and wires it through the native planner.CometKurtosisinspark/src/main/scala/org/apache/comet/serde/aggregates.scalaand registers it inQueryPlanSerde.aggrSerdeMap.supportsMixedPartialFinalis left atfalseto match the conservative policy already used byVarianceandStddevin the same file.docs/source/contributor-guide/expression-audits/agg_funcs.md; flips the support status indocs/source/user-guide/latest/expressions.mdfrom planned to supported.kurtosis(x) OVER (...)) still falls back because the Comet window path doesn't wire kurtosis today. Captured as anexpect_fallbackcase rather than left as an implicit gap.Scaffolding produced by the
implement-comet-expressionskill; theaudit-comet-expressionskill drove the audit and produced the extra fallback and coverage-gap tests.How are these changes tested?
KurtosisAccumulatorcovering empty group, single-value divide-by-zero in bothnullOnDivideByZeromodes, both of Spark's ownExpressionDescriptionexamples (-0.7014368047529627and0.19432323191699075), and a two-partition state merge that reproduces the single-batch result.spark/src/test/resources/sql-tests/expressions/aggregate/kurtosis.sqlrunning underConfigMatrix: parquet.enable.dictionary=false,true. Covers Spark's documented examples, GROUP BY, global aggregate, empty table, integer/decimal/float/bigint inputs, literal argument,FILTER (WHERE ...), mixed with other aggregates, andexpect_fallbackcases forskewnessand window use ofkurtosis. Numerical-stress cases (NaN,Infinity,-Infinity,1e15magnitudes) usespark_answer_onlymode.kurtosis_legacy.sqlgated withConfig: spark.sql.legacy.statisticalAggregate=trueexercises theNaNpath for single-value and all-equal groups../mvnw test -Dsuites="org.apache.comet.CometSqlFileTestSuite kurtosis" -Dtest=nonepasses under both the default (Spark 3.5) and-Pspark-4.0profiles.cd native && cargo clippy --all-targets --workspace -- -D warningspasses.