andygrove commented on code in PR #6475:
URL: https://github.com/apache/datafusion-comet/pull/6475#discussion_r4169049284
##########
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:
Fixed in 4cf0b1cfd. DataFusion finds a `RANGE` frame's `CURRENT ROW` bound
with `ScalarValue::partial_cmp`, which cannot order an array of arrays or
structs, or a struct holding an array (apache/datafusion#24937). That fails
with the default settings too, whatever the element type: before this commit,
`SUM(score) OVER (PARTITION BY player ORDER BY array(named_struct('x',
score)))` over `INT` scores failed with the same error. So instead of widening
the strict-mode guard in `CometSortOrder`, which would also send sorts and
ranks over these keys back to Spark, `CometWindowExec` now falls back for a
`RANGE` frame with a bound other than `UNBOUNDED` when an `ORDER BY` key has
one of those shapes. A struct of structs stays native, because DataFusion
flattens nested structs before comparing. Ranking functions, `ROWS` frames and
an unbounded `RANGE` frame stay native over the same keys, and so does
`CUME_DIST`, whose `RANGE` frame DataFusion never reads.
`nested_float_order_keys_strict.sql` now runs your query and a
struct-of-array key in strict mode: both fall back and match Spark, while
`RANK` and a `ROWS` frame over the same keys stay native. A new section 5.6 in
`window_functions.sql` covers array-of-struct, struct-of-array and
array-of-array keys with `INT` values under the default settings. Before the
change, both fixtures failed with `Uncomparable values`.
--
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]