Skip to content

feat: run array_contains on float elements natively with Spark's equa… - #6599

Open
zhangfengcdt wants to merge 2 commits into
apache:mainfrom
zhangfengcdt:feat/array-contains-float-native
Open

zhangfengcdt wants to merge 2 commits into
apache:mainfrom
zhangfengcdt:feat/array-contains-float-native

Conversation

@zhangfengcdt

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6520

Rationale for this change

array_contains on an array whose elements hold a FLOAT or DOUBLE, at any depth, reports Incompatible, because the native path compares elements by their bits while Spark's genEqual treats -0.0 and 0.0 as equal and all NaNs as equal. By default those calls go through the JVM codegen dispatcher, which is correct but leaves the native path. With allowIncompatible=true they get the bitwise kernel and the wrong answer.

This adds a native kernel with Spark's equality, built the way spark_array_remove was in #6518, so these calls stay native and report Compatible.

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 shared spark_equality.
  • CometArrayContains routes float-element arrays to the kernel and reports them Compatible. Other element types keep the current path through datafusion-spark's array_contains. The Incompatible reason and the dispatcher fallback are removed.
  • The one-pass prefix bit count that array_remove used for its offsets moves to the module so both kernels share it. array_remove's behavior is unchanged.
  • The float semantics page covers 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:

Case len 8 len 8, nulls len 50 len 50, nulls
Constant value 27 vs 31 µs 35 vs 67 µs 60 vs 67 µs 76 vs 120 µs
Value per row 18 vs 30 µs 32 vs 116 µs 18 vs 52 µs 58 vs 242 µs

This is only for fixing array_contains. array_remove, array_position and array_contains, which the issue mentions, is left for a follow-up.

How are these changes tested?

  • Unit tests for the kernel: every edge value looked up in rows of edge values through both the constant and per-row paths, null rows, null values, empty arrays, scalar inputs, nested array elements, and struct elements.
  • array_contains_floating_point.sql, with every query marked expect_native, covers signed zeros, NaNs, null elements, a null array or value, a value per row, nested arrays and structs, and literals.
  • The two routing files now expect native execution on the float column in strict floating-point mode, with and without the dispatcher. They failed before this change.
  • The array_contains benchmark cases are added to float_arrays.

…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
@github-actions github-actions Bot added enhancement New feature or request area:expressions Expression evaluation labels Oct 4, 2026
…loat-native

# Conflicts:
#	docs/source/user-guide/latest/compatibility/floating-point.md

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Float-bearing array_contains used Spark’s dispatcher by default. Opting into the native kernel exposed incorrect signed-zero and NaN comparisons.
  • Design approach: Add spark_array_contains and 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_equality and 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.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa. 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])) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[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.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 t

The 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)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[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.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Float-bearing array_contains used Spark’s dispatcher by default. Opting into native execution exposed incorrect signed-zero and NaN comparisons.
  • Design approach: Add spark_array_contains with 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 false where Spark returns true, and throwing lookup expressions execute for null arrays, aborting queries where Spark returns null.
  • Key design decisions: Reusing spark_equality and extracting set_bits_before keeps the implementation localized. Constant float lookups scan flattened values, while per-row lookups stop at a match. The helper extraction preserves array_remove behavior.
  • 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.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa, 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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:expressions Expression evaluation enhancement New feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Run array_contains on float elements natively with Spark's equality instead of through the codegen dispatcher

3 participants