sunchao commented on code in PR #6475:
URL: https://github.com/apache/datafusion-comet/pull/6475#discussion_r4161824362
##########
spark/src/main/scala/org/apache/comet/serde/CometSortOrder.scala:
##########
@@ -20,38 +20,64 @@
package org.apache.comet.serde
import org.apache.spark.sql.catalyst.expressions.{Ascending, Attribute,
Descending, NullsFirst, NullsLast, SortOrder}
-import org.apache.spark.sql.types.{DoubleType, FloatType}
+import org.apache.spark.sql.types.{ArrayType, DataType, StructType}
import org.apache.comet.CometConf
import org.apache.comet.serde.QueryPlanSerde.exprToProtoInternal
+/**
+ * The key of a Sort, TopK, Window, WindowGroupLimit or range partitioning.
+ *
+ * Arrow orders floats by IEEE 754 total order, so the native side normalizes
a `FLOAT` or
+ * `DOUBLE` key, and an array or struct key with a float at any depth, before
comparing: NaN
+ * payloads fold together and signed zeros tie, as in Spark's
`SQLOrderingUtil`. Only the
+ * comparison key is normalized; returned values keep their original NaN
representation and zero
+ * sign. So every key is compatible, in strict floating-point mode too, except
one that has a
+ * float and whose type can hold a null element or field.
+ *
+ * Spark orders a null element or field below every other value, whatever the
key's null order.
+ * The native sort takes its place from the key's null order, and a `RANGE`
window frame orders it
+ * above every other value, so those keys can sort or frame differently from
Spark, with floats or
+ * without. Strict mode keeps the fallback it applied to every nested float
key before they were
+ * normalized, for the keys that can hit this.
+ * https://github.com/apache/datafusion-comet/issues/6476
+ * https://github.com/apache/datafusion-comet/issues/6477
+ */
object CometSortOrder extends CometExpressionSerde[SortOrder] {
/**
* Subject of both the runtime fallback reason and the generated
compatibility docs. Shared so
* the two cannot describe the policy differently.
*/
- private val nestedFloatingPointSort =
- "Sorting on floating-point values nested in arrays, structs, or maps"
+ private val nullableNestedFloatingPointSort =
+ "Sorting on floating-point values nested in an array or struct that can
hold a null " +
+ "element or field"
override def getIncompatibleReasons(): Seq[String] = Seq(
- s"$nestedFloatingPointSort is not 100% compatible with Spark when " +
+ s"$nullableNestedFloatingPointSort is not 100% compatible with Spark when
" +
s"`${CometConf.COMET_EXEC_STRICT_FLOATING_POINT.key}=true`")
- override def getSupportLevel(expr: SortOrder): SupportLevel =
expr.child.dataType match {
- // Scalar FLOAT/DOUBLE comparison keys are normalized natively (NaN
payloads folded together,
- // signed zeros tied) for Sort, TopK, Window, WindowGroupLimit, and range
partitioning, which
- // matches Spark's SQLOrderingUtil. Only the comparison key is normalized;
returned values keep
- // their original NaN representation and zero sign. So these are
compatible even in strict mode.
- case _: FloatType | _: DoubleType => Compatible()
- // Floating-point values nested in arrays, structs, or maps are still
compared with Arrow's raw
- // total ordering, under which -0.0 sorts below 0.0 and a sign-bit NaN
sorts below -Infinity.
- // https://github.com/apache/datafusion-comet/issues/5507
- case dt =>
+ override def getSupportLevel(expr: SortOrder): SupportLevel = {
+ val dataType = expr.child.dataType
+ if (canHoldNestedNull(dataType)) {
SupportLevel
- .strictFloatingPointReason(dt, nestedFloatingPointSort)
+ .strictFloatingPointReason(dataType, nullableNestedFloatingPointSort)
.map(reason => Incompatible(Some(reason)))
.getOrElse(Compatible())
+ } else {
+ Compatible()
Review Comment:
[P2] Preserve fallback for deeply nested keys in `RANGE` windows. With
`spark.comet.exec.strictFloatingPoint=true` and
`spark.comet.expression.SortOrder.allowIncompatible=false`, a Parquet table
`t(id INT, d DOUBLE)` containing `(1,1.0),(2,3.0)` exposes this through `SELECT
id, COUNT(id) OVER (ORDER BY array(named_struct('x', coalesce(d, 0.0D))), id)
AS running FROM t ORDER BY id`. Spark returns `[(1,1),(2,2)]`, but the native
window fails with `Uncomparable values: List([{x: 1.0}]), List([{x: 1.0}])`.
This key cannot contain nested nulls, so the new branch declares it compatible.
DataFusion’s frame comparator still cannot compare arrays of structs, even
after float normalization. The base strict-mode guard kept this valid query on
Spark. Could we add a window-specific fallback for these unsupported comparison
shapes, or fix their frame comparisons before admitting them?
Evidence: Freshly compiled
`/tmp/review6475-active-u51_pmld/physical_range.rs` includes this head’s
`normalize.rs` and constructs normalized `SortExec` keys followed by
`BoundedWindowAggExec` with `InputOrderMode::Sorted`, matching the current
planner. With locked Arrow 59.3.0/DataFusion 55.1.0, the null-free
array-of-struct RANGE case fails as quoted. Struct-of-array and array-of-array
keys also fail. All corresponding ROWS cases and the flat array<double> RANGE
control pass. A fresh Spark 3.5.9 run returned [(1,1),(2,2)] and confirmed
containsNull=false and struct-field nullable=false. Base
`CometSortOrder.getSupportLevel` rejects this nested DOUBLE key in strict mode.
Current `CometWindowExec` accepts its serialized order, and DataFusion’s
`ScalarValue::partial_cmp_list` returns None when the underlying Arrow
comparison rejects struct elements.
--
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]