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]

Reply via email to