Skip to content

array_distinct and array_union diverge from Spark on -0.0 for Spark versions without SPARK-54918 #5701

Description

@andygrove

Describe the bug

array_distinct and array_union run natively by default (ArrayDistinct via a plain CometScalarFunction, ArrayUnion via CometArrayUnion; neither overrides getSupportLevel). DataFusion collapses -0.0 and 0.0, so on a Spark version that does not normalize signed zeros in array functions, Comet silently returns a different answer:

SELECT array_distinct(array(0.0, double('-0.0'), 1.0))
Spark:  [0.0, -0.0, 1.0]
Comet:  [0.0, 1.0]

The plan is CometProject with no fallback. NormalizeFloatingNumbers only rewrites grouping keys, join keys, window partition specs and Distinct, so a plain SELECT diverges for literals as well as column-sourced values.

SPARK-54918 makes Spark normalize signed zeros in array_distinct, array_union, array_intersect, array_except and arrays_overlap, which makes [0.0, 1.0] the correct answer. Its fix versions are 4.2.0, 4.1.4 and 4.0.5. Comet also supports 3.4 and 3.5, which will never get it, and any 4.0.x before 4.0.5 or 4.1.x before 4.1.4.

array_intersect and array_except already report Incompatible, so they only reach the native path under allowIncompatible. array_distinct and array_union do not.

Steps to reproduce

Run the query above on Spark 3.5, or on 4.1.3 (the version the build currently pins).

Expected behavior

On a Spark version without SPARK-54918, either match Spark or fall back. On a version with it, keep the native path.

Additional context

The split is at the patch level, so a compile-time shim keyed on shims.minorVerSrc is not sufficient: the build pins 4.0.4 and 4.1.3, but Comet runs against any patch release of those lines. It needs a runtime check.

Note that CometSparkSessionExtensions.isSparkNNPlus compares org.apache.spark.SPARK_VERSION as a string, so SPARK_VERSION >= "4.0.5" would be wrong once 4.0.10 ships. A numeric comparison helper is needed.

Suggested shape: a getSupportLevel override on the ArrayDistinct and ArrayUnion serdes returning Incompatible when the element type contains FloatType/DoubleType and the running Spark version predates the fix, in the shape CometArrayExcept already uses for its unsupported element types. That keeps the fast path for every other element type and every fixed Spark version.

The signed-zero cases in array_distinct.sql, array_union.sql, array_except.sql and array_intersect.sql are currently query ignore(...); they should be restored once behavior is version-correct.

Found while reviewing #5262.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions