feat: Lambda function support from DataFusion, illustrated with array_filter - #4744
kazantsev-maksim wants to merge 147 commits into
Conversation
This reverts commit 768b3e9.
# Conflicts: # native/core/src/execution/planner.rs
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: General
filterlambdas used JVM dispatch. This PR moves supported predicates into DataFusion to avoid that execution overhead. - Design approach: Adds HOF protobuf messages, expression-ID-based lambda scopes, native planning, and native → JVM → Spark fallback.
- Correctness / compatibility analysis: The existing P2 short-circuit report remains unresolved. For Parquet input
a=[0,1],filter(a, x -> x = 0 OR monotonically_increasing_id() = 0)returns[0,1]in Spark 4.1.3 but[0]in the native component reproduction.CometHighOrderFunction.scala:195still admits this predicate, and DataFusion evaluates the skipped RHS, advancing its counter. The strict masking described in the discussion is absent from this head. Spark's relevantArrayFiltersemantics agree across supported 3.4–4.2 sources. - Key design decisions: Reusing upstream lambda projection and array filtering limits custom machinery. The expression-ID scope stack handles binding, and
EmptyBatchGuardExprpreserves empty-input behavior. The conditional fallibility check remains insufficient for Spark semantics. - Implementation sketch: Reviewed all 15 changed files, including serde registration, configuration, native planner, protobuf additions, SQL fixtures, and benchmarks.
- Behavioral changes worth calling out: Native unary filters are enabled by default. Indexed lambdas retain dispatch fallback. Empty, mixed empty/null, and all-null component checks passed. Benchmarks cover captures, strings, nulls, and nested arrays, but were not rerun.
- Suggested improvements: Resolve the existing blocker with strict per-element evaluation masks or conservative dispatch, and assert native execution for its regression tests. No additional introduced P1/P2 issues were found within this review.
Reviewed full SHA 670847d7a3ca27e32d0eb23e1f8f5d26c60fce03 against 466e3fe3a7200f35113b523249352b5a33a21365. Routed skills: review-comet-pr and review-comet-expression-pr. Read existing discussion and threads, excluding Copilot reviews.
Exact-head CI: Comet CI and CodeQL report action_required. Only labeling passed.
Validation limits: Native probes compiled this head's guard and counter against cached DataFusion 55.1.0 components, with relevant sources verified byte-identical to upstream. Spark reference execution passed with ANSI both enabled and disabled. Full Comet/JNI execution and SQL suites were not run: Maven bootstrap failed DNS resolution, the system Maven rejected repository configuration syntax, and the native build lacked JNI headers. The Spark-version test matrix was not run. git diff --check passed.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: General
filterlambdas used JVM dispatch. This PR introduces native evaluation to reduce that overhead. - Design approach: Adds HOF protobuf messages, expression-ID-based scopes, DataFusion lambda planning, and native → JVM → Spark fallback.
- Correctness / compatibility analysis: The existing P2 short-circuit concern remains unresolved. For Parquet input
a=[0,1],filter(a, x -> x = 0 OR monotonically_increasing_id() = 0)returns[0,1]in Spark 4.1.3 with ANSI enabled or disabled. The native component reproduction returns[0]because evaluating the skipped branch advances the counter. The strict masking described in the discussion is absent from this head. Spark's relevantArrayFiltersemantics agree across supported 3.4–4.2 sources. - Key design decisions: Reusing DataFusion's lambda projection and filtering limits custom machinery. Expression-ID scopes handle binding, and
EmptyBatchGuardExprcorrectly preserves empty-input behavior. The conditional fallibility check does not cover the existing blocker. - Implementation sketch: Reviewed the full base-relative diff across all 15 changed files, surrounding code, existing reviews and discussion, excluding Copilot.
- Behavioral changes worth calling out: Native unary filters are enabled by default. Indexed lambdas retain dispatch fallback. Empty, mixed empty/null, and all-null component checks passed. Benchmarks cover captures, strings, nulls and nested arrays, but were not rerun.
- Suggested improvements: Resolve the existing blocker with strict per-element evaluation masks or conservative dispatch, and assert native execution in its regression tests. No additional introduced P1/P2 issues found within this review.
Reviewed full SHA 670847d7a3ca27e32d0eb23e1f8f5d26c60fce03 against 466e3fe3a7200f35113b523249352b5a33a21365. Routed skills: review-comet-pr and review-comet-expression-pr.
Exact-head CI: Comet CI and CodeQL report action_required. Only labeling passed.
Validation: Recompiled this head's guard and counter in the native component harness against cached DataFusion 55.1.0 artifacts. Verified relevant dependency sources against upstream. Reran Spark 4.1.3 reference execution and git diff --check successfully. Full locked native build, Comet/JNI SQL suites, benchmarks and the Spark-version execution matrix were not rerun. Prior exact-head build logs record Maven bootstrap/configuration and JNI-header failures, so full integration remains unverified.
|
Thanks for catching that, @sunchao. You were completely right — the To address your point about asserting native execution and prevent any silent JVM codegen fallback, I have also explicitly disabled the codegen dispatcher ( All tests now pass natively without fallback, including:
Could you please take another look when you have a chance? |
sunchao
left a comment
There was a problem hiding this comment.
Re-reviewed 8f5bc62341e3f4394e1bd35b87ca2ac15f977288 against 605051ad239ef704f5f25d67910a446a6b6d7c70 with five independent scopes covering native planning/binding, Spark serde and fallback, runtime errors and memory, compatibility, and performance/tests. Findings were independently reproduced and challenged.
The strict boolean adapter fixes the earlier AND/OR examples: guarded abs(INT_MIN), legacy element_at(..., 0), seeded rand, monotonically_increasing_id, and a reused ANSI Dataset after changing the session setting now agree with Spark. Empty/mixed-null arrays and nested capture/shadowing controls also pass. Reusing DataFusion's lambda infrastructure and keeping the masking adapter local is a reasonable design. Indexed/unsupported lambdas retain JVM dispatch when it is enabled; I found no additional unintended Spark row fallback.
Three P2 issues remain, detailed inline: CASE metadata checks consume state, COALESCE duplicates nondeterministic child evaluation, and captured arrays can cause multiplicative allocation growth. The two correctness issues have distinct causes and both reproduce through the complete current-head native execution path with JVM dispatch disabled. I am requesting changes for these issues.
Validation:
- Built native code with
cargo build --locked -p datafusion-comet --no-default-features, then ran Maven from the root reactor. All fourCometSqlFileTestSuite array_filterfixtures passed (4 tests, 1 suite, no failures). - Ran fresh Spark 4.1.3 / JVM-dispatch / native comparisons on Parquet inputs, including the two wrong-result cases, controls, the earlier reported failures, and captured-array routing. Native probes required
CometProjectwith dispatch disabled. - Focused native checks covered nested binding, scalar/array three-valued boolean masks, empty/null inputs, and capture allocation growth.
Current-head CI: Comet CI and CodeQL remain action_required. Only the label check has passed.
Limits: full upstream Spark SQL and the Spark-version execution matrix were not run. Release benchmarks were not rerun. Allocation numbers are from a focused native component measurement, not a complete-query memory or timing benchmark.
| fn nullable(&self, input_schema: &Schema) -> Result<bool> { | ||
| self.inner.nullable(input_schema) |
There was a problem hiding this comment.
[P2] Keep lambda nullability checks from consuming predicate state
Could the lambda wrapper report nullability without evaluating stateful predicates? With one Parquet row containing a=[1,2,3], this query returns [1] in Spark and JVM dispatch, but [] through the native path at this head:
SELECT filter(a, x ->
(CASE WHEN monotonically_increasing_id() = 0 THEN x ELSE 0 END) > 0)
FROM t;HigherOrderFunctionExpr asks the lambda body for its return field during construction and evaluation. This delegation reaches DataFusion's CaseExpr::nullable(), which evaluates the WHEN condition on a synthetic one-row batch and advances the same counter later used for real elements. The conditional guard admits the expression because its result branches are not fallible. The previous whole-filter JVM dispatcher does not perform these evaluations. Please keep metadata inspection free of these side effects and add a native regression with a stateful CASE condition. I reproduced this through full Spark 4.1.3/Comet execution with JVM dispatch disabled.
| // COALESCE: tail arguments | ||
| case Coalesce(children) if children.length > 1 => | ||
| children.tail.exists(isFallibleExpr) |
There was a problem hiding this comment.
[P2] Preserve single evaluation of nondeterministic COALESCE children
Could native admission also account for nondeterministic nonfinal COALESCE arguments? With one Parquet row a=[1,2,3,4], this returns [2,4] in Spark and JVM dispatch, but [4] with native HOF execution and dispatch disabled:
SELECT filter(a, x ->
coalesce(IF(rand(42L) < 0.6, x, CAST(NULL AS INT)), -1) = x)
FROM t;CometCoalesce serializes the child separately as WHEN IS NOT NULL(child) and THEN child. The second random-expression instance runs only on the selected subset, so its draws correspond to different elements. This remains wrong even when CASE metadata evaluation is suppressed, so it is separate from the nullability issue. Plain IF with the same seeded predicate and deterministic COALESCE controls both agree with Spark. Please ensure each child is evaluated once per element, or retain whole-filter JVM dispatch for these shapes until the native lowering preserves that contract. This mismatch is reproduced through full current-head Spark 4.1.3/Comet execution.
| let mut args = value_args; | ||
| args.extend(lambdas); | ||
|
|
||
| Ok(Arc::new(HigherOrderFunctionExpr::try_new_with_schema( |
There was a problem hiding this comment.
[P2] Avoid copying captured arrays once per lambda element
Could we prevent deep replication of captured arrays before selecting this path by default? For filter(a, x -> x >= 0 AND size(b) > 0), DataFusion's evaluate_single_list_lambda broadcasts captured b with take_arrays, copying its entire child buffer once for each element of a. With one row containing 4,096 integers in each array, a focused native allocation probe measured 67,705,058 bytes allocated from 33,200 bytes of input. The repeated integer payload alone is 64 MiB, although this predicate only needs b's length. The cost grows with both array lengths and batch row count.
Spark retains the outer-row array reference, and the base implementation dispatches this filter to Spark's evaluator. I confirmed that Catalyst keeps size(b) inside the lambda and that the current Comet query takes the native path with dispatch disabled. Please avoid expanding complex captures this way, or retain the existing JVM dispatcher for affected shapes until native captures can represent them efficiently. An allocation benchmark with a captured array would cover this regression; the current capture benchmark uses a scalar integer. These allocation numbers are component measurements, not complete-query memory or timing results.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: General
filterlambdas used JVM dispatch. This PR introduces native evaluation to reduce that overhead. - Design approach: Adds HOF protobuf messages, expression-ID-based lambda scopes, DataFusion planning, and native → JVM → Spark fallback.
- Correctness / compatibility analysis: Found one new P2 regression: null-intolerant comparisons evaluate operands that Spark skips, changing stateful predicate results. Spark’s relevant semantics agree across supported 3.4–4.2 sources. The existing CASE and COALESCE wrong-result reports also reproduce at this head.
- Key design decisions: Reusing DataFusion’s lambda machinery limits custom code. The strict boolean adapter correctly handles tested three-valued logic and masked evaluation, but its scope does not cover implicit null short-circuiting. The existing captured-array replication concern remains substantiated by the
take_arrayscapture path. - Implementation sketch: Reviewed the entire base-relative diff across all 16 files, including surrounding planner and serde code, configuration, protobuf, SQL fixtures, benchmarks, and documentation.
- Behavioral changes worth calling out: Supported unary filters execute natively by default. Indexed and unsupported shapes retain fallback. Nested capture, shadowing, empty-input, and boolean-mask controls passed.
- Suggested improvements: Preserve null-dependent operand evaluation and add a native regression for the new finding. Resolve the three blockers in the existing review, including captured-array allocation growth.
Reviewed full SHA 8f5bc62341e3f4394e1bd35b87ca2ac15f977288 against 605051ad239ef704f5f25d67910a446a6b6d7c70. Routed skills: review-comet-pr and review-comet-expression-pr. Read existing reviews, discussion, inline comments, and threads, excluding Copilot.
Exact-head CI: Comet CI and CodeQL remain action_required. Only labeling passed.
Validation: Locked native build with --no-default-features and root Maven compilation passed. All four array_filter SQL fixtures passed. Ran 36 Spark/JVM-dispatch/native comparisons with ANSI enabled and disabled. Native comparisons required CometProject with dispatch disabled. Focused serializer and native boolean-mask checks also passed.
Validation limits: Full upstream Spark SQL, the Spark-version execution matrix, lint, and release benchmarks were not run. The existing capture-allocation measurement was not rerun.
| let op = match binary.op() { | ||
| Operator::And => Some(ShortCircuitBinaryOp::And), | ||
| Operator::Or => Some(ShortCircuitBinaryOp::Or), | ||
| _ => None, |
There was a problem hiding this comment.
[P2] Preserve null short-circuiting for comparison operands. With one Parquet row containing a=[NULL,0] as ARRAY<BIGINT>, filter(a, x -> x = monotonically_increasing_id()) returns [0] in Spark and JVM dispatch but [] natively. Spark’s BinaryExpression.eval skips the RHS when the LHS is null. This rewrite leaves equality as DataFusion BinaryExpr, which evaluates the counter for both elements, so the comparison against zero sees counter value 1. The newly enabled native path therefore silently drops a valid element. Could it preserve the null-dependent evaluation mask for affected expressions, or dispatch those lambda bodies until that contract is supported?
Evidence: Reproduced through full exact-head Comet execution on Spark 4.1.3 in local[1]. Write SELECT CAST(array(NULL, 0) AS ARRAY<BIGINT>) AS a to Parquet, register it as t, then run SELECT filter(a, x -> x = monotonically_increasing_id()) FROM t. With spark.comet.exec.higherOrderFunction.native.enabled=true and spark.comet.exec.scalaUDF.codegen.enabled=false, the asserted CometProject returns [], with ANSI either enabled or disabled. Spark and whole-filter JVM dispatch both return [0]. Adding x IS NOT NULL AND restores [0] natively. A component reproduction using this head’s lambda adapter and counter independently produced the same mismatch.
There was a problem hiding this comment.
Thanks @sunchao for the thorough review and for pointing out the memory replication in take_arrays as well as the boolean conditional discrepancies, @sunchao!
I have addressed both blockers by refining the degradation rules in CometHighOrderFunction.scala:
1. Complex Outer Captures (capturesComplexOuterAttribute)
- Lambda bodies that capture outer columns of complex types (
ArrayType,MapType,StructType) now safely degrade to JVM codegen dispatch. - This avoids quadratic memory replication from DataFusion's
take_arrays(such as the 64 MiB allocation on 4k elements). - Capturing scalar columns (
Int,String, etc.) remains fully native and fast.
2. Boolean Conditional Expressions (hasUnsupportedConditionals)
- Conditional expressions (
CaseWhen,If,Coalesce) that returnBooleanTypeand serve as filter predicates now degrade to JVM codegen dispatch. This avoids DataFusion's native evaluation discrepancies ([1] -> []and[2, 4] -> [4]). - Scalar conditionals (e.g.
coalesce(x, 0) > 0) continue to execute natively in DataFusion.
3. Tests & Documentation
- Added regression tests covering:
- Complex outer array capture:
filter(a, x -> x >= 0 AND size(b) > 0) - Boolean
COALESCE:filter(a, x -> coalesce(x % 2 = 0, false)) - Boolean
CASE WHEN:filter(a, x -> CASE WHEN x = 1 THEN true ELSE x > 2 END) - Boolean
IF:filter(a, x -> if(x = 2, false, true)) - Verified that scalar conditionals (
filter(arr, x -> coalesce(x, 0) > 0)) continue to execute natively with codegen disabled.
- Complex outer array capture:
- Updated the ScalaDoc on
CometHighOrderFunctionand the corresponding PR documentation.
All native and fallback tests pass. Could you please take another look when you have a chance?
|
This is a light fully automated review since there are so many PRs open.
|
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: General
filterlambdas used JVM dispatch. This PR introduces native evaluation to reduce that overhead. - Design approach: Adds HOF protobuf messages, expression-ID-based lambda scopes, DataFusion planning, and native → JVM → Spark fallback.
- Correctness / compatibility analysis: Reproduced the existing CASE (
[1]→[]), COALESCE ([2,4]→[4]), and null-comparison ([0]→[]) regressions through full native execution, with ANSI enabled and disabled. Spark and JVM dispatch agree. Relevant Spark semantics match across supported 3.4–4.2 sources. - Key design decisions: Reusing DataFusion’s lambda infrastructure limits custom code. The strict boolean adapter passes focused masking and three-valued-logic checks. The existing captured-array allocation concern remains substantiated by the
take_arraysreplication path. - Implementation sketch: Reviewed the entire 16-file base-relative diff and surrounding code, including serialization, binding, configuration, SQL fixtures, benchmarks, and documentation.
- Behavioral changes worth calling out: Supported unary filters run natively by default. Indexed and unsupported shapes retain fallback. Nested captures, shadowing, empty inputs, and boolean-mask controls passed.
- Suggested improvements: Resolve the existing CASE, COALESCE, and capture-allocation blockers and null short-circuiting blocker. Preserve side-effect-free metadata inspection, single evaluation, operand masks, and bounded capture allocation. No additional introduced P1/P2 issues found within this review.
Reviewed full SHA 8f5bc62341e3f4394e1bd35b87ca2ac15f977288 against 605051ad239ef704f5f25d67910a446a6b6d7c70. PR remains non-draft. Routed skills: review-comet-pr, review-comet-expression-pr, and review-comet-memory-pr. Read existing discussions and review threads, excluding Copilot.
Exact-head CI: Comet CI and CodeQL remain action_required. Only labeling passed.
Validation: Locked native build with --no-default-features passed. Root Maven reactor ran all four array_filter fixtures successfully. Completed 36 Spark/JVM-dispatch/native executions with native comparisons requiring CometProject and dispatch disabled. Focused Rust boolean-mask checks passed.
Validation limits: Full upstream Spark SQL, the Spark-version execution matrix, lint, release benchmarks, and the prior allocation measurement were not rerun. Project source remains unchanged.
# Conflicts: # native/core/src/execution/planner.rs
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: General
filterlambdas used JVM dispatch. This PR adds native evaluation to reduce that overhead. - Design approach: Adds HOF protobuf messages, expression-ID-based scopes, DataFusion lambda planning, and native → JVM → Spark fallback.
- Correctness / compatibility analysis: Existing P2 blockers remain. Fresh native component reproductions still produce
[]instead of[1]for the reported CASE example,[4]instead of[2,4]for COALESCE, and[]instead of[0]for the null-comparison example. Spark 4.1.3 Catalyst checks confirmed the expected results. The exact-head admission guards accept all three: the conditionals return scalar values, and the comparison contains no conditional. RelevantArrayFiltersemantics agree across Spark 3.4–4.2 sources. - Key design decisions: Reusing DataFusion’s lambda infrastructure limits custom machinery. The strict AND/OR adapter passed three-valued-logic, scalar, empty-input, stateful-mask, and rewrite checks. The capture guard fixes ordinary complex-column captures but misses captured outer
NamedLambdaVariables. A nested-capture component probe still allocated 67,163,908 bytes from a 16,600-byte input with 4,096 elements, substantiating the existing allocation concern. - Implementation sketch: Reviewed the full 17-file base-relative diff and surrounding code, including serialization, native planning, configuration, protobuf, SQL fixtures, benchmarks, and documentation. Read existing reviews, discussion, inline comments, and threads, excluding Copilot.
- Behavioral changes worth calling out: Supported unary filters run natively by default. Indexed lambdas, Boolean conditionals, complex outer-column captures, and unsupported expressions retain fallback. Scalar conditionals remain native, including the reported failing cases.
- Suggested improvements: Resolve the existing CASE, COALESCE, null-comparison, and capture-allocation blockers. Preserve side-effect-free metadata inspection, single evaluation, null-dependent operand masks, and bounded capture allocation, or dispatch affected bodies. No additional introduced P1/P2 issues found within this review.
Reviewed full SHA bd07f54aac4709c312fd24288539e7f03b43eca9 against e3d52feaf6584295f6870a2b72d95ebf83a18a05. PR remains non-draft. Routed skills: review-comet-pr, review-comet-expression-pr, and review-comet-memory-pr.
Exact-head CI: Comet CI and CodeQL report action_required. Only labeling passed.
Validation: Locked offline native cargo check --no-default-features and git diff --check passed. Recompiled component probes using exact-head adapters and kernels against cached DataFusion 55.1.0 artifacts, verifying seven relevant dependency sources against upstream. Compiled and ran Spark 4.1.3 Catalyst probes with guard methods verified identical to this head.
Validation limits: Full JVM/JNI execution, SQL fixtures, upstream Spark SQL, the Spark-version execution matrix, and release benchmarks were not run. Maven bootstrap failed DNS resolution for repo.maven.apache.org; system Maven rejected the repository’s configuration syntax. Allocation figures measure component allocations, not complete-query peak memory or timing. Project source remains unchanged.
# Conflicts: # native/core/src/execution/planner.rs
|
Thanks for the thorough investigation and reproduction cases, @sunchao! I have resolved both remaining blockers:
All native tests and fallback regressions now pass cleanly. |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: General
filterlambdas used JVM dispatch. This PR introduces native evaluation to reduce that overhead. - Design approach: Adds HOF protobuf messages, expression-ID-based lambda scopes, DataFusion planning, and native → JVM → Spark fallback.
- Correctness / compatibility analysis: The existing P2 null-short-circuiting blocker remains unresolved. For
a=[NULL,0]asARRAY<BIGINT>,filter(a, x -> x = monotonically_increasing_id())returns[0]in Spark but[]in the freshly compiled native component reproduction. Exact-head serialization still admits this expression. The new guard excludesEqualNullSafe(<=>), whereas this query usesEqualTo(=). Native equality evaluates the counter for the null element, advancing its state. Relevant Spark semantics agree across supported 3.4–4.2 sources. - Key design decisions: Reusing DataFusion’s lambda infrastructure limits custom machinery. Strict AND/OR masking passed three-valued-logic, scalar, stateful-order, and rewrite checks. Empty/null-array controls passed. The reported scalar CASE/COALESCE and nested complex-capture shapes now decline native execution.
- Implementation sketch: Reviewed the entire 17-file base-relative diff, surrounding code, serialization, configuration, native planning, SQL fixtures, benchmarks, and documentation. Read existing reviews and discussion, excluding Copilot.
- Behavioral changes worth calling out: Compared with
branch-1.1, general supported unary filters become native by default. Indexed lambdas and excluded shapes retain fallback. The complex-capture guard addresses the reported allocation path. Benchmarks were inspected but not rerun. - Suggested improvements: Resolve the existing null-operand blocker by preserving Spark’s evaluation mask or dispatching affected bodies. Add the original ordinary-equality regression. No additional introduced P1/P2 issues found within this review.
Reviewed full SHA 8df68a38dbb32e8de2768b3052b1cc6c1e224750 against 432523fb59efe75b208ae5688eef6012c208a4d6. PR remains non-draft. Routed skills: review-comet-pr, review-comet-expression-pr, and review-comet-memory-pr.
Exact-head CI: Comet CI and CodeQL report action_required. Only labeling passed.
Validation: Locked offline native cargo check --no-default-features, root Maven test-compile with Java 17/Spark 4.1.3, exact-head serializer probes, Spark Catalyst checks, native component checks, and git diff --check passed. Component dependency versions match the lockfile, and seven relevant DataFusion sources were verified against upstream.
Validation limits: Full native linking/JNI execution, SQL fixtures, upstream Spark SQL, the Spark-version execution matrix, lint, and benchmarks were not run. Project source remains unchanged.
|
All tests are now passing across both native and fallback suites! @sunchao Summary of Final Updates:
|
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: General
filterlambdas used JVM dispatch. This PR introduces native evaluation to reduce that overhead. - Design approach: Adds HOF protobuf messages, expression-ID-based scopes, DataFusion lambda planning, and native → JVM → Spark fallback.
- Correctness / compatibility analysis: The existing P2 null-short-circuiting concern remains unresolved for fallible operands. With ANSI enabled,
local[1], and one Parquet row containinga=[NULL]asARRAY<DOUBLE>,SELECT filter(a, x -> x = 1D / spark_partition_id()) FROM treturns[]in Spark 4.1.3. Exact-head serialization admits native execution with dispatch disabled. A native component reproduction using this head’s lambda adapters andchecked_divthrowsDIVIDE_BY_ZEROin partition 0. Spark skips the RHS whenxis null. The guard atCometHighOrderFunction.scala:197–205excludes selected stateful expressions but still admits this failing RHS. An explicitx IS NOT NULL ANDguard makes the native control succeed. Relevant Spark semantics agree across supported 3.4–4.2 sources. - Key design decisions: Reusing DataFusion’s lambda machinery limits custom code. Strict AND/OR masking passed three-valued-logic, scalar, empty-input, stateful-order, and rewrite checks. The reported counter, CASE/COALESCE, and nested complex-capture examples now decline native execution.
- Implementation sketch: Reviewed the full 17-file base-relative diff and surrounding code, including serialization, native planning, configuration, SQL fixtures, benchmarks, and documentation. Read all existing non-Copilot discussion and review threads.
- Behavioral changes worth calling out: Compared with
branch-1.1, supported unary filters become native by default. Indexed and excluded shapes retain fallback. The eleven benchmark shapes were inspected but not rerun, so this review makes no new performance claim. - Suggested improvements: Complete the existing null-operand fix by preserving Spark’s evaluation masks for fallible operands or dispatching affected bodies. Add the regression above. No additional introduced P1/P2 issues found within this review.
Reviewed full SHA 2179c7a75580908fb12456c8449dd3348b29dac7 against 432523fb59efe75b208ae5688eef6012c208a4d6. PR remains non-draft. Routed skills: review-comet-pr, review-comet-expression-pr, and review-comet-memory-pr.
Exact-head CI: Comet CI and CodeQL report action_required. Only labeling passed.
Validation: Locked offline native cargo check --no-default-features, root Maven test-compile with Java 17/Spark 4.1.3, exact-head serializer checks, Spark Parquet reference execution, native component probes, and git diff --check passed. Component dependency versions match the lockfile, and seven relevant DataFusion sources match previously verified upstream copies.
Validation limits: The failing native result is a component reproduction, not full Comet/JNI execution. Full native linking, Comet SQL fixtures, upstream Spark SQL, the Spark-version execution matrix, lint, and benchmarks were not run. Project source remains unchanged.
|
Thanks for the reproducible example, @sunchao! That was the missing piece — Spark's Rather than trying to detect fallible expressions in Scala, I have resolved this directly in Rust by extending
All native and fallback suites pass. Could you please take another look? |
Which issue does this PR close?
N/A
Rationale for this change
Running higher-order functions through JVM codegen is expensive: each batch incurs a JNI call into Spark's own implementation. Moving the lambda evaluation into the native DataFusion engine removes that overhead and brings the plan closer to fully native execution.
What changes are included in this PR?
1. Protobuf (
native/proto/src/proto/expr.proto)HigherOrderFunc,LambdaFunction, andNamedLambdaVariable.high_order_func(71) andnamed_lambda_variable(72) fields toExpr.2. Lambda Infrastructure & Scope Management (
native/core/src/execution/lambda.rs)NamedLambdaVariableby SparkexprId, preventing name shadowing or column collisions.LambdaParamsCapture(pin_unused_params) to prevent DataFusion's optimizer from pruning unused lambda parameters and preserving physical batch structure.EmptyBatchGuardExpr, a transparent physical expression adapter wrapping lambda bodies. When an input batch has 0 rows (e.g. non-null empty arrays[]or mixed[[], NULL]), it short-circuits evaluation and returns an empty array directly. This preserves Spark's ANSI contract where lambda predicates are never evaluated on empty collections (avoiding runtime errors like scalar division by zero).3. Physical Planner (
native/core/src/execution/planner.rs)PhysicalPlannerto supportHigherOrderFuncexpressions: resolves parameter field types via the HOF UDF contract, plans the lambda body under the resolved scope, and binds physicalLambdaVariableindices.4. Spark Serde & Three-Tier Execution (
CometHighOrderFunction.scala,arrays.scala)Native -> JVM Codegen -> Sparkfallback hierarchy controlled byspark.comet.exec.higherOrderFunction.native.enabledandspark.comet.exec.scalaUDF.codegen.enabled.try-catch NonFatalto gracefully decline the native path if eager evaluation occurs in unreachable branches during planning (e.g.CometCastevaluating literal arguments in guarded branches under ANSI mode).hasGuardedFallibleBranchto decline the native path when conditional operators (AND,OR,CASE WHEN,IF,COALESCE) guard fallible expressions (division, non-try casts, overflow, indexing) under ANSI mode. This accounts for DataFusion's vectorized batch short-circuit threshold (20%) and avoids runtime exceptions on elements skipped by Spark's per-element evaluation.How are these changes tested?
array_filteroperations with captures, literals, string functions, and nested structures.[[], NULL]under ANSI mode (spark_partition_id()division by zero).ANDandORpredicates on[0, 1]under ANSI mode.CASE WHENwith malformed casts on[-1, 0]under ANSI mode.EmptyBatchGuardExprshort-circuiting on empty batches.