From 8e64b0972e18ba1d088fbee962a852d1f8712e95 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Wed, 30 Sep 2026 13:48:43 -0600 Subject: [PATCH 1/5] fix: match Spark's float ordering for sort and window keys nested in arrays and structs Spark orders floats with SQLOrderingUtil.compareDoubles at every depth of an array or struct. Comet normalized scalar float sort and window keys so that Arrow's total order agrees with Spark's, but compared nested keys raw, so -0.0 sorted below 0.0 and a NaN with the sign bit set below -Infinity. Wrap array and struct keys with a float leaf in NormalizeNestedFloats in create_normalized_key_expr, which builds the keys of Sort, TopK, Window and WindowGroupLimit. Only the comparison key is normalized, so returned values keep their bits. CometSortOrder is compatible for every key type, so strict floating-point mode no longer falls back for these keys. Closes #5507. --- .ai/skills/review-comet-shuffle-pr/SKILL.md | 11 +- .../contributor-guide/native_shuffle.md | 4 +- .../latest/compatibility/floating-point.md | 21 ++-- .../latest/compatibility/operators.md | 9 +- .../user-guide/latest/tuning/operators.md | 14 +-- native/core/src/execution/planner.rs | 60 ++++++++-- .../scala/org/apache/comet/CometConf.scala | 8 +- .../apache/comet/serde/CometSortOrder.scala | 38 ++---- .../windows/nested_float_order_keys.sql | 108 ++++++++++++++++++ .../apache/comet/CometExpressionSuite.scala | 86 +++++++++----- 10 files changed, 250 insertions(+), 109 deletions(-) create mode 100644 spark/src/test/resources/sql-tests/windows/nested_float_order_keys.sql diff --git a/.ai/skills/review-comet-shuffle-pr/SKILL.md b/.ai/skills/review-comet-shuffle-pr/SKILL.md index bc4eb9c9468..f003e120e86 100644 --- a/.ai/skills/review-comet-shuffle-pr/SKILL.md +++ b/.ai/skills/review-comet-shuffle-pr/SKILL.md @@ -55,12 +55,11 @@ keys**, and it is not the same rule for the two partitionings: and rejects collated strings, because Comet compares raw bytes. That is the whole rule. Scalar float and double are accepted, under `spark.comet.exec.strictFloatingPoint` as well: since [#5981](https://github.com/apache/datafusion-comet/pull/5981) the native range partitioner - normalizes its comparison keys and its sampled boundary rows the same way the native sort does, - and `CometSortOrder.getSupportLevel` returns `Compatible()` for them regardless of strict mode. - What strict mode still governs is floating point _nested_ in an array, struct, or map - ([#5507](https://github.com/apache/datafusion-comet/issues/5507)), which is already out as a - range key for being nested. `CometNativeShuffleSuite` runs "range partitioning on floating-point - uses native shuffle" under both settings of the config. + normalizes its comparison keys and its sampled boundary rows the same way the native sort does. + `CometSortOrder` is `Compatible()` for every key type regardless of strict mode, because the + native sort also normalizes floats nested in arrays and structs, but a nested key is already out + as a range key for being nested. `CometNativeShuffleSuite` runs "range partitioning on + floating-point uses native shuffle" under both settings of the config. - **`HashPartitioning` is primitive-only by default only.** With `spark.comet.shuffle.native.partitioning.hash.nested.enabled=true` (default `false`), `supportedHashPartitioningDataType` admits structs and arrays recursively, and maps on Spark diff --git a/docs/source/contributor-guide/native_shuffle.md b/docs/source/contributor-guide/native_shuffle.md index e1cabfd6c61..b882acd25a0 100644 --- a/docs/source/contributor-guide/native_shuffle.md +++ b/docs/source/contributor-guide/native_shuffle.md @@ -63,8 +63,8 @@ Native shuffle (`CometExchange`) is selected when all of the following condition compares raw bytes. Scalar float and double are supported, including when `spark.comet.exec.strictFloatingPoint` is enabled, because the native range partitioner normalizes its comparison keys and its sampled boundary rows the same way the native sort - does. Strict floating point only affects floating-point values nested in arrays, structs, or - maps, which are rejected as range keys for being nested anyway. + does. The native sort normalizes floating-point values nested in arrays and structs as well, + but those keys are rejected as range keys for being nested. - `HashPartitioning` keys must be primitive **by default**. Setting `spark.comet.shuffle.native.partitioning.hash.nested.enabled` to `true` admits structs and arrays as keys, checked recursively to their leaves, and maps on Spark 4.0 and later, where diff --git a/docs/source/user-guide/latest/compatibility/floating-point.md b/docs/source/user-guide/latest/compatibility/floating-point.md index 57a9ffe8a83..f9506b64f3b 100644 --- a/docs/source/user-guide/latest/compatibility/floating-point.md +++ b/docs/source/user-guide/latest/compatibility/floating-point.md @@ -48,20 +48,15 @@ Spark's `ORDER BY`, `RANK`, `DENSE_RANK`, and window frame comparisons route thr `SQLOrderingUtil.compareDoubles` / `compareFloats`, which equate all NaN representations and define `-0.0 == 0.0`. NaN sorts above every non-NaN value. -For scalar `FLOAT` and `DOUBLE` keys, Comet normalizes NaNs and signed zeros before native -sorting, window peer comparisons, and `WindowGroupLimitExec` rank comparisons. Native range -partitioning normalizes its keys and sampled boundaries in the same way. Only comparison keys -are normalized; returned values retain their original NaN representations and zero signs. +For `FLOAT` and `DOUBLE` keys, and for keys that nest them in arrays and structs at any depth, +Comet normalizes NaNs and signed zeros before native sorting, window peer comparisons, and +`WindowGroupLimitExec` rank comparisons. Native range partitioning normalizes its keys and +sampled boundaries in the same way; it only accepts scalar keys. Only comparison keys are +normalized; returned values retain their original NaN representations and zero signs. -Native sorting of floating-point values nested in arrays or structs still uses Arrow's raw total -ordering. Nested keys can therefore produce different ordering or rank results from Spark; see -[#5507](https://github.com/apache/datafusion-comet/issues/5507). - -Because those scalar comparison keys match Spark, `spark.comet.exec.strictFloatingPoint=true` no -longer forces a fallback for them: scalar `FLOAT` and `DOUBLE` sort keys, window and rank order -keys, and range partitioning keys all stay native under strict mode. Floating-point values nested -in arrays, structs, or maps still fall back under strict mode, because their ordering is the raw -total ordering described above. +Because those comparison keys match Spark, `spark.comet.exec.strictFloatingPoint=true` does not +force a fallback for them: sort keys, window and rank order keys, and range partitioning keys all +stay native under strict mode, whether the floats in them are scalar or nested. `array_min` and `array_max` use Spark-compatible native comparisons in both strict and non-strict floating-point modes. Signed zeros compare equal, and all NaN representations compare equal and diff --git a/docs/source/user-guide/latest/compatibility/operators.md b/docs/source/user-guide/latest/compatibility/operators.md index c637c197790..9101d4a5a35 100644 --- a/docs/source/user-guide/latest/compatibility/operators.md +++ b/docs/source/user-guide/latest/compatibility/operators.md @@ -103,13 +103,8 @@ runs natively; it is controlled by `spark.comet.exec.windowGroupLimit.enabled` ( (e.g. `UTF8_LCASE`). The native operator detects partitions and order-key peer groups by comparing Arrow row-encoded keys for byte equality, which splits peers that Spark ties. -**Known incompatibilities:** - -- Floating-point values nested in array or struct `ORDER BY` keys are compared with Arrow's raw - total ordering, so ranks can differ from Spark when the data mixes `-0.0` and `+0.0` or more - than one NaN representation ([#5507](https://github.com/apache/datafusion-comet/issues/5507)). - Scalar `FLOAT` and `DOUBLE` keys are normalized and match Spark; see - [floating-point ordering](./floating-point.md). +Floating-point `ORDER BY` keys, including floats nested in arrays and structs, are normalized +and match Spark's ranks; see [floating-point ordering](./floating-point.md). ## Round-Robin Partitioning diff --git a/docs/source/user-guide/latest/tuning/operators.md b/docs/source/user-guide/latest/tuning/operators.md index 8045d5855e8..77545d4ded2 100644 --- a/docs/source/user-guide/latest/tuning/operators.md +++ b/docs/source/user-guide/latest/tuning/operators.md @@ -151,16 +151,10 @@ See [TopK metrics](../metrics.md#local-topk). ## Optimizing Sorting on Floating-Point Values -Comet normalizes NaN payloads and signed zeros in scalar `FLOAT` and `DOUBLE` ordering keys, so `ORDER BY`, window -ordering and range partitioning on them match Spark and stay native even with -`spark.comet.exec.strictFloatingPoint=true`. Only the comparison key is normalized; returned values keep their original -NaN representation and zero sign. - -Floating-point values nested in arrays, structs, or maps are compared with Arrow's raw total ordering instead, which can -differ from Spark when the data contains both zero and negative zero, or more than one NaN representation. This is likely -an edge case that is not of concern for many users. Setting `spark.comet.exec.strictFloatingPoint=true` makes those -nested cases fall back to Spark, and they can be forced back onto the native path with -`spark.comet.expression.SortOrder.allowIncompatible=true`. +Comet normalizes NaN payloads and signed zeros in `FLOAT` and `DOUBLE` ordering keys, including floating-point values +nested in arrays and structs, so `ORDER BY`, window ordering and range partitioning on them match Spark and stay native +even with `spark.comet.exec.strictFloatingPoint=true`. Only the comparison key is normalized; returned values keep their +original NaN representation and zero sign. `sort_array` is separate. It sorts array elements rather than ordering rows, and its elements are compared with Arrow's raw total ordering, so `spark.comet.exec.strictFloatingPoint=true` makes it fall back even for a scalar floating-point diff --git a/native/core/src/execution/planner.rs b/native/core/src/execution/planner.rs index 2f992c4843f..5ee0ae2660d 100644 --- a/native/core/src/execution/planner.rs +++ b/native/core/src/execution/planner.rs @@ -149,9 +149,9 @@ use datafusion_comet_spark_expr::{ create_case_when, jvm_udf::JvmScalarUdfExpr, normalize_floats, spark_in_list, ApproxPercentile, ArrayInsert, Avg, AvgDecimal, Cast, CheckOverflow, Correlation, Covariance, CreateNamedStruct, DecimalRescaleCheckOverflow, GetArrayStructFields, GetStructField, HllPlusPlus, HllSketchAgg, - HllUnionAgg, IfExpr, ListExtract, MaxMinBy, Mode, NormalizeNaNAndZero, Regr, RegrType, - SparkCastOptions, Stddev, SumDecimal, ToJson, UnboundColumn, Variance, WideDecimalBinaryExpr, - WideDecimalOp, + HllUnionAgg, IfExpr, ListExtract, MaxMinBy, Mode, NormalizeNaNAndZero, NormalizeNestedFloats, + Regr, RegrType, SparkCastOptions, Stddev, SumDecimal, ToJson, UnboundColumn, Variance, + WideDecimalBinaryExpr, WideDecimalOp, }; use itertools::Itertools; use jni::objects::{Global, JObject}; @@ -1005,16 +1005,18 @@ impl PhysicalPlanner { } } - /// Normalize scalar floating-point comparison keys without changing output values. - /// Sort, Window, and WindowGroupLimit must use identical expressions so DataFusion - /// can recognize the ordering of window partition keys. + /// Normalize floating-point comparison keys without changing output values: a `FLOAT` or + /// `DOUBLE` key, and an array or struct key with a float at any depth, whose order Arrow + /// otherwise takes from the raw bits. Sort, Window, and WindowGroupLimit must use identical + /// expressions so DataFusion can recognize the ordering of window partition keys. fn create_normalized_key_expr( &self, spark_expr: &Expr, input_schema: SchemaRef, ) -> Result, ExecutionError> { let child = self.create_expr(spark_expr, Arc::clone(&input_schema))?; - Ok(NormalizeNaNAndZero::wrap_if_needed( + let child = NormalizeNaNAndZero::wrap_if_needed(child, input_schema.as_ref())?; + Ok(NormalizeNestedFloats::wrap_if_needed( child, input_schema.as_ref(), )?) @@ -5477,6 +5479,38 @@ mod tests { .collect() } + /// `floating_sort_batches` with each key wrapped in a one-element list and in a one-field + /// struct, which have to order and tie exactly as the bare key does. A null key becomes a + /// null list or struct. + fn nested_floating_sort_batches() -> Vec<(i32, RecordBatch)> { + use arrow::array::StructArray; + use arrow::buffer::OffsetBuffer; + floating_sort_batches() + .into_iter() + .flat_map(|(type_id, batch)| { + let values = Arc::clone(batch.column(0)); + let nulls = values.nulls().cloned(); + let element = Field::new("item", values.data_type().clone(), true); + let list: ArrayRef = Arc::new(ListArray::new( + Arc::new(element), + OffsetBuffer::from_lengths(vec![1; values.len()]), + Arc::clone(&values), + nulls.clone(), + )); + let fields = Fields::from(vec![Field::new("v", values.data_type().clone(), true)]); + let record: ArrayRef = Arc::new(StructArray::new(fields, vec![values], nulls)); + [list, record].into_iter().map(move |key| { + let mut fields: Vec = batch.schema().fields().to_vec(); + fields[0] = Arc::new(Field::new("ord", key.data_type().clone(), true)); + let mut columns = batch.columns().to_vec(); + columns[0] = key; + let batch = RecordBatch::try_new(Arc::new(Schema::new(fields)), columns); + (type_id, batch.unwrap()) + }) + }) + .collect() + } + #[test] fn floating_window_partition_keys_preserve_ordering() { let planner = PhysicalPlanner::default(); @@ -5544,7 +5578,11 @@ mod tests { async fn floating_sort_keys_preserve_window_group_limit_peers() { let planner = PhysicalPlanner::default(); let context = SessionContext::new_with_config(SessionConfig::new().with_batch_size(3)); - for (type_id, batch) in floating_sort_batches() { + // Floats nested in a list or a struct must rank exactly as the bare floats do. + let batches = floating_sort_batches() + .into_iter() + .chain(nested_floating_sort_batches()); + for (type_id, batch) in batches { for (descending, kind, fetch, expected) in [ (true, WindowFnKind::Rank, 3, vec![0, 1, 2, 8, 9]), (true, WindowFnKind::DenseRank, 2, vec![0, 1, 2, 8, 9]), @@ -5596,8 +5634,10 @@ mod tests { } ids.sort_unstable(); assert_eq!( - ids, expected, - "type={type_id}, descending={descending}, {kind:?}" + ids, + expected, + "type={type_id}, key={}, descending={descending}, {kind:?}", + batch.schema().field(0).data_type() ); // The zero peer group and the second NaN peer group straddle size-3 // sort output batches. Tie state must survive those boundaries. diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index 4550f73b6c2..ae14d965709 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -1061,10 +1061,10 @@ object CometConf extends ShimCometConf { .category(CATEGORY_EXEC) .doc( "When enabled, fall back to Spark for floating-point operations that may differ from " + - "Spark, such as comparing -0.0 and 0.0, sorting floating-point values nested in " + - "arrays, structs, or maps, or sorting the elements of a floating-point array with " + - "`sort_array`. Scalar `ORDER BY`, window ordering and range partitioning keys are " + - "unaffected, because Comet normalizes those comparison keys to match Spark. " + + "Spark, such as comparing -0.0 and 0.0, or sorting the elements of a floating-point " + + "array with `sort_array`. `ORDER BY`, window ordering and range partitioning keys " + + "are unaffected, including floating-point values nested in arrays and structs, " + + "because Comet normalizes those comparison keys to match Spark. " + s"$COMPAT_GUIDE.") .booleanConf .createWithDefault(false) diff --git a/spark/src/main/scala/org/apache/comet/serde/CometSortOrder.scala b/spark/src/main/scala/org/apache/comet/serde/CometSortOrder.scala index f9a31972748..b872297228e 100644 --- a/spark/src/main/scala/org/apache/comet/serde/CometSortOrder.scala +++ b/spark/src/main/scala/org/apache/comet/serde/CometSortOrder.scala @@ -20,40 +20,20 @@ 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.comet.CometConf import org.apache.comet.serde.QueryPlanSerde.exprToProtoInternal +/** + * The key of a Sort, TopK, Window, WindowGroupLimit or range partitioning. + * + * Every key type is compatible, in strict floating-point mode too. 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. + */ 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" - - override def getIncompatibleReasons(): Seq[String] = Seq( - s"$nestedFloatingPointSort 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 => - SupportLevel - .strictFloatingPointReason(dt, nestedFloatingPointSort) - .map(reason => Incompatible(Some(reason))) - .getOrElse(Compatible()) - } - override def convert( expr: SortOrder, inputs: Seq[Attribute], diff --git a/spark/src/test/resources/sql-tests/windows/nested_float_order_keys.sql b/spark/src/test/resources/sql-tests/windows/nested_float_order_keys.sql new file mode 100644 index 00000000000..fff436ea69f --- /dev/null +++ b/spark/src/test/resources/sql-tests/windows/nested_float_order_keys.sql @@ -0,0 +1,108 @@ +-- Licensed to the Apache Software Foundation (ASF) under one +-- or more contributor license agreements. See the NOTICE file +-- distributed with this work for additional information +-- regarding copyright ownership. The ASF licenses this file +-- to you under the Apache License, Version 2.0 (the +-- "License"); you may not use this file except in compliance +-- with the License. You may obtain a copy of the License at +-- +-- http://www.apache.org/licenses/LICENSE-2.0 +-- +-- Unless required by applicable law or agreed to in writing, +-- software distributed under the License is distributed on an +-- "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +-- KIND, either express or implied. See the License for the +-- specific language governing permissions and limitations +-- under the License. + +-- Sort and window keys that nest floats in arrays and structs follow Spark's ordering at every +-- depth: -0.0 equals 0.0, every NaN equals every other NaN, and NaN sorts above every other +-- value. Only the comparison keys are normalized, so returned values keep their bits. +-- +-- The keys use IF(s, -d, d): rows 1 and 2 hold -0.0 and 0.0, and rows 4 and 5 a canonical NaN +-- and a NaN with the sign bit set, the NaN that arithmetic produces on x86-64, which Arrow's +-- total order sorts below -Infinity. Each ORDER BY ends with a unique tiebreaker, so peers come +-- out in the tiebreaker's order. + +-- Strict floating-point mode no longer needs to fall back for these keys. The test harness +-- admits incompatible sort orders by default, so turn that off to check the shipped policy. +-- ConfigMatrix: spark.comet.exec.strictFloatingPoint=false,true +-- Config: spark.comet.expression.SortOrder.allowIncompatible=false + +statement +CREATE TABLE nested_float_keys(id INT, g INT, d DOUBLE, f FLOAT, s BOOLEAN) USING parquet + +statement +INSERT INTO nested_float_keys VALUES + (1, 1, 0.0D, float('0.0'), true), + (2, 1, 0.0D, float('0.0'), false), + (3, 1, 1.0D, float('1.0'), false), + (4, 2, double('NaN'), float('NaN'), false), + (5, 2, double('NaN'), float('NaN'), true), + (6, 2, double('Infinity'), float('Infinity'), false), + (7, 2, -1.0D, float('-1.0'), false), + (8, 1, NULL, NULL, false), + (9, 2, double('-Infinity'), float('-Infinity'), false) + +-- Sort. Row 8's keys hold a null element, which Spark orders below every other value whatever +-- the key's null order. Arrow ties the order of nested nulls to NULLS FIRST or LAST, so these +-- queries keep the default null order, where the two agree. +query +SELECT id FROM nested_float_keys ORDER BY array(IF(s, -d, d)), id DESC + +query +SELECT id FROM nested_float_keys ORDER BY array(IF(s, -f, f)) DESC, id + +query +SELECT id FROM nested_float_keys ORDER BY named_struct('x', IF(s, -d, d)) DESC, id DESC + +query +SELECT id FROM nested_float_keys ORDER BY named_struct('x', IF(s, -f, f)), id + +-- Nested more deeply, and more than one nested key +query +SELECT id FROM nested_float_keys +ORDER BY array(named_struct('x', IF(s, -d, d))), named_struct('a', array(IF(s, -f, f))) DESC, id + +-- Returning the keys as well. CometExpressionSuite checks that their bits come back unchanged. +query +SELECT id, array(IF(s, -d, d)) AS k, named_struct('x', IF(s, -f, f)) AS t +FROM nested_float_keys ORDER BY k, id DESC + +-- TopK +query +SELECT id FROM nested_float_keys ORDER BY array(IF(s, -d, d)) DESC, id LIMIT 4 + +query +SELECT id FROM nested_float_keys ORDER BY named_struct('x', IF(s, -f, f)), id DESC LIMIT 4 + +-- Window order keys: peers share a rank, and the default RANGE frame of a running sum spans +-- all of them. The running sums leave out row 8: DataFusion finds a RANGE frame's end by ordering +-- a null element above every value, while the sort puts it first, so the frame of every row +-- after it would run to the end of the partition, with or without floats. +query +SELECT id, + RANK() OVER (ORDER BY array(IF(s, -d, d))) AS r, + DENSE_RANK() OVER (PARTITION BY g ORDER BY named_struct('x', IF(s, -f, f)) DESC) AS dr +FROM nested_float_keys + +-- Comet declines to sort on a lone struct column, which Arrow cannot sort, so every struct key +-- comes with a second key or a partition. +query +SELECT id, + RANK() OVER (PARTITION BY g ORDER BY named_struct('x', IF(s, -d, d))) AS r, + SUM(id) OVER (ORDER BY array(IF(s, -d, d))) AS running, + SUM(id) OVER (PARTITION BY g ORDER BY array(IF(s, -f, f)) DESC) AS by_group +FROM nested_float_keys WHERE id <> 8 + +-- A rank limit keeps every peer of the last rank it admits +query +SELECT id FROM ( + SELECT id, RANK() OVER (ORDER BY array(IF(s, -d, d))) AS r FROM nested_float_keys +) WHERE r <= 3 + +query +SELECT id FROM ( + SELECT id, DENSE_RANK() OVER (PARTITION BY g ORDER BY named_struct('x', IF(s, -f, f)) DESC) AS r + FROM nested_float_keys +) WHERE r <= 1 diff --git a/spark/src/test/scala/org/apache/comet/CometExpressionSuite.scala b/spark/src/test/scala/org/apache/comet/CometExpressionSuite.scala index 0540b8df259..11ef05bf49c 100644 --- a/spark/src/test/scala/org/apache/comet/CometExpressionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometExpressionSuite.scala @@ -114,9 +114,8 @@ class CometExpressionSuite extends CometTestBase with AdaptiveSparkPlanHelper { df.withColumn("id", monotonically_increasing_id()).write.parquet(path) spark.read.parquet(path).createOrReplaceTempView("tbl") - // Scalar floating-point sort keys are normalized natively, so strict mode admits them even - // with allowIncompatible off. Nested floating-point keys are not, and the two tests below - // still assert the strict-mode fallback. + // Floating-point sort keys, scalar or nested, are normalized natively, so strict mode admits + // them even with allowIncompatible off. withSQLConf( CometConf.getExprAllowIncompatConfigKey("SortOrder") -> "false", CometConf.COMET_EXEC_STRICT_FLOATING_POINT.key -> "true") { @@ -243,51 +242,82 @@ class CometExpressionSuite extends CometTestBase with AdaptiveSparkPlanHelper { } } - test("sort array of floating point with negative zero") { - val schema = StructType( - Seq( - StructField("c0", DataTypes.createArrayType(DataTypes.FloatType), true), - StructField("c1", DataTypes.createArrayType(DataTypes.DoubleType), true))) + // Floats nested in array and struct sort keys are normalized natively as well, so strict mode + // keeps these sorts native. A unique `id` sorted last makes the ordering total, as above. + private def checkStrictNestedFloatingPointSort(schema: StructType): Unit = { val df = FuzzDataGenerator.generateDataFrame( new Random(42), spark, schema, 1000, DataGenOptions(generateNegativeZero = true)) - df.createOrReplaceTempView("tbl") - withSQLConf( - CometConf.getExprAllowIncompatConfigKey("SortOrder") -> "false", - CometConf.COMET_EXEC_STRICT_FLOATING_POINT.key -> "true") { - checkSparkAnswerAndFallbackReason( - "select * from tbl order by 1, 2", - "unsupported range partitioning sort order") + withTempDir { dir => + val path = new Path(dir.toString, "tbl").toString + df.withColumn("id", monotonically_increasing_id()).write.parquet(path) + spark.read.parquet(path).createOrReplaceTempView("tbl") + + withSQLConf( + CometConf.getExprAllowIncompatConfigKey("SortOrder") -> "false", + CometConf.COMET_EXEC_STRICT_FLOATING_POINT.key -> "true") { + checkSparkAnswerAndOperator( + sql("select * from tbl order by 1, 2, 3"), + Seq(classOf[CometSortExec])) + } } } + test("sort array of floating point with negative zero") { + checkStrictNestedFloatingPointSort( + StructType( + Seq( + StructField("c0", DataTypes.createArrayType(DataTypes.FloatType), true), + StructField("c1", DataTypes.createArrayType(DataTypes.DoubleType), true)))) + } + test("sort struct containing floating point with negative zero") { - val schema = StructType( - Seq( + checkStrictNestedFloatingPointSort( + StructType(Seq( StructField( "float_struct", StructType(Seq(StructField("c0", DataTypes.FloatType, true)))), StructField( "float_double", - StructType(Seq(StructField("c0", DataTypes.DoubleType, true)))))) - val df = FuzzDataGenerator.generateDataFrame( - new Random(42), - spark, - schema, - 1000, - DataGenOptions(generateNegativeZero = true)) - df.createOrReplaceTempView("tbl") + StructType(Seq(StructField("c0", DataTypes.DoubleType, true))))))) + } + + test("strict floating point: nested sort keeps NaN payloads and zero signs unchanged") { + // As in the scalar test, a local relation keeps the raw bits that Parquet would canonicalize. + val negNan = java.lang.Double.longBitsToDouble(0xfff8000000000002L) + val posNan = java.lang.Double.longBitsToDouble(0x7ff8000000000002L) + val rows = Seq((0, negNan), (1, posNan), (2, -0.0d), (3, 0.0d), (4, 1.0d)) + val expected = rows.map { case (id, d) => + id -> java.lang.Double.doubleToRawLongBits(d) + }.toMap + + rows.toDF("id", "d").createOrReplaceTempView("strict_fp_nested_bits") withSQLConf( CometConf.getExprAllowIncompatConfigKey("SortOrder") -> "false", CometConf.COMET_EXEC_STRICT_FLOATING_POINT.key -> "true") { - checkSparkAnswerAndFallbackReason( - "select * from tbl order by 1, 2", - "unsupported range partitioning sort order") + for ((key, value) <- Seq[(String, Row => Double)]( + "array(d)" -> (_.getSeq[Double](1).head), + "named_struct('v', d)" -> (_.getStruct(1).getDouble(0)))) { + val query = s"SELECT id, $key AS k FROM strict_fp_nested_bits ORDER BY k, id" + checkSparkAnswerAndOperator( + sql(query), + Seq(classOf[CometSortExec]), + classOf[LocalTableScanExec]) + + val actual = sql(query).collect().toSeq + actual.foreach { row => + assert( + java.lang.Double.doubleToRawLongBits(value(row)) == expected(row.getInt(0)), + s"row ${row.getInt(0)} had its floating-point bits rewritten by the sort on $key") + } + // The zeros are peers, and so are the two NaNs, which sort last. + assert(actual.map(_.getInt(0)) == Seq(2, 3, 4, 0, 1), s"sort on $key") + } } } From e0ec3e48b55ee30aa89d4aba120b5cd95d9d6c45 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Wed, 30 Sep 2026 14:05:06 -0600 Subject: [PATCH 2/5] test: link the nested-null issues from the nested float key fixture --- .../resources/sql-tests/windows/nested_float_order_keys.sql | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/spark/src/test/resources/sql-tests/windows/nested_float_order_keys.sql b/spark/src/test/resources/sql-tests/windows/nested_float_order_keys.sql index fff436ea69f..83e82e8bb73 100644 --- a/spark/src/test/resources/sql-tests/windows/nested_float_order_keys.sql +++ b/spark/src/test/resources/sql-tests/windows/nested_float_order_keys.sql @@ -46,7 +46,7 @@ INSERT INTO nested_float_keys VALUES -- Sort. Row 8's keys hold a null element, which Spark orders below every other value whatever -- the key's null order. Arrow ties the order of nested nulls to NULLS FIRST or LAST, so these --- queries keep the default null order, where the two agree. +-- queries keep the default null order, where the two agree (#6476). query SELECT id FROM nested_float_keys ORDER BY array(IF(s, -d, d)), id DESC @@ -79,7 +79,7 @@ SELECT id FROM nested_float_keys ORDER BY named_struct('x', IF(s, -f, f)), id DE -- Window order keys: peers share a rank, and the default RANGE frame of a running sum spans -- all of them. The running sums leave out row 8: DataFusion finds a RANGE frame's end by ordering -- a null element above every value, while the sort puts it first, so the frame of every row --- after it would run to the end of the partition, with or without floats. +-- after it would run to the end of the partition, with or without floats (#6477). query SELECT id, RANK() OVER (ORDER BY array(IF(s, -d, d))) AS r, From 6471565d82e0d54cd53047da0bdb28ab91ccb9d2 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Thu, 1 Oct 2026 13:18:20 -0600 Subject: [PATCH 3/5] fix: decline nullable nested float keys in strict mode and drop the nested float rank-limit guard Strict floating-point mode fell back for every float nested in a sort or window key, before those keys were normalized. Removing that fallback also admitted keys that can hold a null element or field, which the native sort and RANGE window frames order differently from Spark (#6476, #6477). CometSortOrder now keeps the fallback for a key that has a float and whose type can hold a nested null, and every other key stays native. #6468 made RANK and DENSE_RANK limits fall back for nested float order keys, because the native operator compared them raw. The native planner now normalizes those keys, so drop that guard and make its test expect native execution. Split the nested float order key fixture: the existing one runs with the default settings, and a new one runs in strict mode and pins both halves of the policy. --- .ai/skills/review-comet-shuffle-pr/SKILL.md | 11 +- .../latest/compatibility/floating-point.md | 10 ++ .../latest/compatibility/operators.md | 8 +- .../user-guide/latest/tuning/operators.md | 7 ++ .../scala/org/apache/comet/CometConf.scala | 4 +- .../apache/comet/serde/CometSortOrder.scala | 56 +++++++++- .../sql/comet/CometWindowGroupLimitExec.scala | 34 +----- .../windows/nested_float_order_keys.sql | 6 +- .../nested_float_order_keys_strict.sql | 103 ++++++++++++++++++ .../apache/comet/CometExpressionSuite.scala | 23 +++- .../comet/exec/CometWindowExecSuite.scala | 22 ++-- 11 files changed, 225 insertions(+), 59 deletions(-) create mode 100644 spark/src/test/resources/sql-tests/windows/nested_float_order_keys_strict.sql diff --git a/.ai/skills/review-comet-shuffle-pr/SKILL.md b/.ai/skills/review-comet-shuffle-pr/SKILL.md index f003e120e86..285c61b3628 100644 --- a/.ai/skills/review-comet-shuffle-pr/SKILL.md +++ b/.ai/skills/review-comet-shuffle-pr/SKILL.md @@ -56,10 +56,13 @@ keys**, and it is not the same rule for the two partitionings: float and double are accepted, under `spark.comet.exec.strictFloatingPoint` as well: since [#5981](https://github.com/apache/datafusion-comet/pull/5981) the native range partitioner normalizes its comparison keys and its sampled boundary rows the same way the native sort does. - `CometSortOrder` is `Compatible()` for every key type regardless of strict mode, because the - native sort also normalizes floats nested in arrays and structs, but a nested key is already out - as a range key for being nested. `CometNativeShuffleSuite` runs "range partitioning on - floating-point uses native shuffle" under both settings of the config. + `CometSortOrder` is `Compatible()` for scalar floats regardless of strict mode. It is for floats + nested in arrays and structs too, because the native sort normalizes them, unless the key's type + can hold a null element or field + ([#6476](https://github.com/apache/datafusion-comet/issues/6476), + [#6477](https://github.com/apache/datafusion-comet/issues/6477)), which strict mode declines. + A nested key is already out as a range key for being nested. `CometNativeShuffleSuite` runs + "range partitioning on floating-point uses native shuffle" under both settings of the config. - **`HashPartitioning` is primitive-only by default only.** With `spark.comet.shuffle.native.partitioning.hash.nested.enabled=true` (default `false`), `supportedHashPartitioningDataType` admits structs and arrays recursively, and maps on Spark diff --git a/docs/source/user-guide/latest/compatibility/floating-point.md b/docs/source/user-guide/latest/compatibility/floating-point.md index f9506b64f3b..46765f3aa2e 100644 --- a/docs/source/user-guide/latest/compatibility/floating-point.md +++ b/docs/source/user-guide/latest/compatibility/floating-point.md @@ -58,6 +58,16 @@ Because those comparison keys match Spark, `spark.comet.exec.strictFloatingPoint force a fallback for them: sort keys, window and rank order keys, and range partitioning keys all stay native under strict mode, whether the floats in them are scalar or nested. +The exception is a key that nests floats in an array or struct whose type can hold a null element +or field. Spark orders such a null below every other value, whatever the key's `NULLS FIRST` or +`NULLS LAST`. The native sort places it by that null order, so `ASC NULLS LAST` and +`DESC NULLS FIRST` can differ from Spark, and a `RANGE` window frame orders it above every other +value, so a running aggregate can span the whole partition +([#6476](https://github.com/apache/datafusion-comet/issues/6476), +[#6477](https://github.com/apache/datafusion-comet/issues/6477)). Strict mode makes those keys fall +back to Spark. A key whose type cannot hold a null, such as `array(coalesce(x, 0.0D))`, stays +native. + `array_min` and `array_max` use Spark-compatible native comparisons in both strict and non-strict floating-point modes. Signed zeros compare equal, and all NaN representations compare equal and greater than non-NaN values. The original first equal element is retained: for example, diff --git a/docs/source/user-guide/latest/compatibility/operators.md b/docs/source/user-guide/latest/compatibility/operators.md index 2d12a04c703..0be169f5ce9 100644 --- a/docs/source/user-guide/latest/compatibility/operators.md +++ b/docs/source/user-guide/latest/compatibility/operators.md @@ -102,14 +102,10 @@ runs natively; it is controlled by `spark.comet.exec.windowGroupLimit.enabled` ( - Any `PARTITION BY` or `ORDER BY` key whose type carries a non-default `StringType` collation (e.g. `UTF8_LCASE`). The native operator detects partitions and order-key peer groups by comparing Arrow row-encoded keys for byte equality, which splits peers that Spark ties. -- `RANK` and `DENSE_RANK` whose `ORDER BY` key has a `FLOAT` or `DOUBLE` nested in an array or - struct. The same byte equality decides their ties, and nested floating-point values aren't - normalized, so `-0.0` and `+0.0`, or two NaN representations, would get different ranks and the - cutoff would drop rows that Spark keeps - ([#5507](https://github.com/apache/datafusion-comet/issues/5507)). Floating-point `ORDER BY` keys, including floats nested in arrays and structs, are normalized -and match Spark's ranks; see [floating-point ordering](./floating-point.md). +and match Spark's ranks; see [floating-point ordering](./floating-point.md), which also covers +strict floating-point mode. ## Round-Robin Partitioning diff --git a/docs/source/user-guide/latest/tuning/operators.md b/docs/source/user-guide/latest/tuning/operators.md index 27afd239d1a..4dae2ad7ead 100644 --- a/docs/source/user-guide/latest/tuning/operators.md +++ b/docs/source/user-guide/latest/tuning/operators.md @@ -154,6 +154,13 @@ nested in arrays and structs, so `ORDER BY`, window ordering and range partition even with `spark.comet.exec.strictFloatingPoint=true`. Only the comparison key is normalized; returned values keep their original NaN representation and zero sign. +The exception is a key that nests floating-point values in an array or struct whose type can hold a null element or +field. Spark orders such a null below every other value, and the native sort and `RANGE` window frames do not +([#6476](https://github.com/apache/datafusion-comet/issues/6476), +[#6477](https://github.com/apache/datafusion-comet/issues/6477)), so +`spark.comet.exec.strictFloatingPoint=true` makes those keys fall back to Spark. They can be forced back onto the +native path with `spark.comet.expression.SortOrder.allowIncompatible=true`. + `sort_array` is separate. It sorts array elements rather than ordering rows, and its elements are compared with Arrow's raw total ordering, so `spark.comet.exec.strictFloatingPoint=true` makes it fall back even for a scalar floating-point element type. Use `spark.comet.expression.SortArray.allowIncompatible=true` to keep it native. diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index 96251ce270f..99a7a440d40 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -1082,7 +1082,9 @@ object CometConf extends ShimCometConf { "Spark, such as comparing -0.0 and 0.0, or sorting the elements of a floating-point " + "array with `sort_array`. `ORDER BY`, window ordering and range partitioning keys " + "are unaffected, including floating-point values nested in arrays and structs, " + - "because Comet normalizes those comparison keys to match Spark. " + + "because Comet normalizes those comparison keys to match Spark. The exception is a " + + "nested key whose type can hold a null element or field, which falls back until the " + + "native sort and window frames order those nulls as Spark does. " + s"$COMPAT_GUIDE.") .booleanConf .createWithDefault(false) diff --git a/spark/src/main/scala/org/apache/comet/serde/CometSortOrder.scala b/spark/src/main/scala/org/apache/comet/serde/CometSortOrder.scala index b872297228e..36b8f87bba9 100644 --- a/spark/src/main/scala/org/apache/comet/serde/CometSortOrder.scala +++ b/spark/src/main/scala/org/apache/comet/serde/CometSortOrder.scala @@ -20,20 +20,66 @@ package org.apache.comet.serde import org.apache.spark.sql.catalyst.expressions.{Ascending, Attribute, Descending, NullsFirst, NullsLast, SortOrder} +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. * - * Every key type is compatible, in strict floating-point mode too. 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. + * 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 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"$nullableNestedFloatingPointSort is not 100% compatible with Spark when " + + s"`${CometConf.COMET_EXEC_STRICT_FLOATING_POINT.key}=true`") + + override def getSupportLevel(expr: SortOrder): SupportLevel = { + val dataType = expr.child.dataType + if (canHoldNestedNull(dataType)) { + SupportLevel + .strictFloatingPointReason(dataType, nullableNestedFloatingPointSort) + .map(reason => Incompatible(Some(reason))) + .getOrElse(Compatible()) + } else { + Compatible() + } + } + + /** + * Whether a value of `dataType` can hold a null below its top level: an array element or a + * struct field, at any depth. The key being null is not one: both engines place that by the + * key's null order. + */ + private def canHoldNestedNull(dataType: DataType): Boolean = dataType match { + case ArrayType(elementType, containsNull) => containsNull || canHoldNestedNull(elementType) + case StructType(fields) => fields.exists(f => f.nullable || canHoldNestedNull(f.dataType)) + case _ => false + } + override def convert( expr: SortOrder, inputs: Seq[Attribute], diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/CometWindowGroupLimitExec.scala b/spark/src/main/scala/org/apache/spark/sql/comet/CometWindowGroupLimitExec.scala index 5ad9c9022cd..ee3cbd61c63 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/CometWindowGroupLimitExec.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/CometWindowGroupLimitExec.scala @@ -24,13 +24,12 @@ import scala.jdk.CollectionConverters._ import org.apache.spark.sql.catalyst.expressions.{Attribute, Expression, SortOrder} import org.apache.spark.sql.catalyst.plans.physical.Partitioning import org.apache.spark.sql.execution.SparkPlan -import org.apache.spark.sql.types.{DoubleType, FloatType} import com.google.common.base.Objects import org.apache.comet.{CometConf, ConfigEntry} import org.apache.comet.CometSparkSessionExtensions.withFallbackReason -import org.apache.comet.serde.{CometOperatorSerde, OperatorOuterClass, SupportLevel} +import org.apache.comet.serde.{CometOperatorSerde, OperatorOuterClass} import org.apache.comet.serde.OperatorOuterClass.{Operator, RankLikeFunction} import org.apache.comet.serde.QueryPlanSerde.{exprToProto, hasNonDefaultStringCollation} import org.apache.comet.shims.ShimCometWindowGroupLimit @@ -100,32 +99,11 @@ object CometWindowGroupLimitExec extends CometOperatorSerde[SparkPlan] { return None } - // The same byte equality decides RANK and DENSE_RANK ties. Scalar FLOAT and DOUBLE order keys - // are normalized natively first, but a float nested in an array or struct keeps its raw bits, - // so -0.0 and 0.0, or two NaN encodings, would split a tie that Spark keeps and the cutoff - // would drop rows. ROW_NUMBER never compares peers, and Spark normalizes floating-point - // partition keys before the limit is planned. - // https://github.com/apache/datafusion-comet/issues/5507 - if (fields.rankLikeFunction != RankLikeFunction.RowNumber) { - val nestedFloat = fields.orderSpec.map(_.child).filter { e => - e.dataType match { - case _: FloatType | _: DoubleType => false - case dt => SupportLevel.containsType(dt, classOf[FloatType], classOf[DoubleType]) - } - } - if (nestedFloat.nonEmpty) { - withFallbackReason( - op, - nestedFloat - .map(_.sql) - .mkString( - "WindowGroupLimit: RANK and DENSE_RANK compare nested floating-point values " + - "exactly, so they fall back for order key(s): ", - ", ", - "")) - return None - } - } + // The same byte equality decides RANK and DENSE_RANK ties, but it runs on keys that the native + // planner has normalized: a `FLOAT` or `DOUBLE` order key, scalar or nested in an array or + // struct, has its NaNs folded together and its signed zeros tied first, as Spark does, so no + // float key needs a guard here. Strict floating-point mode declines the nested keys that can + // hold a null element or field, through `CometSortOrder`. val childOutput = op.children.head.output val partitionProtos = fields.partitionSpec.map(e => e -> exprToProto(e, childOutput)) diff --git a/spark/src/test/resources/sql-tests/windows/nested_float_order_keys.sql b/spark/src/test/resources/sql-tests/windows/nested_float_order_keys.sql index 83e82e8bb73..e6441847661 100644 --- a/spark/src/test/resources/sql-tests/windows/nested_float_order_keys.sql +++ b/spark/src/test/resources/sql-tests/windows/nested_float_order_keys.sql @@ -24,10 +24,8 @@ -- total order sorts below -Infinity. Each ORDER BY ends with a unique tiebreaker, so peers come -- out in the tiebreaker's order. --- Strict floating-point mode no longer needs to fall back for these keys. The test harness --- admits incompatible sort orders by default, so turn that off to check the shipped policy. --- ConfigMatrix: spark.comet.exec.strictFloatingPoint=false,true --- Config: spark.comet.expression.SortOrder.allowIncompatible=false +-- Strict floating-point mode declines these keys, because their types can hold a null element or +-- field: see nested_float_order_keys_strict.sql. statement CREATE TABLE nested_float_keys(id INT, g INT, d DOUBLE, f FLOAT, s BOOLEAN) USING parquet diff --git a/spark/src/test/resources/sql-tests/windows/nested_float_order_keys_strict.sql b/spark/src/test/resources/sql-tests/windows/nested_float_order_keys_strict.sql new file mode 100644 index 00000000000..713b0cd4bf6 --- /dev/null +++ b/spark/src/test/resources/sql-tests/windows/nested_float_order_keys_strict.sql @@ -0,0 +1,103 @@ +-- Licensed to the Apache Software Foundation (ASF) under one +-- or more contributor license agreements. See the NOTICE file +-- distributed with this work for additional information +-- regarding copyright ownership. The ASF licenses this file +-- to you under the Apache License, Version 2.0 (the +-- "License"); you may not use this file except in compliance +-- with the License. You may obtain a copy of the License at +-- +-- http://www.apache.org/licenses/LICENSE-2.0 +-- +-- Unless required by applicable law or agreed to in writing, +-- software distributed under the License is distributed on an +-- "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +-- KIND, either express or implied. See the License for the +-- specific language governing permissions and limitations +-- under the License. + +-- Strict floating-point mode and sort or window keys that nest floats in arrays and structs. +-- nested_float_order_keys.sql checks that Comet normalizes those floats to match Spark. Strict +-- mode still declines a key whose type can hold a null element or field, because Spark orders +-- such a null below every value whatever the key's null order, the native sort places it by the +-- null order (#6476), and a RANGE window frame orders it above every value (#6477). A key whose +-- type cannot hold a null stays native. +-- +-- The test harness admits incompatible sort orders by default, so turn that off to check the +-- shipped policy. +-- Config: spark.comet.exec.strictFloatingPoint=true +-- Config: spark.comet.expression.SortOrder.allowIncompatible=false + +statement +CREATE TABLE nested_float_strict(id INT, g INT, d DOUBLE, f FLOAT, s BOOLEAN) USING parquet + +statement +INSERT INTO nested_float_strict VALUES + (1, 1, 0.0D, float('0.0'), true), + (2, 1, 0.0D, float('0.0'), false), + (3, 1, 1.0D, float('1.0'), false), + (4, 2, double('NaN'), float('NaN'), false), + (5, 2, double('NaN'), float('NaN'), true), + (6, 2, double('Infinity'), float('Infinity'), false), + (7, 2, -1.0D, float('-1.0'), false), + (8, 1, NULL, NULL, false), + (9, 2, double('-Infinity'), float('-Infinity'), false) + +-- Keys over the nullable columns can hold a null element or field, whether or not any row does, +-- so they fall back, and Spark's answers come back. Row 8 holds one. Native, ORDER BY array(d) +-- NULLS LAST puts row 8 last, where Spark puts it first, and the running sum spans the whole +-- partition on row 8 and every row after it. +query expect_fallback(can hold a null element or field) +SELECT id FROM nested_float_strict ORDER BY array(d) NULLS LAST, id + +query expect_fallback(can hold a null element or field) +SELECT id, SUM(id) OVER (ORDER BY array(d)) AS running FROM nested_float_strict + +query expect_fallback(can hold a null element or field) +SELECT id FROM nested_float_strict ORDER BY named_struct('x', f) DESC NULLS FIRST, id + +query expect_fallback(can hold a null element or field) +SELECT id FROM nested_float_strict ORDER BY array(IF(s, -d, d)), id LIMIT 4 + +query expect_fallback(can hold a null element or field) +SELECT id, + RANK() OVER (PARTITION BY g ORDER BY named_struct('x', IF(s, -d, d))) AS r, + DENSE_RANK() OVER (ORDER BY array(IF(s, -f, f)) DESC) AS dr +FROM nested_float_strict + +query expect_fallback(can hold a null element or field) +SELECT id FROM ( + SELECT id, RANK() OVER (ORDER BY array(IF(s, -d, d))) AS r FROM nested_float_strict +) WHERE r <= 3 + +-- coalesce with a literal is never null, so an array or struct of it cannot hold a null, and +-- these keys stay native. Row 8 becomes a third zero. Each ORDER BY ends with a unique +-- tiebreaker, so peers come out in its order, and the running sums need no null row to leave out. +query +SELECT id FROM nested_float_strict ORDER BY array(coalesce(IF(s, -d, d), 0.0D)), id DESC + +query +SELECT id FROM nested_float_strict ORDER BY named_struct('x', coalesce(IF(s, -f, f), 0.0F)) DESC, id + +query +SELECT id FROM nested_float_strict ORDER BY array(coalesce(IF(s, -d, d), 0.0D)) DESC, id LIMIT 4 + +query +SELECT id, + RANK() OVER (ORDER BY array(coalesce(IF(s, -d, d), 0.0D))) AS r, + SUM(id) OVER (ORDER BY array(coalesce(IF(s, -d, d), 0.0D))) AS running, + DENSE_RANK() OVER (PARTITION BY g ORDER BY named_struct('x', coalesce(IF(s, -f, f), 0.0F)) DESC) AS dr +FROM nested_float_strict + +-- A rank limit keeps every peer of the last rank it admits +query +SELECT id FROM ( + SELECT id, RANK() OVER (ORDER BY array(coalesce(IF(s, -d, d), 0.0D))) AS r + FROM nested_float_strict +) WHERE r <= 3 + +query +SELECT id FROM ( + SELECT id, + DENSE_RANK() OVER (PARTITION BY g ORDER BY named_struct('x', coalesce(IF(s, -f, f), 0.0F)) DESC) AS r + FROM nested_float_strict +) WHERE r <= 1 diff --git a/spark/src/test/scala/org/apache/comet/CometExpressionSuite.scala b/spark/src/test/scala/org/apache/comet/CometExpressionSuite.scala index 11ef05bf49c..af8aa7f0b9e 100644 --- a/spark/src/test/scala/org/apache/comet/CometExpressionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometExpressionSuite.scala @@ -114,8 +114,8 @@ class CometExpressionSuite extends CometTestBase with AdaptiveSparkPlanHelper { df.withColumn("id", monotonically_increasing_id()).write.parquet(path) spark.read.parquet(path).createOrReplaceTempView("tbl") - // Floating-point sort keys, scalar or nested, are normalized natively, so strict mode admits - // them even with allowIncompatible off. + // Scalar floating-point sort keys are normalized natively, so strict mode admits them even + // with allowIncompatible off. The nested tests below cover array and struct keys. withSQLConf( CometConf.getExprAllowIncompatConfigKey("SortOrder") -> "false", CometConf.COMET_EXEC_STRICT_FLOATING_POINT.key -> "true") { @@ -242,8 +242,11 @@ class CometExpressionSuite extends CometTestBase with AdaptiveSparkPlanHelper { } } - // Floats nested in array and struct sort keys are normalized natively as well, so strict mode - // keeps these sorts native. A unique `id` sorted last makes the ordering total, as above. + // Floats nested in array and struct sort keys are normalized natively as well, but strict mode + // declines a key whose type can hold a null element or field: Spark orders that null below + // every value whatever the key's null order, and the native sort places it by the null order + // (#6476). Every field the generator makes is nullable, so these sorts fall back. A unique `id` + // sorted last makes the ordering total, as above. private def checkStrictNestedFloatingPointSort(schema: StructType): Unit = { val df = FuzzDataGenerator.generateDataFrame( new Random(42), @@ -260,6 +263,15 @@ class CometExpressionSuite extends CometTestBase with AdaptiveSparkPlanHelper { withSQLConf( CometConf.getExprAllowIncompatConfigKey("SortOrder") -> "false", CometConf.COMET_EXEC_STRICT_FLOATING_POINT.key -> "true") { + checkSparkAnswerAndFallbackReason( + sql("select * from tbl order by 1, 2, 3"), + "can hold a null element or field") + } + + // The default null order agrees with Spark's, so opting in keeps the sort native and right. + withSQLConf( + CometConf.getExprAllowIncompatConfigKey("SortOrder") -> "true", + CometConf.COMET_EXEC_STRICT_FLOATING_POINT.key -> "true") { checkSparkAnswerAndOperator( sql("select * from tbl order by 1, 2, 3"), Seq(classOf[CometSortExec])) @@ -288,6 +300,9 @@ class CometExpressionSuite extends CometTestBase with AdaptiveSparkPlanHelper { test("strict floating point: nested sort keeps NaN payloads and zero signs unchanged") { // As in the scalar test, a local relation keeps the raw bits that Parquet would canonicalize. + // `d` is a primitive `Double`, so it is not nullable, and neither are the element of + // `array(d)` and the field of `named_struct('v', d)`: these keys cannot hold a null, and + // strict mode admits them. val negNan = java.lang.Double.longBitsToDouble(0xfff8000000000002L) val posNan = java.lang.Double.longBitsToDouble(0x7ff8000000000002L) val rows = Seq((0, negNan), (1, posNan), (2, -0.0d), (3, 0.0d), (4, 1.0d)) diff --git a/spark/src/test/scala/org/apache/comet/exec/CometWindowExecSuite.scala b/spark/src/test/scala/org/apache/comet/exec/CometWindowExecSuite.scala index 1e82f614ea7..700748b6edc 100644 --- a/spark/src/test/scala/org/apache/comet/exec/CometWindowExecSuite.scala +++ b/spark/src/test/scala/org/apache/comet/exec/CometWindowExecSuite.scala @@ -212,10 +212,12 @@ class CometWindowExecSuite extends CometTestBase { test("window group limit: floating-point values nested in the order key") { assume(isSpark35Plus, "WindowGroupLimit was added in Spark 3.5") - // The native RANK and DENSE_RANK find ties by byte equality, and nested floats aren't - // normalized, so [-0.0] and [0.0] would get different ranks and the cutoff would drop a row - // that Spark keeps (#5507). The limit falls back to Spark for them instead. ROW_NUMBER never - // compares peers, and Spark normalizes nested floating-point partition keys itself. + // The native RANK and DENSE_RANK find ties by byte equality, on order keys that the native + // planner has normalized, nested floats included. So [-0.0] and [0.0] are peers, and the + // cutoff keeps both rows, as Spark does (#5507). Strict floating-point mode declines a nested + // key whose type can hold a null element or field (#6476, #6477), and these columns, read + // back from Parquet, can. ROW_NUMBER never compares peers, and Spark normalizes nested + // floating-point partition keys itself. withTempDir { dir => val path = new Path(dir.toString, "nested_float_order").toString Seq((1, -0.0f, -0.0d), (2, 0.0f, 0.0d), (3, 1.0f, 1.0d)) @@ -235,14 +237,20 @@ class CometWindowExecSuite extends CometTestBase { // A struct needs a second key, because a lone struct sort key falls back on its own. orderBy <- Seq("array(f)", "array(d)", "named_struct('x', d), id > 0") } { - val query = sql(s""" + def rankLimit(): DataFrame = sql(s""" |SELECT id FROM ( | SELECT id, $rankFunction() OVER (ORDER BY $orderBy) AS rnk | FROM nested_float_order |) WHERE rnk <= 1 |""".stripMargin) - checkSparkAnswerAndFallbackReason(query, "compare nested floating-point values exactly") - assert(query.collect().map(_.getInt(0)).sorted.toSeq == Seq(1, 2)) + + checkSparkAnswerAndOperator(rankLimit(), Seq(classOf[CometWindowGroupLimitExec])) + assert(rankLimit().collect().map(_.getInt(0)).sorted.toSeq == Seq(1, 2)) + + withSQLConf(CometConf.COMET_EXEC_STRICT_FLOATING_POINT.key -> "true") { + checkSparkAnswerAndFallbackReason(rankLimit(), "can hold a null element or field") + assert(rankLimit().collect().map(_.getInt(0)).sorted.toSeq == Seq(1, 2)) + } } checkSparkAnswerAndOperator( From 6fa5836c39fba34985f0d3425305e595318fd2c2 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Fri, 2 Oct 2026 04:43:16 -0600 Subject: [PATCH 4/5] test: drop the sweep suite's known gap for nested float sort keys With floats nested in array and struct sort keys normalized, the nested ORDER BY cases match Spark (#5507), so the suite now requires them to. --- .../scala/org/apache/comet/CometFloatSemanticsSuite.scala | 4 ---- 1 file changed, 4 deletions(-) diff --git a/spark/src/test/scala/org/apache/comet/CometFloatSemanticsSuite.scala b/spark/src/test/scala/org/apache/comet/CometFloatSemanticsSuite.scala index 1e6aa60fdf1..96db282c46d 100644 --- a/spark/src/test/scala/org/apache/comet/CometFloatSemanticsSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometFloatSemanticsSuite.scala @@ -362,10 +362,6 @@ object CometFloatSemanticsSuite { c.group == "comparison" && (c.context.startsWith("array operands") || c.context.startsWith("struct operands")) && !Set("=", "!=").contains(c.variant)), - KnownGap( - issue(5507), - "Sort keys that nest floats in an array are compared raw.", - in("key", "nested ORDER BY")), KnownGap( issue(6385), "hash, xxhash64 and the native shuffle's hash partitioner hash a NaN's raw bits.", From 4cf0b1cfdf2cee68b6ad75df73d9b1666d079169 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Fri, 2 Oct 2026 05:24:05 -0600 Subject: [PATCH 5/5] fix: fall back for RANGE frames over ORDER BY keys DataFusion cannot compare DataFusion finds a RANGE frame's CURRENT ROW bound by comparing ORDER BY values with ScalarValue::partial_cmp, which cannot order an array of arrays or structs, or a struct holding an array (apache/datafusion#24937), so the query fails with "Uncomparable values". Strict floating-point mode declined these keys while nested float keys were compared raw, and once null-free nested float keys became compatible they reached this failure in strict mode too. The same shapes fail with the default settings whatever their element type. CometWindowExec now falls back for a RANGE frame bounded by CURRENT ROW over such a key. Ranking functions, CUME_DIST, ROWS frames and an unbounded RANGE frame never compare values this way and stay native. --- .../latest/compatibility/operators.md | 4 ++ .../spark/sql/comet/CometWindowExec.scala | 38 ++++++++++++- .../nested_float_order_keys_strict.sql | 22 ++++++++ .../sql-tests/windows/window_functions.sql | 56 +++++++++++++++++++ 4 files changed, 119 insertions(+), 1 deletion(-) diff --git a/docs/source/user-guide/latest/compatibility/operators.md b/docs/source/user-guide/latest/compatibility/operators.md index 0be169f5ce9..657a32de532 100644 --- a/docs/source/user-guide/latest/compatibility/operators.md +++ b/docs/source/user-guide/latest/compatibility/operators.md @@ -88,6 +88,10 @@ incorrect result. When any single window expression in a `WindowExec` falls back overflow instead of returning Spark's `NULL`. - `RANGE` frame with an explicit offset when the `ORDER BY` column is `DATE` or `DECIMAL` ([#4834](https://github.com/apache/datafusion-comet/issues/4834)). +- `RANGE` frame bounded by `CURRENT ROW` when an `ORDER BY` key is an array of arrays or structs, or a struct + holding an array, such as `array(named_struct('x', x))`. DataFusion cannot compare those values to find the + frame's bounds ([apache/datafusion#24937](https://github.com/apache/datafusion/issues/24937)). Ranking functions + and `ROWS` frames over the same keys run natively. - `first_value` / `last_value` on a `RANGE` frame with a literal offset ([#4835](https://github.com/apache/datafusion-comet/issues/4835)). - `lag` / `lead` with a non-literal default value ([#4268](https://github.com/apache/datafusion-comet/issues/4268)). diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/CometWindowExec.scala b/spark/src/main/scala/org/apache/spark/sql/comet/CometWindowExec.scala index 21496755294..20de01f8c30 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/CometWindowExec.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/CometWindowExec.scala @@ -28,7 +28,7 @@ import org.apache.spark.sql.execution.SparkPlan import org.apache.spark.sql.execution.metric.{SQLMetric, SQLMetrics} import org.apache.spark.sql.execution.window.WindowExec import org.apache.spark.sql.internal.SQLConf -import org.apache.spark.sql.types.{DataType, DateType, Decimal, DecimalType, DoubleType, LongType, NumericType} +import org.apache.spark.sql.types.{ArrayType, DataType, DateType, Decimal, DecimalType, DoubleType, LongType, MapType, NumericType, StructType} import com.google.common.base.Objects @@ -433,6 +433,27 @@ object CometWindowExec extends CometOperatorSerde[WindowExec] { case _ => } + // DataFusion finds a RANGE frame bound other than UNBOUNDED by comparing ORDER BY values, + // and its comparison cannot order an array of arrays or structs, or a struct holding an + // array, so execution fails with "Uncomparable values". Ranking functions and ROWS frames + // never compare values that way and stay native over the same keys. That includes CUME_DIST, + // whose RANGE frame DataFusion never reads. + // https://github.com/apache/datafusion/issues/24937 + f match { + case SpecifiedWindowFrame(RangeFrame, lb, ub) + if (lb != UnboundedPreceding || ub != UnboundedFollowing) && + !windowExpr.windowFunction.isInstanceOf[CumeDist] => + windowExpr.windowSpec.orderSpec.map(_.dataType).find(!isRangeComparable(_)) match { + case Some(dt) => + withFallbackReason( + windowExpr, + s"RANGE frame on ${dt.catalogString} ORDER BY is not supported") + return None + case None => + } + case _ => + } + val (frameType, lowerBound, upperBound) = f match { case SpecifiedWindowFrame(frameType, lBound, uBound) => val frameProto = frameType match { @@ -601,6 +622,21 @@ object CometWindowExec extends CometOperatorSerde[WindowExec] { SerializedPlan(None)) } + // Whether DataFusion can compare two ORDER BY values of `dataType` to find a RANGE frame + // bound. `ScalarValue::partial_cmp` compares an array's elements, and a struct's fields with + // nested structs flattened, using Arrow comparison kernels that reject nested values. + private def isRangeComparable(dataType: DataType): Boolean = dataType match { + case ArrayType(_: ArrayType | _: StructType | _: MapType, _) => false + case StructType(fields) => + fields.forall { field => + field.dataType match { + case _: ArrayType | _: MapType => false + case dt => isRangeComparable(dt) + } + } + case _ => true + } + // Folds a RANGE frame bound expression to a constant and serializes its // magnitude as a typed Literal proto. Spark encodes PRECEDING/FOLLOWING via // the sign of the literal (negative => PRECEDING, positive => FOLLOWING), diff --git a/spark/src/test/resources/sql-tests/windows/nested_float_order_keys_strict.sql b/spark/src/test/resources/sql-tests/windows/nested_float_order_keys_strict.sql index 713b0cd4bf6..995470b2273 100644 --- a/spark/src/test/resources/sql-tests/windows/nested_float_order_keys_strict.sql +++ b/spark/src/test/resources/sql-tests/windows/nested_float_order_keys_strict.sql @@ -101,3 +101,25 @@ SELECT id FROM ( DENSE_RANK() OVER (PARTITION BY g ORDER BY named_struct('x', coalesce(IF(s, -f, f), 0.0F)) DESC) AS r FROM nested_float_strict ) WHERE r <= 1 + +-- DataFusion cannot compare an array of structs or arrays, or a struct holding an array, to find +-- the CURRENT ROW bound of a RANGE frame (apache/datafusion#24937), so a running aggregate over +-- such a key falls back even when the key cannot hold a null. window_functions.sql checks the +-- same without floats. +query expect_fallback(RANGE frame on array> ORDER BY is not supported) +SELECT id, COUNT(id) OVER (ORDER BY array(named_struct('x', coalesce(d, 0.0D))), id) AS running +FROM nested_float_strict + +query expect_fallback(RANGE frame on struct> ORDER BY is not supported) +SELECT id, + SUM(id) OVER (PARTITION BY g + ORDER BY named_struct('a', array(coalesce(IF(s, -f, f), 0.0F))) DESC) AS by_group +FROM nested_float_strict + +-- Ranks and ROWS frames over the same keys stay native +query +SELECT id, + RANK() OVER (PARTITION BY g ORDER BY array(named_struct('x', coalesce(IF(s, -d, d), 0.0D)))) AS r, + SUM(id) OVER (PARTITION BY g ORDER BY array(named_struct('x', coalesce(IF(s, -d, d), 0.0D))), id + ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS running +FROM nested_float_strict diff --git a/spark/src/test/resources/sql-tests/windows/window_functions.sql b/spark/src/test/resources/sql-tests/windows/window_functions.sql index d8fb8a4df24..b56ba4cac01 100644 --- a/spark/src/test/resources/sql-tests/windows/window_functions.sql +++ b/spark/src/test/resources/sql-tests/windows/window_functions.sql @@ -673,6 +673,62 @@ SELECT player, game, score, RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS range_sum FROM scores +-- ============================================================ +-- 5.6: RANGE frame over array and struct ORDER BY keys +-- DataFusion finds a CURRENT ROW bound of a RANGE frame by comparing ORDER BY +-- values. It compares an array's elements and a struct's fields, flattening +-- nested structs, but its comparison rejects anything nested below that, so +-- over an array of arrays or structs, or a struct holding an array, the query +-- fails with "Uncomparable values" +-- (https://github.com/apache/datafusion/issues/24937). Comet falls back for +-- those keys. Ranking functions, ROWS frames and an unbounded RANGE frame +-- never compare values that way and stay native over the same keys. +-- ============================================================ + +query +SELECT player, game, score, + SUM(score) OVER (PARTITION BY player ORDER BY array(score)) AS by_array, + SUM(score) OVER (PARTITION BY player + ORDER BY named_struct('a', named_struct('b', score)) DESC) AS by_struct +FROM scores + +query expect_fallback(RANGE frame on array> ORDER BY is not supported) +SELECT player, game, score, + SUM(score) OVER (PARTITION BY player ORDER BY array(named_struct('x', score))) AS run_sum +FROM scores + +query expect_fallback(RANGE frame on struct> ORDER BY is not supported) +SELECT player, game, score, + COUNT(score) OVER (PARTITION BY player ORDER BY named_struct('a', array(score)) DESC + RANGE BETWEEN CURRENT ROW AND UNBOUNDED FOLLOWING) AS run_count +FROM scores + +query expect_fallback(RANGE frame on array> ORDER BY is not supported) +SELECT player, game, score, + LAST_VALUE(game) OVER (PARTITION BY player ORDER BY array(array(score)), game) AS lv +FROM scores + +query +SELECT player, game, score, + RANK() OVER (PARTITION BY player ORDER BY array(named_struct('x', score))) AS rk, + SUM(score) OVER (PARTITION BY player ORDER BY array(named_struct('x', score)) + RANGE BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING) AS total +FROM scores + +-- Spark gives CUME_DIST a RANGE frame, but DataFusion computes it from peer +-- groups without reading the frame. +query tolerance=1e-6 +SELECT player, game, score, + PERCENT_RANK() OVER (PARTITION BY player ORDER BY named_struct('a', array(score))) AS pr, + CUME_DIST() OVER (PARTITION BY player ORDER BY named_struct('a', array(score))) AS cd +FROM scores + +query +SELECT player, game, score, + SUM(score) OVER (PARTITION BY player ORDER BY named_struct('a', array(score)), game + ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS rows_sum +FROM scores + -- ############################################################ -- Section 6: Other window / aggregate functions -- ############################################################