Aggregate Expressions#
ApproximatePercentile#
The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:
The percentage argument must be a foldable literal.
The accuracy argument must be a foldable literal.
Only byte, short, int, long, float, and double input types are supported.
Average#
The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:
YearMonthIntervalType and DayTimeIntervalType inputs are not supported
CollectSet#
The following differences from Spark are always present and do not require any additional configuration:
On Spark 4.2+,
collect_setinputs are normalized with Spark’s recursive floating-point normalizer to match SPARK-57298. For array inputs containing floating-point values the normalizer producesArrayTransform, which Comet executes through the JVM codegen dispatcher instead of fully natively (scalar and struct inputs stay native). Whenspark.comet.exec.scalaUDF.codegen.enabled=false, those array cases fall back to Spark.
The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:
collect_setwithRESPECT NULLSfalls back to Spark, since the native implementation always drops null inputs.
First#
The following differences from Spark are always present and do not require any additional configuration:
This function is not deterministic. Results may not match Spark.
HllSketchAgg#
The following incompatibilities cause HllSketchAgg to fall back to Spark by default. Set spark.comet.expression.HllSketchAgg.allowIncompatible=true to enable Comet acceleration despite these differences.
Comet uses a Rust DataSketches port; HLL sketch bytes and estimates may differ slightly from Spark.
The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:
The lgConfigK argument must be a foldable literal.
The lgConfigK argument must be in the range [4, 21].
Only int, long, string, and binary input types are supported.
HllUnionAgg#
The following incompatibilities cause HllUnionAgg to fall back to Spark by default. Set spark.comet.expression.HllUnionAgg.allowIncompatible=true to enable Comet acceleration despite these differences.
Comet uses a Rust DataSketches port; HLL sketch bytes and estimates may differ slightly from Spark.
Errors surface as plain Comet execution errors rather than Spark’s SparkRuntimeException with condition HLL_UNION_DIFFERENT_LG_K / HLL_INVALID_INPUT_SKETCH_BUFFER (sqlState 22000), so the message and error class differ even though both engines fail.
An input sketch in the updatable HLL_4 form carrying auxiliary-map entries is rejected with an error, where Spark reads it: the bundled Rust decoder reads the compact auxiliary layout in both forms and would otherwise return a silently wrong estimate. Comet only ever writes HLL_8, so this affects sketch columns produced elsewhere.
The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:
The allowDifferentLgConfigK argument must be a foldable literal.
HyperLogLogPlusPlus#
The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:
Input type must be boolean, integral, floating-point, decimal, date/time, string, or binary
DecimalTypewith precision > 18 is not supportedCollated (non-UTF8_BINARY) strings are not supported
Last#
The following differences from Spark are always present and do not require any additional configuration:
This function is not deterministic. Results may not match Spark.
ListAgg#
The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:
BinaryTypeinputs are not supported.WITHIN GROUP (ORDER BY ...)is not supported.Non-literal delimiters are not supported.
Non-default string collations are not supported.
MaxBy#
The following differences from Spark are always present and do not require any additional configuration:
This function is non-deterministic when multiple rows share the maximum ordering value. Results may differ from Spark in that case.
The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:
The value and ordering must both be fixed-length types (boolean, integral, floating-point, decimal, date, or timestamp). A variable-length or nested type such as string, binary, or struct falls back to Spark.
MinBy#
The following differences from Spark are always present and do not require any additional configuration:
This function is non-deterministic when multiple rows share the minimum ordering value. Results may differ from Spark in that case.
The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:
The value and ordering must both be fixed-length types (boolean, integral, floating-point, decimal, date, or timestamp). A variable-length or nested type such as string, binary, or struct falls back to Spark.
Mode#
The following incompatibilities cause Mode to fall back to Spark by default. Set spark.comet.expression.Mode.allowIncompatible=true to enable Comet acceleration despite these differences.
mode breaks ties non-deterministically in Spark (the result depends on JVM hash-map iteration order); Comet returns the smallest of the tied values instead (https://github.com/apache/datafusion-comet/issues/3970)
Percentile#
The following cases are not supported by Comet and always fall back to Spark, regardless of any allowIncompatible setting:
An array of percentages is not supported.
The percentage argument must be a literal.
A frequency argument is not supported.
Descending order in
WITHIN GROUP (ORDER BY ... DESC)is not supported.Only numeric input types are supported.
RegrIntercept#
The following incompatibilities cause RegrIntercept to fall back to Spark by default. Set spark.comet.expression.RegrIntercept.allowIncompatible=true to enable Comet acceleration despite these differences.
Comet merges the partial aggregates of
regr_interceptin a different floating-point operation order from Spark. When a group’s rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423)
RegrR2#
The following incompatibilities cause RegrR2 to fall back to Spark by default. Set spark.comet.expression.RegrR2.allowIncompatible=true to enable Comet acceleration despite these differences.
Comet merges the partial aggregates of
regr_r2in a different floating-point operation order from Spark. When a group’s rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423)
RegrReplacement#
The following incompatibilities cause RegrReplacement to fall back to Spark by default. Set spark.comet.expression.RegrReplacement.allowIncompatible=true to enable Comet acceleration despite these differences.
Comet merges the partial aggregates of
regr_sxxandregr_syyin a different floating-point operation order from Spark. When a group’s rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423)
RegrSXY#
The following incompatibilities cause RegrSXY to fall back to Spark by default. Set spark.comet.expression.RegrSXY.allowIncompatible=true to enable Comet acceleration despite these differences.
Comet merges the partial aggregates of
regr_sxyin a different floating-point operation order from Spark. When a group’s rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423)
RegrSlope#
The following incompatibilities cause RegrSlope to fall back to Spark by default. Set spark.comet.expression.RegrSlope.allowIncompatible=true to enable Comet acceleration despite these differences.
Comet merges the partial aggregates of
regr_slopein a different floating-point operation order from Spark. When a group’s rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423)