andygrove opened a new pull request, #6691:
URL: https://github.com/apache/datafusion-comet/pull/6691
**Stacked on #6447.** Review 494af960d; I'll rebase onto main and mark this
ready once #6447 lands.
## Which issue does this PR close?
Closes #6590.
## Rationale for this change
Since #6447, every native comparison of `FLOAT` or `DOUBLE` operands follows
Spark's SQL ordering, in which `-0.0` equals `0.0`, all NaNs are equal and NaN
sorts above every other value. It gets there by normalizing each operand that
isn't a literal into a copy and running Arrow's kernels on the copies, which
made a comparison slower than DataFusion's own. For arrays and structs, `<=>`,
`<`, `<=`, `>` and `>=` deep-copy both operands in every batch, while nested
`=` and `<>` already compare in place. The issue has the details.
## What changes are included in this PR?
`spark_comparison` now builds a `SparkComparison` for Float32 or Float64
operands of the same type, and for `<`, `<=`, `>`, `>=`, `<=>` and `IS DISTINCT
FROM` on lists and structs with a float leaf. It compares the operands as they
are:
- **Flat floats** go through new kernels in `float_semantics/kernels.rs`,
column against column or against a scalar. Each operator is one branch-free
test over the two value buffers in the style of `float_gt` (`(a > b) |
(a.is_nan() & !b.is_nan())`), and the loops pack 64 results into a word at a
time. Against a scalar, whether the scalar is NaN picks the test before the
loop, so most operators come down to a single IEEE comparison. SQL null rules,
including `<=>` and `IS DISTINCT FROM`, are applied with bitmap operations
afterwards.
- **Lists and structs** use `spark_comparator` for the ordering operators
and `spark_equality` for the null-safe ones, so `<=>` keeps the length check
that tells lists of different lengths apart without comparing their elements.
When either side has nulls, the comparison walks only the rows where neither
side is null, through the validity bitmap, instead of testing both null buffers
on every row. Nested `=`, `<>` and `IN` (`NestedPredicate`) get the same
treatment.
Planning doesn't change. A literal is still normalized while the plan is
built. `FloatOperands::Raw` still leaves a float column compared with a literal
to a plain `BinaryExpr` in a scan's data filters, so Parquet pruning recognizes
it. Operands of other types, such as dictionary-encoded floats, keep the
normalize-then-compare path. An operand whose runtime layout doesn't match its
type, which the planner never produces, falls back to normalizing and comparing
with Arrow rather than failing.
The floating-point compatibility page no longer says the operands are
normalized. The `float_comparison` bench now covers `<`, `=` and `<=>`, and
runs a `normalized` engine that builds #6447's expression, so the two can be
compared within one run. The `nested_comparison` bench gains `lt` and
`not_distinct` modes.
## How are these changes tested?
- **Rust unit tests:**
- The kernels are checked against `compare_floats` on every pair of edge
values (both zeros, canonical NaN, a sign-bit NaN, a NaN with a payload, ±1,
±Infinity and NULL), for every operator, Float32 and Float64, column against
column and against a scalar. The inputs span whole and partial words and are
sliced at several offsets.
- Through `spark_comparison`, sliced columns are compared with scalars
that only appear at evaluation time, so planning can't normalize them, as well
as scalar against scalar.
- Nested comparisons are tested against list and struct literals on either
side, against a null literal and from an offset.
- Further tests cover the null-skipping comparator and the layout fallback.
- The existing `float_operands_follow_spark_ordering` and
`nested_float_operands_follow_spark_ordering` now run on the new path.
- Dropping the NaN term from `<=` or from `>` fails four of these tests.
- **Differential test (not committed):** a throwaway test compared the new
path with #6447's normalize-then-compare on randomized inputs, with identical
results in every case:
- 19,200 flat comparisons: arrays of up to 1,000 rows, Float32 and
Float64, sliced, with nulls, all eight operators, and column, literal and
runtime-scalar operands.
- 12,000 nested comparisons: lists, structs, lists of lists and lists of
structs, with outer and inner nulls and values under null slots.
- **Rust suites:** all 1,132 `spark-expr` tests and the 55 core planner
tests pass, workspace clippy (`--all-targets -D warnings`) passes, and `cargo
fmt` is clean.
- **JVM, Spark 4.1 (the default profile):**
- `CometFloatSemanticsSuite` and the `CometNativeReaderSuite` row-group
pruning tests, including "row-group statistics pruning fires for a
floating-point comparison": 415 passed.
- `CometSqlFileTestSuite` for `expressions/conditional/`,
`expressions/math/`, `expressions/aggregate/`, `expressions/array/`, `windows/`
and `join/`: 233 tests over 189 fixtures passed, including
`float_comparisons.sql`, `nan_divisor.sql` and `arithmetic.sql`.
- **Not run locally:** the other Spark profiles.
### Benchmarks
These were measured in release mode on an M3 Max that other builds were also
using. For "before", the same bench sources were built against #6447's head
f2c2fc4e5, and the two builds' runs were interleaved.
Flat comparisons, per batch of 8,192 doubles, median of three runs:
| Comparison | DataFusion `BinaryExpr` | #6447 | This PR |
| --- | ---: | ---: | ---: |
| column `<` column | 7.34 µs | 11.28 µs | 2.73 µs |
| column `<` literal | 4.16 µs | 6.06 µs | 1.20 µs |
| column `=` column | 5.83 µs | 9.74 µs | 2.95 µs |
| column `=` literal | 3.37 µs | 5.29 µs | 1.19 µs |
| column `<=>` column | 5.99 µs | 9.94 µs | 2.99 µs |
| column `<=>` literal | 3.42 µs | 5.41 µs | 1.17 µs |
That's 3.3x to 5.1x faster than #6447. With `-0.0`, NaN or a sign-bit NaN in
every tenth row, the numbers stay within a few percent. In #6447's build, the
`normalized` engine matches `spark_comparison` to within 1%.
The DataFusion column comes from #6447's build. In this PR's build,
DataFusion's `BinaryExpr` measured 40-60% slower on the same code:
`normalize_float_zero` scans each column for `-0.0` with an `any()` loop that
doesn't vectorize, and how fast that runs depended on the build. It's the same
in both builds once the data holds a `-0.0`, because the scan then stops early.
Nested comparisons of 8,192 rows of `ARRAY<DOUBLE>`, where the null rows
hold real lists:
| List width | Lists | Null rows | `<`: #6447 | `<`: this PR | `<=>`: #6447
| `<=>`: this PR |
| ---: | --- | --- | ---: | ---: | ---: | ---: |
| 1 | differ in the first element | none | 40.2 µs | 26.5 µs | 40.8 µs |
29.8 µs |
| 1 | differ in the first element | 1 in 16 | 36.5 µs | 28.9 µs | 36.4 µs |
33.3 µs |
| 1 | differ in the first element | 1 in 2 | 27.7 µs | 16.2 µs | 27.6 µs |
18.4 µs |
| 1 | equal | none | 37.5 µs | 36.5 µs | 35.5 µs | 33.5 µs |
| 1 | equal | 1 in 16 | 37.3 µs | 29.5 µs | 38.0 µs | 35.2 µs |
| 1 | equal | 1 in 2 | 27.8 µs | 16.5 µs | 28.5 µs | 19.5 µs |
| 16 | differ in the first element | none | 95.3 µs | 29.3 µs | 88.7 µs |
31.2 µs |
| 16 | differ in the first element | 1 in 16 | 88.9 µs | 29.7 µs | 89.3 µs |
34.1 µs |
| 16 | differ in the first element | 1 in 2 | 79.8 µs | 16.6 µs | 81.1 µs |
19.2 µs |
| 16 | equal | none | 223.2 µs | 211.3 µs | 223.6 µs | 213.8 µs |
| 16 | equal | 1 in 16 | 222.8 µs | 199.5 µs | 224.7 µs | 204.3 µs |
| 16 | equal | 1 in 2 | 150.9 µs | 107.5 µs | 150.2 µs | 109.0 µs |
| 1024 | differ in the first element | none | 3.69 ms | 73.6 µs | 3.72 ms |
82.6 µs |
| 1024 | differ in the first element | 1 in 16 | 3.76 ms | 70.5 µs | 3.75 ms
| 80.5 µs |
| 1024 | differ in the first element | 1 in 2 | 3.80 ms | 56.6 µs | 3.70 ms
| 63.0 µs |
| 1024 | equal | none | 13.19 ms | 11.87 ms | 13.12 ms | 11.90 ms |
| 1024 | equal | 1 in 16 | 12.49 ms | 11.27 ms | 12.45 ms | 11.29 ms |
| 1024 | equal | 1 in 2 | 8.96 ms | 6.67 ms | 9.07 ms | 6.58 ms |
Nested `=` already compared in place before this PR. Skipping null rows
makes it 1.1x to 1.7x faster when there are nulls. Without nulls its code path
does the same work as before, and it measured between 0.87x and 1.13x of #6447
across widths. A build of this branch that still had #6447's `=` code showed
the same spread against #6447's build, so I read that as build-to-build noise
rather than a change in the comparison.
My first version compared every row and tested the null buffers inside the
comparator. Width-1 lists with sparse nulls were slower than #6447 that way,
and so were dense nulls once the null rows held data. Walking the validity
bitmap fixed both.
To reproduce:
```shell
cargo bench -p datafusion-comet-spark-expr --bench float_comparison
cargo bench -p datafusion-comet-spark-expr --bench nested_comparison --
'nested_comparison/(lt|not_distinct|eq)/'
```
--
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]