andygrove opened a new issue, #6738: URL: https://github.com/apache/datafusion-comet/issues/6738
### Describe the bug `mode` counts values in an `OpenHashMap[AnyRef, Long]`. Before [SPARK-45599](https://github.com/apache/spark/pull/45036), Spark's `OpenHashSet` matched a key with `==` but hashed its `doubleToLongBits`, so two NaNs never match and every NaN row becomes a key of its own with a count of 1. That holds on every Spark 3.4 release and on 3.5.0 and 3.5.1. The fix, which matches keys with `equals`, is in 3.5.2 and 4.0.0. Comet's native `Mode` canonicalizes NaN on every Spark version (`normalize_key` in `native/spark-expr/src/agg_funcs/mode.rs`). So with `spark.comet.expression.Mode.allowIncompatible=true`, Comet can return a different value from Spark 3.4 even when there is no tie: | Values of `d` in one group | Spark 3.4.3 | Spark 4.1.3 | Comet on Spark 3.4.3 | | --- | --- | --- | --- | | `NaN, NaN, NaN, 1.0, 1.0` | `1.0` | `NaN` | `NaN` | `mode` is `Incompatible` today only because of how it breaks ties (#3970), so this difference isn't documented anywhere. Two comments also claim otherwise. The doc comment on `Mode` in `mode.rs` says Spark compares keys with `equals` and that "every supported version collapses `NaN` via `doubleToLongBits`". `CometMode.convert` says that versions before 4.2 key by `java.lang.Double.equals`. Both are true only from 3.5.2 and 4.0.0. ### Steps to reproduce With the `spark-3.4` profile: ```scala Seq(Double.NaN, Double.NaN, Double.NaN, 1.0, 1.0).toDF("d").coalesce(1).write.parquet(path) spark.read.parquet(path).createOrReplaceTempView("t") spark.sql("SELECT mode(d) FROM t").show() ``` Spark returns `1.0`. Comet, with `spark.comet.expression.Mode.allowIncompatible=true`, returns `NaN`. ### Expected behavior On a Spark version whose `OpenHashSet` predates SPARK-45599, Comet should either count each NaN as a value of its own, as Spark does, or not run `mode` on `FLOAT` or `DOUBLE` input natively and give that as the reason. The check needs the runtime version down to the patch (`Utils.majorMinorPatchVersion`, as `ArraySetSupport` in `arrays.scala` does) to cover 3.5.0 and 3.5.1. Signed zeros can't be matched exactly. Before the fix, `-0.0` and `0.0` match only when probing from one reaches the other in Spark's hash table. Keeping them apart, as Comet does now, matches the usual case. ### Additional context Found while correcting the floating-point contributor guide in #6447. Part of #6385, whose plan keeps Spark version differences such as SPARK-45599 in one Scala policy. `array_distinct`, `array_union`, `array_intersect` and `array_except` also sit on `OpenHashSet`, but Comet runs them natively on floating-point input only on Spark 4.2.0 (`ArraySetSupport`). `collect_set` uses a Scala `HashSet` instead, and that difference is #5312. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
