andygrove opened a new pull request, #6447:
URL: https://github.com/apache/datafusion-comet/pull/6447
## Which issue does this PR close?
Closes #6157.
Part of #6385: the "normalize comparisons in the native comparison builder
and remove `CometExecRule.normalize`" step. The divisor half of
`CometExecRule.normalize` stays until #6519 is fixed.
This PR also covers the last item on #6385's checklist, the contributor
guide page on which Spark rule an expression follows and how to test for it. It
doesn't close the epic.
## Rationale for this change
Spark compares floats with `SQLOrderingUtil.compareDoubles`: `-0.0` equals
`0.0`, all NaNs are equal, and NaN sorts above every other value. Arrow's
comparison kernels use IEEE 754 total order instead. `CometExecRule.normalize`
bridged the two by wrapping comparison operands in `NormalizeNaNAndZero`, but
only inside `ProjectExec` and `FilterExec`, so a comparison anywhere else
compared raw Arrow values. DataFusion 55 folds `-0.0` in scalar comparisons,
which fixed the zero cases, but a NaN with the sign bit set still sorts below
every other value, and on x86-64 every NaN that arithmetic produces has that
bit set. From #6385, where `-d` turns the NaN in row 3 into a sign-bit NaN on
any platform:
| Query | Spark | Comet before |
| --- | --- | --- |
| `SELECT sum(if(-d > 0.0, 1, 0)) FROM t` | `2` | `1` |
| `SELECT count(*) FILTER (WHERE -d >= d) FROM t` | `4` | `3` |
| `SELECT a.id FROM t a JOIN t b ON a.id = b.id AND -a.d >= b.d` | `1, 2, 3,
5` | `1, 2, 5` |
| `SELECT b.id FROM t a JOIN t b ON -a.d > b.d WHERE a.id = 3` | `1, 2, 4,
5` | (no rows) |
| `SELECT id FROM t ORDER BY -d > 0.0, id` | `1, 2, 4, 3, 5` | `1, 2, 3, 4,
5` |
| `SELECT id, x FROM t LATERAL VIEW explode(array(-d > 0.0)) v AS x WHERE id
= 3` | `3, true` | `3, false` |
Arrays and structs had the same gap for `<=>`, `<`, `<=`, `>` and `>=` in
every operator, Project and Filter included, because the rewrite only wrapped
scalar operands (#6157).
## What changes are included in this PR?
- `spark_comparison`, which builds every native comparison, now normalizes
float operands for all eight comparison operators: a `FLOAT` or `DOUBLE`
operand through `NormalizeNaNAndZero`, and an array or struct with a float leaf
through `NormalizeNestedFloats`. Once `-0.0` is folded and NaN canonicalized,
Arrow's total order agrees with Spark's. A literal operand is normalized while
the plan is built, so it stays a literal. Nested `=` and `<>` keep their
`spark_equality` path, which compares without copying.
- The data filters that a native scan pushes into the Parquet reader depend
on row-level pushdown (`spark.comet.parquet.rowFilterPushdown.enabled`, off by
default). The Filter above the scan still evaluates every data filter with
Spark's semantics.
- With it off, the reader only prunes with them. A `FLOAT` or `DOUBLE`
column compared with a literal other than NaN keeps both as they are, because
DataFusion's pruning only recognizes a column compared with a literal. A bloom
filter probe hashes the literal's bits, so `=` against either zero becomes `=
-0.0 OR = 0.0`. A NaN literal, and every other comparison, is normalized: a
stored NaN can have other bits than the literal, and a writer can leave NaNs
out of the column statistics.
- With it on, the reader drops the rows a data filter rejects, and a raw
`-d > 0.0` dropped the NaN row that Spark keeps. So every operand is
normalized, and float comparisons in data filters give up pruning in that mode
(#6702).
- `PhysicalPlanner` becomes `Clone` so that `create_data_filter` can plan
the filters with a copy set to `FloatOperands::Raw`.
- `CometExecRule.normalize` no longer wraps comparisons. What remains is
`normalizeDivisors`, which still wraps the divisor of a floating-point `Divide`
or `Remainder` in a Project or Filter. Native comparisons don't need it any
more, but it makes the NaN quotient of `1.0D / (-d)` canonical, and
`percentile_approx` still orders a sign-bit NaN below every other value
(#6519). Because of it the TPC-DS plans keep main's expression counts, so no
golden file changes.
- `NormalizeNaNAndZero` keeps a scalar child a scalar instead of expanding
it into a column.
- `is_nested_with_float_leaf` replaces `needs_spark_equality`, so the
comparison builder and the normalizer share one test for a list or struct with
a float leaf.
- A `float_comparison` benchmark times `spark_comparison` against
DataFusion's `BinaryExpr`, and `normalize_floats` on its own.
- A new contributor guide page, `contributor-guide/floating_point.md`, under
Project Architecture after Timezone Handling, and a testing tip in
`adding_a_new_expression.md` that points to it. Most of the fixes under #6385
came from the same two gaps: a Spark function inherits its float equality from
whichever Java or Scala API its implementation calls, and there are six of
them; and a test that only uses canonical NaNs passes on Apple Silicon while
the same query fails on x86-64, where arithmetic produces NaNs with the sign
bit set. The page covers:
- Spark's six rules and which functions use them, and the releases in
which `collect_set`, `OpenHashSet` and `OpenHashMap` (and with them `mode`),
`array_distinct`/`array_union` and map construction changed rules.
- How Arrow and DataFusion differ.
- Where Comet applies each rule: the `float_semantics` helpers, key
normalization in the planner, the comparison builder, `IN`, hashing, and the
Scala version policies.
- How to choose the rule for a new expression, guidelines for native code,
and how to test, including how to make a sign-bit NaN that survives a Parquet
round trip and how `CometFloatSemanticsSuite`'s known gaps work.
- The floating-point compatibility guide describes the new behavior. It also
says that Spark's own Parquet reader keeps `-0.0` and `0.0` apart in its
dictionary and bloom filters, so for a comparison with a zero it can skip a row
group that Comet reads, and Comet can return rows that Spark does not. The
`spark.comet.parquet.rowFilterPushdown.enabled` doc and the scan tuning page
say that float comparisons in pushed filters don't prune while that flag is on.
## How are these changes tested?
- A new `float_comparisons.sql` fixture compares `DOUBLE` and `FLOAT` edge
values, including negated NaNs, with every operator, `IS DISTINCT FROM`
included, against columns and against literals on either side. It covers
Project, Filter, an aggregate argument, a `FILTER` clause, a grouping key,
broadcast hash, shuffled hash, sort-merge and nested loop join conditions, a
sort key and `explode`. It also compares arrays and structs of floats under all
six operators and `IS DISTINCT FROM` (#6157), one level deep and as
`ARRAY<STRUCT<v: DOUBLE>>` and `STRUCT<a: ARRAY<DOUBLE>>` columns, and filters
stored edge values with constants. The file runs with row-level pushdown off
and on; with it on, `WHERE -d > 0.0D` loses the NaN row if data filters leave
computed operands raw. On `main` it fails at the first query outside Project
and Filter.
- `arithmetic.sql` divides and takes the remainder by `-0.0D` and
`double('-0.0')`.
- A new `nan_divisor.sql` fixture takes `hash`, `xxhash64`, `least`,
`greatest`, `max`, `min` and `percentile_approx` of `1.0D / (-d)`, `1.0D %
(-d)` and `1.0F % (-f)` for a stored NaN, with ANSI off and on. Without
`normalizeDivisors`, `percentile_approx` returns `0.0` for the maximum where
Spark returns NaN.
- `CometNativeReaderSuite`:
- Row-group statistics pruning still fires for `d > 500.0D` on a `DOUBLE`
column, and a row group holding a NaN, whose other values are all below 500,
still returns the NaN row.
- With bloom filters on `d`, one file holding `-0.0` and another holding
`0.0`, `d = -0.0D` and `d = 0.0D` each return both rows with both bloom filters
matching, and `d = 0.5D` is pruned by both bloom filters. Spark's own reader
probes the bloom filter with the literal's bits and skips the file holding the
other zero, so the test writes out the expected rows instead of comparing with
Spark's.
- Native tests:
- A `parquet_exec` test scans an arrow-rs file with a bloom filter through
`init_datasource_exec`. The data filter for `d = -0.0` and for `d = 0.0` reads
the row group holding `-0.0` and no `0.0`, and `d = 0.5` is pruned by the bloom
filter. It fails if `=` keeps a single zero.
- Planner tests check that `create_data_filter` leaves the column
unwrapped where `create_expr` wraps it, and that with row-level pushdown it
normalizes the column, so `d = -double('NaN')` and `d > 0.0` keep a stored
sign-bit NaN.
- `spark_comparison` on every pair of edge values (both zeros, canonical,
sign-bit and payload NaNs, infinities and null) under all eight operators, as
two columns and as a column and a literal on either side, for `Float64` and
`Float32`, checked against `compare_floats`.
- `<=>`, `<`, `<=`, `>` and `>=` on lists, including lists of different
lengths, and on structs, checked against `spark_comparator`.
- Literal folding, wrapping each operand once, `FloatOperands::Raw`
leaving a column and a literal as they are on either side, `=` against either
zero giving the bloom filter both zeros, a NaN literal normalizing the column,
and non-comparison operators left alone.
- `NormalizeNaNAndZero` keeping scalars.
- The contributor guide page is docs only. prettier 3.9.9 `--check` passes
on every changed Markdown file, and the files, helpers and fixtures the page
names exist at this commit. Runs on Spark 3.4.3 and 4.1.3 confirm what the two
guides say about Spark:
- `size(collect_set(d))` is 3 for a group holding `-0.0`, `0.0` and two
NaNs, and 4 for three NaNs and two `1.0`s.
- `mode` over three NaNs and two `1.0`s returns `1.0` on 3.4.3, whose
`OpenHashMap` never matches two NaNs, and NaN on 4.1.3.
- On dictionary-encoded files, Spark's `d = -0.0D`, `d <=> -0.0D`, `d IN
(-0.0D, 5.0D)`, `d >= 0.0D` and `d <= -0.0D` miss rows holding the other zero,
with either Parquet reader. Spark returns them with
`parquet.filter.dictionary.enabled=false` or
`spark.sql.parquet.filterPushdown=false`, and Comet returns them in every case.
- Results on macOS aarch64 with the default Spark 4.1 profile, at 449059ca7
(the commits after it change only docs):
- Unit tests: 1122 in spark-expr, 665 in core. clippy and rustfmt pass.
- `CometFloatSemanticsSuite` and the SQL fixtures under `expressions/`
`conditional`, `math`, `aggregate`, `array` and `hash`, `windows`, `join` and
`operators`: 653.
- `CometNativeReaderSuite`, `CometExecRuleSuite` and TPC-DS plan stability
v1.4 and v2.7: 314, plus one test that cancels itself on this Spark version.
- `CometAggregateSuite`, `CometInMemoryCacheSuite` and
`CometTypedDatasetSuite`, which commits merged from main changed: 222.
Performance: comparisons in Project and Filter cost what they did before,
since the rewrite already wrapped their operands in the same
`NormalizeNaNAndZero`. Elsewhere a comparison now pays for normalizing its
operands. From `benches/float_comparison.rs` on an M3 Max, a batch of 8192
doubles takes 10.4 µs for a column compared with a column, against 6.8 µs for
DataFusion's comparison, and 5.6 µs for a column compared with a literal,
against 3.8 µs. Skipping the copy when an operand holds no NaN or `-0.0` would
save at most about 8%, because checking the values costs about as much as
copying them. A kernel that compares in Spark's order directly, without copying
the operands, could win the difference back later; that follow-up is #6590.
--
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]