mbutrovich commented on code in PR #6447:
URL: https://github.com/apache/datafusion-comet/pull/6447#discussion_r4190049382
##########
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:
> With `spark.comet.parquet.rowFilterPushdown.enabled=true` they normalize
both sides, so in that mode float comparisons give up statistics pruning
instead of dropping rows (2f003c207).
I checked how much pruning that gives up. I planned `d > 500.0D` through
`create_data_filter` at 2f003c2 and built a DataFusion `PruningPredicate` from
the result. With `pushdown_filters` off, the filter is `d@0 > 500` and the
pruning predicate is `d_null_count@1 != row_count@2 AND d_max@0 > 500`. With it
on, the filter is `FloatNormalize [child: d@0] > 500` and
`PruningPredicate::always_true()` returns `true`. So with the flag on, every
`FLOAT` or `DOUBLE` comparison in a data filter loses row-group statistics,
page index and bloom filter pruning, not only statistics pruning.
Choosing correct rows over pruning makes sense to me. As I read DataFusion
55's `ParquetSource`, one predicate drives both pruning and the row filter, so
keeping pruning in this mode would need a DataFusion change or a separate
pruning predicate. Could you open an issue for that and link it from
`floating-point.md`, so people who turn the flag on know what they lose?
--
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]