andygrove commented on PR #5750:
URL: 
https://github.com/apache/datafusion-comet/pull/5750#issuecomment-5959868196

   The gate question from my last review is still open. 
`ArraySetSupport.normalizesSignedZero` returns `true` for 4.0.5, 4.1.4 and 
4.2.1, but on those releases SPARK-59602 normalizes inside `ArraySetLike` 
instead of in the plan, so Comet receives the raw input. I re-checked the 
release candidates today: `v4.0.5-rc1`, `v4.1.4-rc2` and `v4.2.1-rc1` all carry 
it, and `v4.2.0` is the only release with the plan rewrite.
   
   Rather than narrowing the gate to 4.2.0, I prototyped normalizing in the 
native path, which keeps your version gate as it is and makes it correct on 
those releases too. The native side is a thin wrapper that runs 
`normalize_nested_floats` from `float_semantics` on each argument and then 
delegates to DataFusion, the same way `SparkArrayExtrema` wraps 
`array_min_udf()`:
   
   ```rust
   /// Spark's `array_distinct` and `array_union` for elements that hold a 
float at any depth.
   ///
   /// From SPARK-54918, Spark treats `-0.0` and `0.0`, and every NaN 
representation, as one value at
   /// any depth, and returns the normalized value. DataFusion folds `-0.0` 
into `0.0` only in a flat
   /// float array and compares NaNs by their bits, so normalize the input 
before delegating.
   #[derive(Debug, Hash, Eq, PartialEq)]
   pub struct SparkArraySetOp {
       name: &'static str,
       datafusion_udf: Arc<ScalarUDF>,
   }
   
   impl SparkArraySetOp {
       pub fn distinct() -> Self {
           Self {
               name: "spark_array_distinct",
               datafusion_udf: array_distinct_udf(),
           }
       }
   
       pub fn union() -> Self {
           Self {
               name: "spark_array_union",
               datafusion_udf: array_union_udf(),
           }
       }
   }
   
   impl ScalarUDFImpl for SparkArraySetOp {
       fn name(&self) -> &str {
           self.name
       }
   
       fn signature(&self) -> &Signature {
           self.datafusion_udf.signature()
       }
   
       fn return_type(&self, arg_types: &[DataType]) -> Result<DataType> {
           self.datafusion_udf.return_type(arg_types)
       }
   
       fn invoke_with_args(&self, mut args: ScalarFunctionArgs) -> 
Result<ColumnarValue> {
           args.args = args
               .args
               .into_iter()
               .map(|arg| match arg {
                   ColumnarValue::Array(array) => {
                       Ok(ColumnarValue::Array(normalize_nested_floats(&array)))
                   }
                   ColumnarValue::Scalar(value) => {
                       let array = normalize_nested_floats(&value.to_array()?);
                       Ok(ColumnarValue::Scalar(ScalarValue::try_from_array(
                           &array, 0,
                       )?))
                   }
               })
               .collect::<Result<_>>()?;
           self.datafusion_udf.invoke_with_args(args)
       }
   }
   ```
   
   It lives in `native/spark-expr/src/array_funcs/array_set_ops.rs` and is 
registered next to `SparkArrayRemove` in `comet_scalar_funcs.rs`:
   
   ```rust
   Arc::new(ScalarUDF::new_from_impl(SparkArraySetOp::distinct())),
   Arc::new(ScalarUDF::new_from_impl(SparkArraySetOp::union())),
   ```
   
   On the Scala side, the two serdes pick those names for float element types, 
the way `CometArrayRemove` picks `spark_array_remove` since #6518. Every other 
element type still goes straight to DataFusion:
   
   ```scala
     // DataFusion folds -0.0 into 0.0 only in a flat float array and compares 
NaNs by their bits.
     // The spark_ variants normalize floats at any depth first, as Spark does 
from SPARK-54918.
     def function(name: String, dataType: DataType): String =
       if (SupportLevel.containsType(dataType, classOf[FloatType], 
classOf[DoubleType])) {
         s"spark_$name"
       } else {
         name
       }
   ```
   
   ```scala
   object CometArrayDistinct extends CometExpressionSerde[ArrayDistinct] {
     // getIncompatibleReasons and getSupportLevel unchanged
   
     override def convert(
         expr: ArrayDistinct,
         inputs: Seq[Attribute],
         binding: Boolean): Option[ExprOuterClass.Expr] = {
       val childProto = exprToProtoInternal(expr.child, inputs, binding)
       scalarFunctionExprToProto(
         ArraySetSupport.function("array_distinct", expr.dataType),
         childProto)
     }
   }
   ```
   
   `CometArrayUnion.convert` changes the same way, passing 
`ArraySetSupport.function("array_union", expr.dataType)` instead of 
`"array_union"`.
   
   For the test, I extended `array set noncanonical NaN normalization` with 
nested cases and an opt-in arm that runs on every version. That arm is what 
lets CI see the bug. 4.2.0 rewrites the plan to normalize the input first, and 
the SPARK-59602 releases aren't in the matrix yet. Every Spark version returns 
1 for all of these, because older releases canonicalize a flat NaN and compare 
nested floats with the SQL ordering. On 4.1.3 the current head returns 2 for 
them under the opt-in:
   
   ```scala
           sql("SELECT float('NaN') AS f, double('NaN') AS d, 0.0D AS z").write
             .parquet(dir + "/data")
           spark.read.parquet(dir + 
"/data").createOrReplaceTempView("array_set_nan")
           // Negate scanned values because Parquet canonicalizes NaNs on 
write. Every Spark version
           // merges these elements: older ones canonicalize a flat NaN and 
compare nested floats with
           // the SQL ordering, so the native path must normalize them even 
under the opt-in.
           val expressions = Seq("f", "d").flatMap { column =>
             Seq(
               s"array_distinct(array($column, -$column))",
               s"array_union(array($column), array(-$column))")
           } ++ Seq(
             "array_distinct(array(array(z), array(-z)))",
             "array_distinct(array(array(d), array(-d)))",
             "array_distinct(array(named_struct('x', z), named_struct('x', 
-z)))",
             "array_union(array(named_struct('x', d)), array(named_struct('x', 
-d)))")
           expressions.foreach { expression =>
             val query = s"SELECT size($expression) FROM array_set_nan"
             if 
(ArraySetSupport.normalizesSignedZero(org.apache.spark.SPARK_VERSION)) {
               checkSparkAnswerAndOperator(query)
             } else {
               checkSparkAnswerAndFallbackReason(query, "SPARK-54918")
             }
             withSQLConf(
               CometConf.getExprAllowIncompatConfigKey(classOf[ArrayDistinct]) 
-> "true",
               CometConf.getExprAllowIncompatConfigKey(classOf[ArrayUnion]) -> 
"true") {
               checkSparkAnswerAndOperator(query)
             }
           }
   ```
   
   The flat `array(z, -z)` case is left out of the opt-in arm on purpose: older 
Spark keeps those zeros apart, and that difference is what the gate is for. I 
also added four Rust unit tests on `SparkArraySetOp`, covering NaN payloads and 
nulls in a flat array, union order across both sides, nested lists and structs, 
and a scalar argument. All four fail with the normalization removed.
   
   With the normalization in place, the opt-in caveat narrows to signed zeros. 
In the prototype I changed the reason string to end in "for matching 
signed-zero semantics". I rewrote the opt-in paragraph in `floating-point.md` 
to say that native execution treats `-0.0` and `0.0` as one value and returns 
`0.0`, while older Spark keeps both zeros in a flat array and returns whichever 
came first inside nested arrays and structs. I dropped "and NaN" from the two 
`expressions.md` notes. I also replaced the `KnownFloatingPointNormalized` 
comment on `normalizesSignedZero`, since the newer releases don't emit a marker 
at all. With all of that on top of a merge of main, the array SQL fixtures, the 
float sweep's array cases and the `array_distinct` fuzz test pass on 4.1.3 and 
4.2.0.
   
   Main moved again with #6518, but the conflicts are mechanical: in 
`CometFloatSemanticsSuite` both adjacent `KnownGap` entries go, and in 
`floating-point.md` both appended sections stay. Once you push I'll add 
`run-all-spark-profiles`, since the PR tier only runs 4.1 and this behaves 
differently on each profile. Would you be up for taking this route? If it's 
easier, I can push the commit to your branch and you can take it from there.
   


-- 
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]

Reply via email to