andygrove commented on code in PR #6447:
URL: https://github.com/apache/datafusion-comet/pull/6447#discussion_r4190391558
##########
native/spark-expr/src/array_funcs/nested_comparison.rs:
##########
@@ -219,31 +217,83 @@ impl PhysicalExpr for NestedPredicate {
}
}
-/// Build equality after the planner has reconciled nested operand nullability.
+/// How [`spark_comparison`] treats floating-point operands.
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub enum FloatOperands {
+ /// Normalize them, so that the comparison follows Spark's SQL ordering.
+ Normalize,
+ /// Leave a Float32 or Float64 column compared with a literal as it is,
and normalize every
+ /// other operand. Only a scan's pushed-down data filters use this, and
only when the Parquet
+ /// reader prunes with them but does not filter rows: Parquet pruning
recognizes a column
+ /// compared with a literal but not a normalized column, and Spark's
Filter above the scan
+ /// applies Spark's semantics to every row. With row-level pushdown the
reader would drop the
+ /// rows such a comparison rejects, including a stored NaN whose bits
differ from the
+ /// normalized literal, so the data filters use
[`FloatOperands::Normalize`] there instead.
+ Raw,
+}
+
+/// Builds a comparison with Spark's SQL ordering for floats, in which `-0.0`
equals `0.0`, all
+/// NaNs are equal and NaN sorts above every other value, at any depth of a
list or struct.
+///
+/// Arrow compares floats by IEEE 754 total order instead, so each float
operand is normalized
+/// first with [`normalize_comparison_operand`], after which the two orders
agree. Nested `=` and
+/// `<>` compare with `spark_equality` instead, without building normalized
copies of the nested
+/// values. Any other operator, such as `AND`, builds a plain [`BinaryExpr`].
+///
+/// The planner reconciles the nullability of nested operands before calling
this.
pub fn spark_comparison(
left: Arc<dyn PhysicalExpr>,
op: Operator,
right: Arc<dyn PhysicalExpr>,
schema: &Schema,
+ float_operands: FloatOperands,
) -> Result<Arc<dyn PhysicalExpr>> {
+ use Operator::*;
+ if !matches!(
+ op,
+ Eq | NotEq | Lt | LtEq | Gt | GtEq | IsDistinctFrom | IsNotDistinctFrom
+ ) {
+ return Ok(Arc::new(BinaryExpr::new(left, op, right)));
+ }
// An operand whose type does not resolve against this schema falls back
to the plain
// comparison, the way `reconcile_nested_comparison_types` already leaves
such operands alone.
- let nested = matches!(op, Operator::Eq | Operator::NotEq)
- && match (left.data_type(schema), right.data_type(schema)) {
- (Ok(lt), Ok(_)) => needs_spark_equality(<),
- _ => false,
- };
- if nested {
+ let (Ok(left_type), Ok(_)) = (left.data_type(schema),
right.data_type(schema)) else {
+ return Ok(Arc::new(BinaryExpr::new(left, op, right)));
+ };
+ if matches!(op, Eq | NotEq) && is_nested_with_float_leaf(&left_type) {
validate_types(&left, std::slice::from_ref(&right), schema)?;
- Ok(Arc::new(NestedPredicate {
+ return Ok(Arc::new(NestedPredicate {
value: left,
candidates: vec![right],
- negated: op == Operator::NotEq,
+ negated: op == NotEq,
membership: false,
- }))
- } else {
- Ok(Arc::new(BinaryExpr::new(left, op, right)))
+ }));
}
+ let raw = float_operands == FloatOperands::Raw;
+ let (left, right) = if raw && is_float_column(&left, schema) &&
is_literal(&right) {
+ (left, normalize_comparison_operand(right, schema)?)
+ } else if raw && is_literal(&left) && is_float_column(&right, schema) {
+ (normalize_comparison_operand(left, schema)?, right)
Review Comment:
Good catch. I went with expanding it (449059ca7). With `FloatOperands::Raw`,
`=` against either zero is now `d = -0.0 OR d = 0.0`, which `LiteralGuarantee`
turns into one guarantee holding both zeros, so the bloom filter is probed for
each. The other operators keep the literal as it is. DataFusion 55's
`apply_cmp` folds `-0.0` before it compares, so statistics and page index
pruning already treat the two zeros as equal, and only `=` gives the bloom
filter anything to probe. I checked that on Spark-written files, whose
statistics aren't widened (a row group holding only `-0.0` has `min = max =
-0.0`): `=`, `<=>`, `<`, `<=`, `>` and `>=` against a zero returned the same
rows as Spark with filter pushdown off, with and without bloom filters.
A NaN literal had the same problem, since no literal has the bits of every
NaN a file can hold, and arrow-rs leaves NaNs out of the statistics. So a NaN
literal now normalizes the column too, as the `IN` path already does, and gives
up pruning.
The new `CometNativeReaderSuite` test writes two files with bloom filters on
`d`, one holding `-0.0` and the other `0.0`. It checks that `d = -0.0D` and `d
= 0.0D` each return both rows with both bloom filters matching, and that `d =
0.5D` is still pruned by both. It writes out the expected rows instead of
comparing with Spark's, because Spark's own reader probes the bloom filter with
the literal's bits too and skips the file holding the other zero. There's also
a `parquet_exec` test along the lines of your repro, which fails if `=` probes
for a single zero.
##########
native/spark-expr/src/array_funcs/nested_comparison.rs:
##########
@@ -219,31 +217,83 @@ impl PhysicalExpr for NestedPredicate {
}
}
-/// Build equality after the planner has reconciled nested operand nullability.
+/// How [`spark_comparison`] treats floating-point operands.
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub enum FloatOperands {
+ /// Normalize them, so that the comparison follows Spark's SQL ordering.
+ Normalize,
+ /// Leave a Float32 or Float64 column compared with a literal as it is,
and normalize every
+ /// other operand. Only a scan's pushed-down data filters use this:
Parquet pruning recognizes
+ /// a column compared with a literal but not a normalized column. With
row-level pushdown the
+ /// reader also evaluates the filters on each row, and a row it drops
never reaches Spark's
+ /// Filter above the scan, so any other shape, which pruning cannot use
anyway, is normalized.
+ /// A computed operand such as `-d` can hold a NaN with the sign bit set,
which a raw
+ /// comparison sorts below every other value.
+ Raw,
+}
+
+/// Builds a comparison with Spark's SQL ordering for floats, in which `-0.0`
equals `0.0`, all
+/// NaNs are equal and NaN sorts above every other value, at any depth of a
list or struct.
+///
+/// Arrow compares floats by IEEE 754 total order instead, so each float
operand is normalized
+/// first with [`normalize_comparison_operand`], after which the two orders
agree. Nested `=` and
+/// `<>` compare with `spark_equality` instead, without building normalized
copies of the nested
+/// values. Any other operator, such as `AND`, builds a plain [`BinaryExpr`].
+///
+/// The planner reconciles the nullability of nested operands before calling
this.
pub fn spark_comparison(
left: Arc<dyn PhysicalExpr>,
op: Operator,
right: Arc<dyn PhysicalExpr>,
schema: &Schema,
+ float_operands: FloatOperands,
) -> Result<Arc<dyn PhysicalExpr>> {
+ use Operator::*;
+ if !matches!(
+ op,
+ Eq | NotEq | Lt | LtEq | Gt | GtEq | IsDistinctFrom | IsNotDistinctFrom
+ ) {
+ return Ok(Arc::new(BinaryExpr::new(left, op, right)));
+ }
// An operand whose type does not resolve against this schema falls back
to the plain
// comparison, the way `reconcile_nested_comparison_types` already leaves
such operands alone.
- let nested = matches!(op, Operator::Eq | Operator::NotEq)
- && match (left.data_type(schema), right.data_type(schema)) {
- (Ok(lt), Ok(_)) => needs_spark_equality(<),
- _ => false,
- };
- if nested {
+ let (Ok(left_type), Ok(_)) = (left.data_type(schema),
right.data_type(schema)) else {
+ return Ok(Arc::new(BinaryExpr::new(left, op, right)));
+ };
+ if matches!(op, Eq | NotEq) && is_nested_with_float_leaf(&left_type) {
validate_types(&left, std::slice::from_ref(&right), schema)?;
- Ok(Arc::new(NestedPredicate {
+ return Ok(Arc::new(NestedPredicate {
value: left,
candidates: vec![right],
- negated: op == Operator::NotEq,
+ negated: op == NotEq,
membership: false,
- }))
- } else {
- Ok(Arc::new(BinaryExpr::new(left, op, right)))
+ }));
}
+ let raw = float_operands == FloatOperands::Raw;
+ let (left, right) = if raw && is_float_column(&left, schema) &&
is_literal(&right) {
+ (left, normalize_comparison_operand(right, schema)?)
Review Comment:
Filed #6702 and linked it from `floating-point.md`. The issue also notes a
partial fix that needs no DataFusion change: `=`, `<>`, `<=>` and `IS DISTINCT
FROM` against a constant other than NaN could keep the raw column with
row-level pushdown too, since DataFusion compares the two zeros as equal and no
NaN equals such a constant. That would bring back pruning for equality, but not
for `<`, `<=`, `>` and `>=`.
##########
docs/source/user-guide/latest/compatibility/floating-point.md:
##########
@@ -23,20 +23,34 @@ Spark normalizes NaN and zero for floating point numbers
for several cases. See
However, one exception is comparison. Spark does not normalize NaN and zero
when comparing values
because they are handled well in Spark (e.g.,
`SQLOrderingUtil.compareFloats`). But the comparison
functions of arrow-rs used by DataFusion do not normalize NaN and zero (e.g.,
[arrow::compute::kernels::cmp::eq](https://docs.rs/arrow/latest/arrow/compute/kernels/cmp/fn.eq.html#)).
-For top-level `FLOAT` and `DOUBLE` comparisons, Comet normalizes both operands
before native
-execution, including noncanonical NaN literals. Top-level `IN`, `InSet`, and
`NOT IN` membership
-also normalize dynamic candidates and lists containing NaN. When every
candidate is a non-NaN
-literal, Comet keeps DataFusion's static filter and pruning path, enumerating
both signed-zero
-forms when a list contains zero.
+For `FLOAT` and `DOUBLE` comparisons (`=`, `<>`, `<=>`, `<`, `<=`, `>` and
`>=`), Comet
+normalizes both operands before native execution, including noncanonical NaN
literals. This
+applies wherever a comparison appears: projections, filters, aggregate
arguments and `FILTER`
+clauses, join conditions, sort keys, and generator arguments.
+
+A native Parquet scan skips row groups with its data filters. So that
statistics pruning still
+applies, a `FLOAT` or `DOUBLE` column compared with a constant in a data
filter is compared
+without normalizing the column, and the filter above the scan evaluates the
comparison again with
+Spark's semantics. Every other comparison in a data filter is normalized. With
+`spark.comet.parquet.rowFilterPushdown.enabled=true` the scan also filters
rows with these
+comparisons, so a noncanonical NaN stored in the file, such as one with the
sign bit set, can be
+filtered differently from Spark when compared with a constant. Spark's Parquet
writer only writes
+canonical NaNs.
Review Comment:
Thanks, I took this wording in 449059ca7, with two additions for what
changed there: the column is left raw only against a constant other than NaN,
and `=` against either zero checks bloom filters for both zeros. The last
sentence links #6702.
--
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]