andygrove opened a new pull request, #6457:
URL: https://github.com/apache/datafusion-comet/pull/6457
## Which issue does this PR close?
Part of #6385: the `min`/`max` and `greatest`/`least` part of the "fix
`min`/`max`, `greatest`/`least`, `array_remove`, NaN handling in
`array_distinct`/`array_union`, and `sort_array`" step.
## Rationale for this change
Spark orders floats with `SQLOrderingUtil.compareDoubles`: NaN is larger
than every other value and `-0.0` equals `0.0`. Its `max` updates the buffer
with `greatest(max, input)`, and `greatest` replaces its result only with a
strictly greater value, so of equal values the first one wins. `min` and
`least` are the same the other way round. Comet ran DataFusion's versions of
all four, which differ in three ways:
- DataFusion's batch kernels order floats by IEEE 754 total order, where a
NaN with the sign bit set is the smallest value. On x86-64 every NaN that
arithmetic produces has that bit set, so `max(sqrt(x))` over a column with a
negative value returned a number where Spark returns NaN.
- DataFusion's grouped `max` starts each group at `f64::MIN` rather than
`-Infinity`, and `min` at `f64::MAX`. A group whose only value is `-Infinity`
therefore got `max = -1.7976931348623157E308`, and an `Infinity` group got that
value for `min`, on every platform.
- DataFusion's `greatest` and `least` also use total order, fold constant
arguments before the others, and let a later argument win a tie, so
`greatest(-0.0D, 0.0D)` returned `0.0` where Spark returns `-0.0`.
On `main`, where `-d` turns a NaN into a NaN with the sign bit set on any
platform:
| Query | Spark | Comet before |
| --- | --- | --- |
| `SELECT max(-d), min(-d) FROM t` (`d` holds 1.0, NaN, -1.0, ...) | `NaN,
-Infinity` | `Infinity, NaN` |
| `SELECT g, max(d), min(d) FROM t GROUP BY g`, group of only `-Infinity` |
`-Infinity, -Infinity` | `-1.7976931348623157E308, -Infinity` |
| same, group of only `Infinity` | `Infinity, Infinity` | `Infinity,
1.7976931348623157E308` |
| `SELECT greatest(a, b), least(a, b)` for `(0.0, -0.0)` and `(-0.0, 0.0)` |
`0.0, 0.0` and `-0.0, -0.0` | `0.0, -0.0` and `0.0, -0.0` |
Because the native `min` and `max` did not match Spark,
`spark.comet.exec.strictFloatingPoint=true` made them fall back for float
inputs (#2448).
## What changes are included in this PR?
- A `SparkMinMax` aggregate for `FLOAT` and `DOUBLE` that keeps the first of
equal values. The planner uses it for float `min`/`max` in aggregates and in
windows; every other type keeps DataFusion's.
- The ungrouped accumulator folds each batch in order with
`float_gt`/`float_lt` from `float_semantics`.
- The grouped accumulator stores a group's first value as is and replaces
it only with a strictly greater (or smaller) one, so it needs no starting
value. Its loop indexes the groups without bounds checks, relying on the
`GroupsAccumulator` guarantee that every group index is below
`total_num_groups`, as DataFusion's grouped `max` does; the checks cost a third
of the loop.
- The window accumulator supports sliding frames with a monotonic deque
that keeps equal values rather than dropping the older one, so its front is
always the first of the best values in the frame.
- The grouped and sliding loops test with `replaces_max`/`replaces_min`,
the same relation as `float_gt`/`float_lt` written so that one comparison
decides the usual answer. That form is faster in loops that cannot vectorize,
while `float_gt` is faster in the ungrouped fold and in
`array_min`/`array_max`, which do. A test checks that the two forms agree on
every pair of edge values.
- `SparkGreatestLeast` for float arguments, and for arrays and structs with
a float leaf. It walks the arguments in order, skipping nulls, and replaces the
result only with a strictly greater (or smaller) value. Floats are picked in
one pass per argument; nested values compare with `spark_comparator`. Arguments
are cast to the result type first, because Arrow compares field names and
nullability and Spark only guarantees the same type up to nullability. Every
other type keeps DataFusion's `greatest` and `least`.
- Strict floating-point mode no longer makes `min` and `max` fall back.
- The floating-point compatibility guide gains a section on these four. It
notes one difference that remains: Spark treats `greatest` and `least` as
commutative when it matches expressions, so a projection holding both
`greatest(a, b)` and `greatest(b, a)` evaluates one of them for both, while
Comet evaluates each, and for zeros of different signs the two can differ.
- The `greatest` and `least` entries in the math expression audit describe
the new routing.
## How are these changes tested?
- `min_max_floating_point.sql` (new, run with strict floating-point mode off
and on): ungrouped and grouped `min`/`max` over `DOUBLE` and `FLOAT` with NaNs,
sign-bit NaNs, both zeros in both orders, and groups holding only `-Infinity`,
only `Infinity`, only NaN or only NULL; and window frames that grow, slide and
cover the whole partition.
- `greatest_least_floating_point.sql` (new): both zeros in both orders,
sign-bit NaNs on either side, literals before, between and after the columns,
three arguments, calls with only literals, and arrays and structs of floats.
Each query uses one argument order, because of the Spark behavior above.
- Both fixtures fail on `main`.
- The `CometAggregateSuite` test that asserted the strict-mode fallback now
checks that `min` and `max` over both zeros run natively and match Spark, over
data written as a single file so both engines read the rows in the same order.
- Unit tests:
- Each accumulator against a fold in Spark's order (`compare_floats`),
compared by bits, so they check which zero and which NaN is kept, over
pseudo-random sequences of edge values: in batches, as two merged partial
states, grouped (with nulls and a filter, and without either) and emitted with
`EmitTo::First`, and over sliding frames of widths 1 to 4 that retract.
- Groups holding only an infinity, and `convert_to_state` with a filter.
- `greatest` and `least` over every combination of three edge values, with
and without nulls, also with a constant as the middle argument, `FLOAT` ties,
constant-only calls, and lists whose element fields differ in name and
nullability.
- `replaces_max`/`replaces_min` against `float_gt`/`float_lt` on every
pair of edge values.
- Letting a tie replace the current value makes the accumulator tests and
the float `greatest`/`least` tests fail.
- Results on macOS aarch64 with the default Spark 4.1 profile:
- Unit tests: 1036 in spark-expr, 563 in core.
- `CometSqlFileTestSuite`: 590 passed, including both new fixtures (the
`min`/`max` one in both strict modes).
- `CometAggregateSuite` and `CometWindowExecSuite`: 194 passed, 2 ignored
(existing metrics tests).
Performance, on an M3 Max with 8192 doubles per batch:
| | Comet with this PR | DataFusion's version, as before |
| --- | --- | --- |
| `max`, ungrouped | 3.2 µs | 5.1 µs |
| `max` into 1024 groups, eight batches | 32.6 µs | 25.0 µs |
| same, 10% NaN | 32.2 µs | 42.9 µs |
| `greatest` of three columns | 4.2 µs | 36.1 µs |
| same, 10% nulls | 19.9 µs | 97.7 µs |
Grouped `max` over data without NaNs is the one slower case, by about 30%,
which is the cost of the extra NaN test and of recording each group's first
value exactly.
--
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]