Skip to content
Open
14 changes: 8 additions & 6 deletions .ai/skills/review-comet-shuffle-pr/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -55,12 +55,14 @@ 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 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
Expand Down
4 changes: 2 additions & 2 deletions docs/source/contributor-guide/native_shuffle.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
33 changes: 19 additions & 14 deletions docs/source/user-guide/latest/compatibility/floating-point.md
Original file line number Diff line number Diff line change
Expand Up @@ -48,20 +48,25 @@ 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.

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

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.

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
Expand Down
22 changes: 8 additions & 14 deletions docs/source/user-guide/latest/compatibility/operators.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)).
Expand All @@ -102,20 +106,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)).

**Known incompatibilities:**

- `ROW_NUMBER` over such a key still runs natively and follows the native sort, which compares
nested floating-point values with Arrow's raw total ordering. Which of two rows that differ only
in `-0.0` and `+0.0`, or in their NaN representation, gets the lower row number can therefore
differ from Spark ([#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), which also covers
strict floating-point mode.

## Round-Robin Partitioning

Expand Down
21 changes: 11 additions & 10 deletions docs/source/user-guide/latest/tuning/operators.md
Original file line number Diff line number Diff line change
Expand Up @@ -149,16 +149,17 @@ 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.

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` sorts array elements rather than ordering rows. It follows Spark's floating-point ordering as well, so it
also stays native with `spark.comet.exec.strictFloatingPoint=true`.
60 changes: 50 additions & 10 deletions native/core/src/execution/planner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -141,9 +141,9 @@ use datafusion_comet_spark_expr::{
create_case_when, create_if_expr, jvm_udf::JvmScalarUdfExpr, 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};
Expand Down Expand Up @@ -996,16 +996,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<Arc<dyn PhysicalExpr>, 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(),
)?)
Expand Down Expand Up @@ -4950,6 +4952,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<FieldRef> = 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();
Expand Down Expand Up @@ -5017,7 +5051,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]),
Expand Down Expand Up @@ -5069,8 +5107,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.
Expand Down
10 changes: 6 additions & 4 deletions spark/src/main/scala/org/apache/comet/CometConf.scala
Original file line number Diff line number Diff line change
Expand Up @@ -1088,10 +1088,12 @@ 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, or sorting floating-point values nested in " +
"arrays, structs, or maps. Scalar `ORDER BY`, window ordering and range partitioning " +
"keys are unaffected, because Comet normalizes those comparison keys to match Spark, " +
"and so is `sort_array`, which follows Spark's ordering. " +
"Spark, such as comparing -0.0 and 0.0. `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. 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. " +
"`sort_array` is unaffected too, because it follows Spark's ordering. " +
s"$COMPAT_GUIDE.")
.booleanConf
.createWithDefault(false)
Expand Down
Loading
Loading