feat: run array_contains on float elements natively with Spark's equa… - #6599
zhangfengcdt wants to merge 2 commits into
Conversation
…lity A Comet kernel compares float elements as Spark's genEqual does, with -0.0 equal to 0.0 and all NaNs equal, at any depth, so these calls no longer go through the codegen dispatcher. Closes apache#6520
…loat-native # Conflicts: # docs/source/user-guide/latest/compatibility/floating-point.md
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Float-bearing
array_containsused Spark’s dispatcher by default. Opting into the native kernel exposed incorrect signed-zero and NaN comparisons. - Design approach: Add
spark_array_containsand route float-bearing elements to Spark-compatible equality. - Correctness / compatibility analysis: Reviewed Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. Reproduced two introduced P2 regressions: nested string collations and evaluation of throwing lookup expressions for null arrays.
- Key design decisions: Reusing
spark_equalityand the extracted bitmap-count helper keeps the implementation localized. Scalar lookups scan flattened values; column lookups stop at a match. Performance against the replaced JVM dispatcher was not independently measured. - Implementation sketch: Scala routing selects the registered Rust UDF, with native unit tests, SQL fixtures and compatibility documentation.
- Behavioral changes worth calling out: Compared against
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa. Its relevant dispatcher route matches the base. Native execution and corrected opt-in float equality are intended changes; the two reproduced regressions are unintended. - Suggested improvements: Preserve dispatch for collated nested types and unsafe eager evaluation, with SQL regressions for both cases. Recommend request changes.
Reviewed all 11 changed files at 46faab056ba96f002559ae405a10e5a3362ddcf6 against fe5c41aab1159ee8ea995da90b48061a905e92bf. Routed skills: review-comet-pr and review-comet-expression-pr. Existing discussion was empty.
Exact-head CI at 2026-10-04 16:15 UTC: 18 successful, 6 running, 14 skipped, no failures. Native build and lint passed. Runtime suites remained pending; Spark SQL suites, macOS and benchmark checks were skipped.
Validation: 5 array_contains and 7 array_remove tests passed. Three disposable native probes checked sliced buffers/nulls and reproduced both regressions against Spark 4.1.3 interpreted and generated expressions. Full local Spark SQL execution was unavailable because the assembled classpath lacked KVStore. No end-to-end Comet SQL or performance run was completed. Disposable project tests were removed and the checkout is clean.
| // compares them as Spark does, with -0.0 equal to 0.0 and all NaNs equal, at any depth. | ||
| val elementType = expr.left.dataType.asInstanceOf[ArrayType].elementType | ||
| val function = | ||
| if (SupportLevel.containsType(elementType, classOf[FloatType], classOf[DoubleType])) { |
There was a problem hiding this comment.
[P2] Preserve Spark dispatch for nested string collations. On Spark 4.x, let a have type ARRAY<STRUCT<d:DOUBLE,s:STRING COLLATE UTF8_LCASE>> and v the matching struct type. For a=[(0D,'a')] and v=(0D,'A'), array_contains(a,v) must return true. This float-leaf gate now selects spark_array_contains, whose shared comparator compares the string field bytewise and returns false. The previous float fallback evaluated the whole expression through Spark and preserved collation. Please retain CodegenDispatchFallback and guard hasNonDefaultStringCollation(elementType) before declaring this route compatible.
Evidence: Spark 4.1.3 Catalyst ArrayContains returned true through both eval and UnsafeProjection for the collated struct pair. The exact-head Rust UDF returned [Some(false)] for both scalar and column lookup arguments after the existing serde’s conversion to Arrow types. Reproduction files: /tmp/6599-review/CatalystProbe.scala, /tmp/6599-review/review_6599.rs (mixed_struct_string_comparison_probe). Native command: cargo test -p datafusion-comet-spark-expr --test review_6599 --locked --offline -- --nocapture.
There was a problem hiding this comment.
I reproduced this one end to end at 46faab0 on Spark 4.1.3. It doesn't need a collated column, a cast is enough:
SELECT array_contains(
array(named_struct('d', d, 's', CAST(s AS STRING COLLATE UTF8_LCASE))),
named_struct('d', d, 's', CAST(upper(s) AS STRING COLLATE UTF8_LCASE))) FROM tThe plan is CometProject over CometNativeScan, and it returns false for s = 'a' where Spark returns true. With main's serde the same query goes through the dispatcher and matches Spark, so this is new with this PR.
| } else { | ||
| "array_contains" | ||
| } | ||
| scalarFunctionExprToProto(function, arrayExprProto, keyExprProto) |
There was a problem hiding this comment.
[P2] Preserve null-array short-circuiting before evaluating a throwing lookup. With ANSI enabled, consider SELECT array_contains(a, CAST(s AS DOUBLE)) FROM t for a ARRAY<DOUBLE>, s STRING containing (NULL,'bad') and ([1D],'1'). Spark returns NULL,true because it skips the cast when a is null. The new native route evaluates both children eagerly and raises CastInvalidValue for 'bad', aborting a previously successful query. The old float route dispatched the whole expression and preserved Spark’s evaluation order. Please keep such expressions dispatched unless lookup evaluation is proven safe, or implement masked evaluation of the lookup.
Evidence: Spark 4.1.3 interpreted and generated ArrayContains returned null for the null-array/invalid-string row and true for the valid row. An exact-head ScalarFunctionExpr using SparkFloatArrayContains and Comet’s ANSI Cast on the same batch returned Err(External(CastInvalidValue { value: "bad", from_type: "STRING", to_type: "DOUBLE" })). Reproduction files: /tmp/6599-review/CatalystProbe.scala, /tmp/6599-review/review_6599.rs (null_array_throwing_lookup_probe). ScalarFunctionExpr::evaluate evaluates every child before invoking the UDF.
There was a problem hiding this comment.
Confirmed this end to end too. With spark.sql.ansi.enabled=true, array_contains(a, CAST(s AS DOUBLE)) over the rows (NULL, 'bad') and (array(1.0D), '1') fails here with CAST_INVALID_INPUT. On main it goes through the dispatcher and returns NULL, true like Spark. The same query over an ARRAY<INT> already fails on main, so that part is older and fits #6006, but for float arrays this PR is what removes the protection. CometArrayJoin.orderInsensitive in the same file handles this by keeping anything other than a literal or a column read on the dispatcher. Could CometArrayContains do the same for the value argument when the array can be null?
andygrove
left a comment
There was a problem hiding this comment.
I ran this locally at 46faab0 on Spark 4.1.3 and 3.4.3, including the new fixtures, the routing files and all 412 cases of CometFloatSemanticsSuite. The float equality matches Spark everywhere I looked. I'm requesting changes for three regressions against 1.1, where these shapes went through the dispatcher. Two are sunchao's findings, which I reproduced through SQL (details on those threads). The third is performance.
The benchmark in the description compares against datafusion-spark's kernel, but float arrays never ran that by default. Users upgrading from 1.1 get the codegen dispatcher today. I measured both on one release build, switching only arrays.scala (2M rows, 50-element arrays, five array_contains calls per row, end-to-end ms over two interleaved runs):
| case | this PR | dispatcher |
|---|---|---|
ARRAY<DOUBLE>, constant value |
1144, 976 | 1395, 1488 |
ARRAY<DOUBLE>, value per row |
1474, 1328 | 1298, 1494 |
ARRAY<STRUCT<x: DOUBLE, y: INT>>, value per row |
3831, 4060 | 2591, 2739 |
The constant case is clearly faster and the per-row case is even. Struct elements cost about twice what the dispatcher does once the scan is subtracted, and ARRAY<ARRAY<DOUBLE>> was slower in both of my runs too. The per-element spark_equality comparator in contains_where looks like the cost. Could nested element types stay on the dispatcher for now, so only flat FLOAT and DOUBLE arrays route to spark_array_contains? That would also close the collation case, since a collated string can only reach this kernel inside a struct.
A few places still describe the old routing. The Notes cell for array_contains in expressions.md says float arrays go through the dispatcher, and GenerateDocs only rewrites the Implementation column, so it won't fix itself. The comments at lines 61 and 90 of array_contains.sql say the same. The comment on map_contains_key over nested floating-point map keys falls back in CometMapExpressionSuite relies on CometArrayContains reporting Incompatible, and its allowIncompatible=true setting does nothing now. floating-point.md also still says "Both expressions" now that the section covers three. Could you update these here?
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Float-bearing
array_containsused Spark’s dispatcher by default. Opting into native execution exposed incorrect signed-zero and NaN comparisons. - Design approach: Add
spark_array_containswith Spark-compatible float equality and enable it for flat and nested float-bearing elements. - Correctness / compatibility analysis: Checked Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. Float equality and ordinary null handling match. Two existing P2 reports remain reproducible: collated struct comparisons return
falsewhere Spark returnstrue, and throwing lookup expressions execute for null arrays, aborting queries where Spark returns null. - Key design decisions: Reusing
spark_equalityand extractingset_bits_beforekeeps the implementation localized. Constant float lookups scan flattened values, while per-row lookups stop at a match. The helper extraction preservesarray_removebehavior. - Implementation sketch: Scala selects the registered Rust UDF. Native tests, SQL fixtures, routing assertions and compatibility documentation accompany the change.
- Behavioral changes worth calling out: Compared with
branch-1.1at992c806a7e38c2e88bd018aa5774164b0850e1fa, native routing and corrected opt-in float equality are intended. The correctness regressions are unintended. The existing performance review also remains unresolved: its struct-array workload measured 3,831–4,060 ms here versus 2,591–2,739 ms with the dispatcher. Those measurements were not rerun in this pass. - Suggested improvements: Address the existing requests to preserve dispatch for nested elements and unsafe eager lookup evaluation, with regression coverage. No additional introduced P1/P2 issues found within this review beyond the existing unresolved reports. The existing request-changes concerns remain valid.
Reviewed all 11 changed files at 46faab056ba96f002559ae405a10e5a3362ddcf6 against fe5c41aab1159ee8ea995da90b48061a905e92bf. Confirmed non-draft status. Read existing reviews and threads. Routed skills: review-comet-pr and review-comet-expression-pr.
Exact-head CI: 25 passed, 15 skipped, no failed or running checks. Rust tests, Linux Spark 4.1 suites and TPC-H/TPC-DS verification passed. Spark SQL suites, macOS and benchmark checks were skipped.
Validation: five array_contains tests, seven array_remove tests and three disposable native probes passed. The probes checked sliced buffers and reproduced both existing correctness regressions. Spark 4.1.3 interpreted and generated Catalyst probes confirmed the expected results. No fresh end-to-end Comet SQL, full Spark SQL suite or performance run was completed. Disposable project tests were removed, and the checkout is clean.
Which issue does this PR close?
Closes #6520
Rationale for this change
array_containson an array whose elements hold aFLOATorDOUBLE, at any depth, reportsIncompatible, because the native path compares elements by their bits while Spark'sgenEqualtreats-0.0and0.0as equal and all NaNs as equal. By default those calls go through the JVM codegen dispatcher, which is correct but leaves the native path. WithallowIncompatible=truethey get the bitwise kernel and the wrong answer.This adds a native kernel with Spark's equality, built the way
spark_array_removewas in #6518, so these calls stay native and reportCompatible.What changes are included in this PR?
spark_array_contains, a Comet kernel for arrays with a float leaf at any depth. A constant float value is one test per element in a single pass over all the values. A value per row uses Spark's float comparison. Nested array and struct elements go through the sharedspark_equality.CometArrayContainsroutes float-element arrays to the kernel and reports themCompatible. Other element types keep the current path through datafusion-spark'sarray_contains. TheIncompatiblereason and the dispatcher fallback are removed.array_removeused for its offsets moves to the module so both kernels share it.array_remove's behavior is unchanged.array_contains, and the audit entry for it is corrected.Benchmark, Comet kernel versus the native path other element types take, 8,192 rows, medians:
This is only for fixing
array_contains.array_remove,array_positionandarray_contains, which the issue mentions, is left for a follow-up.How are these changes tested?
array_contains_floating_point.sql, with every query markedexpect_native, covers signed zeros, NaNs, null elements, a null array or value, a value per row, nested arrays and structs, and literals.array_containsbenchmark cases are added tofloat_arrays.