Floating-point Number Comparison#
Spark normalizes NaN and zero for floating point numbers for several cases. See NormalizeFloatingNumbers optimization rule in Spark.
However, one exception is comparison. Spark does not normalize NaN and zero when comparing values
because they are handled well in Spark (e.g., SQLOrderingUtil.compareFloats). But the comparison
functions of arrow-rs used by DataFusion do not normalize NaN and zero (e.g., arrow::compute::kernels::cmp::eq).
For top-level FLOAT and DOUBLE comparisons, Comet normalizes both operands before native
execution, including noncanonical NaN literals. Top-level IN, InSet, and NOT IN membership
also normalize dynamic candidates and lists containing NaN. When every candidate is a non-NaN
literal, Comet keeps DataFusion’s static filter and pruning path, enumerating both signed-zero
forms when a list contains zero.
This scalar membership handling does not yet recurse into floating-point leaves nested in arrays or structs; see #6019.
Nested equality and membership#
For arrays and structs containing FLOAT or DOUBLE, native =, <>, IN, and NOT IN
compare signed zeros as equal and all NaN representations as equal, matching Spark. This also
covers single-candidate membership that Spark rewrites into equality.
Equality and dynamic membership compare nested elements directly and stop at the first mismatch. Constant membership sets use normalized comparison values for static lookup. These operations preserve SQL null semantics and do not change the values returned by projections.
Ordering: NaN and signed zero (-0.0 vs +0.0)#
Spark’s ORDER BY, RANK, DENSE_RANK, and window frame comparisons route through
SQLOrderingUtil.compareDoubles / compareFloats, which equate all NaN representations and
define -0.0 == 0.0. NaN sorts above every non-NaN value.
For scalar FLOAT and DOUBLE keys, Comet normalizes NaNs and signed zeros before native
sorting, window peer comparisons, and WindowGroupLimitExec rank comparisons. Native range
partitioning normalizes its keys and sampled boundaries in the same way. Only comparison keys
are normalized; returned values retain their original NaN representations and zero signs.
Native sorting of floating-point values nested in arrays or structs still uses Arrow’s raw total ordering. Nested keys can therefore produce different ordering or rank results from Spark; see #5507.
Because those scalar comparison keys match Spark, spark.comet.exec.strictFloatingPoint=true no
longer forces a fallback for them: scalar FLOAT and DOUBLE sort keys, window and rank order
keys, and range partitioning keys all stay native under strict mode. Floating-point values nested
in arrays, structs, or maps still fall back under strict mode, because their ordering is the raw
total ordering described above.
Array distinct and union#
array_distinct and array_union fall back to Spark when their element type contains
FLOAT or DOUBLE, on every Spark version except 4.2.0. Spark 4.2.0 normalizes signed zeros
and NaNs in the arguments of these functions before they run (SPARK-54918), so native execution
returns the same results. Spark 3.4, 3.5, 4.0.0 to 4.0.4, and 4.1.0 to 4.1.3 keep positive and
negative zero distinct in flat arrays. Spark 4.0.5+, 4.1.4+, and 4.2.1+ normalize while these
functions evaluate instead (SPARK-59602), which native execution does not match for NaNs or for
zeros nested in arrays or structs. Other element types remain native.
The check is based on the element type, not the values. It also applies to NULL or empty
floating-point arrays and columns that never contain negative zero. The entire projection
falls back to Spark, introducing a CometColumnarToRow transition and moving unrelated
expressions in the same projection out of Comet. For example,
SELECT id + 1, array_distinct(a), i[0] + 5 evaluates all three expressions in a Spark Project.
This can have a substantial cost. A local Spark 4.1.3 benchmark of
sum(cardinality(array_distinct(d))) over two million array<double> rows found the default
projection fallback about 15 times slower than native opt-in (best of five runs).
The slowdown depends on the workload.
Setting spark.comet.expression.ArrayDistinct.allowIncompatible=true or
spark.comet.expression.ArrayUnion.allowIncompatible=true restores native execution on other
versions, but signed-zero and NaN results may differ from Spark. Native execution can keep
NaNs with different signs or payloads distinct. Signed-zero differences also depend on the
element type: native execution merges positive and negative zero in flat floating-point arrays,
but can keep them distinct inside nested arrays or structs. Only opt in if these differences
are acceptable for your data.
The check uses the Spark version number, not the changes a build contains. A build that reports
any version other than 4.2.0 falls back even if it includes SPARK-54918. A vendor build that
reports 4.2.0 but includes SPARK-59602 still runs natively; set
spark.comet.expression.ArrayDistinct.enabled=false and
spark.comet.expression.ArrayUnion.enabled=false on such a build.