Skip to content

fix: preserve Spark evaluation for next_day and levenshtein - #5972

Open
sunchao wants to merge 9 commits into
apache:mainfrom
sunchao:codex/nextday-levenshtein-evaluation-oss-20260915
Open

sunchao wants to merge 9 commits into
apache:mainfrom
sunchao:codex/nextday-levenshtein-evaluation-oss-20260915

Conversation

@sunchao

@sunchao sunchao commented Sep 15, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Part of #6006. This PR is stacked on #5533 and uses its shared evaluation-mask policy.

Rationale for this change

Spark can skip an ANSI next_day expression below a limit, within a conditional, or in a deferred projection. Evaluating that expression earlier in Comet can raise an invalid-weekday error for a row Spark never evaluates. Spark also disagrees between generated and interpreted evaluation of three-argument levenshtein when its threshold is NULL.

What changes are included in this PR?

NextDay enrolls in RequiresSparkEvaluationMask, sharing #5533's traversal, early exit, configuration switch and aggregate-buffer restoration. Main's current aggregate safety and rule composition remain intact. Valid literal weekdays and literal NULL do not need this protection.

Serializer and deferred-projection analysis share the eager-child rule. Eager date arithmetic, concat/greatest, sort keys and unfiltered max(next_day(...)) retain native execution. Aggregate protection applies to FILTER and First/Last update masks. Whole-expression JVM dispatch uses the general serde policy rather than a hardcoded expression name. A shared helper handles interpreted projection configuration across Spark versions.

Three-argument levenshtein with a nullable threshold stays in Spark in every factory mode, including default FALLBACK; this also covers runtime fallback from generated to interpreted projections. Two-argument calls and known non-null thresholds retain native execution. The existing raw-bytes collation explanation is preserved, and the compatibility guide describes both behaviors.

Filter input analysis follows Spark's generated predicate order, including null-check reordering. Values required by the first emitted predicate retain native execution; later or null-gated values keep their deferred-evaluation protection.

This stack includes #5533's fused Top-K/AQE fix: a local selector only owns sorting, while the outer Top-K applies the result projection and offset once. The deferred-projection logic remains in this follow-up and is exercised against both next_day admission controls and the complete shared unbase64 mask suite. The upstream Spark nullable-threshold issue still needs a JIRA filing; a reproducible report is supplied in the review follow-up because the available Jira connection is not authenticated.

How are these changes tested?

Current rebased implementation (c7c8c73): Spark 4.1.3 passed all 106 tests in the full CometExecRuleSuite and the next_day_safe_evaluation SQL matrix. Spark 3.5.9 passed strict main/test compilation and 105 tests from the same suites; the existing map-grouping test is guarded to Spark 4+. Both whole-stage codegen modes are covered.

The eager-parent SQL regression fails with the preceding serializer, reproducing unnecessary DateSub, NextDay and CreateArray fallback. With the fix these parents retain native next_day; invalid eager inputs still raise Spark's error, a leading-null array still evaluates later elements, and nullable-left date expressions preserve masked-right fallback. Planner controls verify that deferred parent outputs still retain Spark projections.

The full planner suite includes the filter-order, null-check, reused-plan, aggregate-buffer and Top-K/AQE regressions. Native/protobuf inputs match main 7425780; these JVM tests used its verified CI library from run 37346691204 (artifact 11363961998). Spotless and whitespace checks pass. Broader Spark SQL and Spark-profile runtime checks are requested by the existing CI labels and remain pending.

@github-actions github-actions Bot added bug Something isn't working area:expressions Expression evaluation labels Sep 15, 2026
Comment thread spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala Outdated
@rich7420

Copy link
Copy Markdown
Contributor

@sunchao thanks for the patch!

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

Thanks for digging into this. I confirmed the levenshtein divergence against the Spark 4.0.1 source and it is real: eval unboxes a Some(null) threshold to 0 while doGenCode wraps the body in ctx.nullSafeExec and leaves the result NULL. Our native kernel returns NULL, so it matches the generated path and diverges wherever Spark interprets. There is definitely something to fix here. My concerns are about the shape of the fix, plus one regression I measured that I would like resolved before this goes in.

This duplicates #5533

preserveNextDayEvaluationMasks is close to line-for-line identical to preserveEvaluationMasks in #5533. The firstMatch, startsLimit, materializesInput, limitAncestor, finalAggregate, limitName, aggregateBufferName, restartsNative and prepared blocks are the same apart from whitespace. The one real substitution is that #5533 looks the expression up through RequiresSparkEvaluationMask and this hardcodes case nextDay: NextDay if nextDay.failOnError.

Both PRs are open, both are yours, and both edit the same lines, so whichever lands second is going to be an unpleasant reconciliation. Could this rebase onto #5533 and enroll NextDay in the policy instead of growing a second traversal? That picks up the config and the compatibility-guide section for free. It would also be good to reference #5533 and #6006 in the description, since #6006 reads like the tracking issue for exactly this family.

Three things this copy dropped that I think are worth keeping:

The early exit. #5533 returns the plan unchanged when no node carries an enrolled expression or a sticky tag. Without it protect walks every node of every plan and at each one computes original.expressions, a collectFirst per expression tree, supportsWholeStage (another full expression walk plus two isTooManyFields calls) and eagerReferences. Every query pays that, next_day or not.

The config. #5533 added spark.comet.exec.preserveEvaluationMasks.enabled. There is no way to turn this one off, which matters given the next section.

The shared restoreNativeAggregateBuffers. #5533 factors one version shared with the Celeborn fallback. This adds a third local copy, which also drops the withFallbackReason tagging the shared one does.

The fallback blast radius is wider than it needs to be

I built the branch and main and ran the same probe on each, Spark 4.1 default profile, spark.sql.ansi.enabled=true, parquet table, AQE off:

query main this PR
next_day(d, dow) CometProject CometProject
date_add(next_day(d, dow), 1) CometProject Spark Project
datediff(next_day(d, dow), d) CometProject Spark Project
concat(cast(next_day(d, dow) AS string), 'x') CometProject Spark Project
greatest(next_day(d, dow), d) CometProject Spark Project
ORDER BY next_day(d, dow) CometSort + CometExchange Spark Sort + Spark Exchange
WHERE k = 1 AND next_day(d, dow) > date'2020-01-01' CometFilter Spark Filter
if(k = 1, next_day(d, dow), null) CometProject CometProject
GROUP BY next_day(d, dow) 2x CometHashAggregate 2x CometHashAggregate
max(next_day(d, dow)) 2x CometHashAggregate Spark partial + Comet final

Only the AND row is a genuine mask. Spark's generated code for DateAdd, DateDiff, Concat and Greatest evaluates the next_day operand on every row, and a sort key is computed for every row.

The cause is the default arm of unsupportedNextDayEvaluation, case _ => pushChildren(Some(current.nodeName)). Any expression with two or more children outside the five-case whitelist marks all of its children unsafe, including the first child that a nullIntolerant parent always evaluates. eagerReferences in this same PR already gets that right with take(if (expr.children.head.nullable) 1 else 2), so the two halves disagree about date_add. Could they share the rule? And for the ORDER BY case, could SortOrder push None for its child and skip sameOrderExpressions, which are never serialized anyway?

ANSI defaults to on in Spark 4, so this is the default configuration, and none of these shapes has a test or a mention in the description.

The aggregate guard is unconditioned

The check in aggExprToProto declines every Partial and Complete aggregate whose arguments contain an ANSI next_day. The reason text names FILTER, null-skipping and state-dependent updates, but the code checks none of them, and aggExpr.filter is never read. That is what costs max(next_day(...)) its partial aggregate above, and Max's update expression is greatest(max, child), which has no branch. The Levenshtein check directly below it is properly scoped to ImperativeAggregate. Could the NextDay arm be narrowed the same way, to aggExpr.filter.isDefined plus the aggregates that actually skip their child such as First and Last?

levenshtein: worth reporting upstream, and the FALLBACK path is still open

Since eval and doGenCode genuinely disagree, this is a Spark bug rather than a Comet one. Could you file it and link the JIRA here so the workaround has an expiry?

The guard covers CODEGEN_FACTORY_MODE=NO_CODEGEN, which is internal and test-only. What about the default FALLBACK, where generated code that fails to compile drops silently to InterpretedUnsafeProjection? Spark returns 0 there and we return NULL, and nothing catches it. Given how narrow a nullable three-argument levenshtein is, is a plain Unsupported for that shape the better trade than trying to predict which contexts interpret?

Smaller things

wholeExpressionDispatch singles out UnBase64. The reasoning holds for any CodegenDispatchFallback serde whose instance reports Unsupported, which includes CometNextDay and CometLevenshtein themselves. Could that be the general test rather than one expression name hardcoded inside a function about another?

SQLConf.CODEGEN_FACTORY_MODE.toString compared against "NO_CODEGEN" now appears in both CometExecRule and QueryPlanSerde, with the Spark 4.2 explanation comment on only one of them. Can that move into a shim or a single shared helper?

The CometLevenshtein collation reason change reverts the wording #5720 introduced sixteen days ago. That PR's description says the longer text was chosen so it would not contradict the GenerateDocs "no native implementation and always run in the JVM" header. The new string restates the header and drops the raw-bytes explanation. Was that deliberate?

No user-facing documentation. #5533 added an "Errors from rows Spark skips" section to compatibility/index.md for this behavior class. This changes when ANSI next_day is accelerated under Spark 4 defaults with nothing to tell users about it. Extending #5533's section would be better than adding a second one.

What I would suggest

Split it. The levenshtein work is self-contained and could land on its own once the FALLBACK gap is settled. The ANSI next_day work should rebase onto #5533 and go through RequiresSparkEvaluationMask. The deferred-projection analysis is the genuinely new and riskiest piece, it applies to every enrolled expression rather than just next_day, and I think it deserves its own PR against the general policy where it can be justified and tested on its own terms.

On my side: this needs a rebase, and note that main has gained revertUnsafePartialAggregates and the CometRule composition from #6082 inside the code you are editing, so the clean auto-merge is not proof of much. I will add run-spark-sql since this touches both the serde and the planner and the PR tier does not cover it.

@andygrove andygrove added the run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue label Sep 22, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member Author

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: Comet could change interpreted levenshtein results or evaluate ANSI next_day errors that Spark skips or catches.
  • Design approach: Preserve affected Spark evaluation boundaries through serializer guards, dispatcher selection, and physical-plan fallback.
  • Correctness / compatibility analysis: Reviewed the entire 10-file diff against 4c2ab9686526f30bc782462353e3dc97e3954b3f. Compared Spark implementations across supported 3.4–4.2 profiles. No additional introduced P1/P2 issues found within this review beyond the previously reported fallback regressions.
  • Key design decisions: Separate interpreted from generated evaluation, preserve deferred expressions, and restore compatible aggregate buffers before AQE materialization. These decisions span the serializer and a dedicated planner traversal.
  • Implementation sketch: Add expression-tree checks, restrict the array_compact fast path, tag protected operators, restore reused native subtrees, and add SQL/planner regressions.
  • Behavioral changes worth calling out: Existing P2 concerns remain substantiated. Valid-literal next_day below LIMIT, eager parent expressions, sort keys, and unfiltered max(next_day(...)) unnecessarily lose native execution. The guards still contain the blanket conditions identified in the existing reviews. This confirms lost acceleration without claiming an unmeasured timing regression.
  • Suggested improvements: Resolve the existing safe-literal guard concern and eager-expression/aggregate fallback concerns. Narrow those guards and retain admission controls. These are existing blockers, not new findings.

Reviewed SHA: 935acb146deb1d645c58826afe98683e04925670. PR remains non-draft. Routed skills: review-comet-pr, audit-comet-expression. Read existing reviews, issue comments, inline comments, and threads.

Exact-head CI: 24 successful checks, 13 skipped. Linux Spark 4.1 suites, native tests, and build/lint checks passed. Upstream Spark SQL suites, macOS, and benchmarks were skipped.

Validation: Eleven Spark 4.1.3 reference queries and one expected-error control confirmed the relevant semantics. git diff --check passed. Local Comet tests could not run: Maven wrapper bootstrap failed resolving Maven Central, system Maven rejected the repository configuration, and exact-head CI artifacts had expired. No local Comet runtime parity or performance measurement is claimed. Project code remains unchanged.

@sunchao sunchao added the run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue label Oct 4, 2026
@sunchao
sunchao force-pushed the codex/nextday-levenshtein-evaluation-oss-20260915 branch from 935acb1 to 8f033a6 Compare October 4, 2026 17:58
@sunchao

sunchao commented Oct 4, 2026

Copy link
Copy Markdown
Member Author

Addressed the review in the updated branch, stacked on #5533:

  • Enrolled NextDay in the shared policy; removed the duplicate traversal, retained the early exit/configuration and main's shared aggregate restoration.
  • Added valid-weekday admission below LIMIT, eager parent/sort admission, and native unfiltered MAX controls. Serializer and deferred-reference analysis share their eager-input rule; aggregate fallback is scoped to FILTER and First/Last.
  • Nullable three-argument levenshtein now stays in Spark for FALLBACK, CODEGEN_ONLY and NO_CODEGEN, including containing JVM-dispatched expressions. Safe forms remain native.
  • Generalized whole-expression dispatch, shared the interpreted-mode helper, restored the raw-bytes explanation, and extended the existing compatibility section.

The deferred-projection analysis remains in this follow-up to the shared policy rather than introducing another open dependent PR. Its scope is covered by six focused planner tests and the nine-test shared mask suite, in addition to nine next_day and eight levenshtein SQL cases. All 32 focused tests pass on Spark 4.1.3 across the current-source runs; other profiles and full Spark SQL suites are requested in CI.

The Spark JIRA filing is still outstanding: the available Jira connector is not authenticated. I checked current Spark master and the same eval/doGenCode discrepancy remains. Here is a ready-to-file report so that limitation is explicit:

Title: Nullable levenshtein threshold returns different results in interpreted and generated projections

Run with Comet disabled, using stored input to avoid constant folding:

CREATE TABLE levenshtein_null_threshold (s STRING, t STRING, threshold INT) USING parquet;
INSERT INTO levenshtein_null_threshold VALUES ('same', 'same', NULL), ('a', 'b', NULL);
SET spark.sql.codegen.wholeStage=false;
SET spark.sql.codegen.factoryMode=NO_CODEGEN;
SELECT levenshtein(s, t, threshold) FROM levenshtein_null_threshold;
SET spark.sql.codegen.factoryMode=CODEGEN_ONLY;
SELECT levenshtein(s, t, threshold) FROM levenshtein_null_threshold;

Levenshtein.eval passes Some(null) through asInstanceOf[Int], producing a zero threshold; generated evaluation checks threshold nullness and returns NULL. Expected: both paths return NULL. The default FALLBACK mode can encounter either projection implementation after a compilation failure. The Comet workaround must remain until upstream behavior is consistent across supported Spark versions.

@sunchao

sunchao commented Oct 4, 2026

Copy link
Copy Markdown
Member Author

Additional current-head validation: the complete strict Spark 3.5 / Scala 2.12 root-reactor main/test compilation passes, including Spotless and Scalastyle. Spark 4.1 runtime results remain the 32 focused tests described above; this does not claim Spark 3.5 runtime execution.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member Author

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: Batch execution could evaluate unbase64 and ANSI next_day on rows Spark skips. Nullable levenshtein thresholds also produce different interpreted and generated results in Spark.
  • Design approach: Share expression enrollment and planner protection, preserve affected Spark pipelines, and reject nullable three-argument levenshtein before whole-expression dispatch.
  • Correctness / compatibility analysis: Compared relevant Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The named valid-weekday, eager-date, sort and unfiltered-MAX cases are addressed. One additional P2 regression is reproduced below: filters unnecessarily disable native execution when their predicates always consume the protected value.
  • Key design decisions: Shared enrollment, the early exit and existing aggregate restoration reduce duplication. However, CodegenSupport.usedInputs is not a complete description of eagerly consumed inputs. Synthetic native operators also need restoration rules matching their actual behavior.
  • Implementation sketch: Cache enrolled serdes, inspect expression trees, propagate deferred-input information, tag protected operators, and restore affected native ancestors and aggregate buffers.
  • Behavioral changes worth calling out: Compared with branch-1.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa, masked pipelines and nullable-threshold calls intentionally lose acceleration to preserve Spark evaluation. The existing eager-parent fallback concern remains reproducible for date_sub(next_day(d,dow),1), nested next_day, and array(next_day(d,dow)): these use Spark projections where the base uses CometProject. The prerequisite's fused top-K failure also remains reproducible with valid input and AQE.
  • Suggested improvements: Fix the filter analysis below. Finish addressing the existing eager-parent fallback concern and the top-K restoration blocker in #5533. Those existing concerns are not duplicated as new inline findings. Request changes.

Reviewed full SHA 8f033a6952a621a9ede546f53a78985154936db9. Confirmed non-draft. Reviewed all five commits and 23 files against base 09194436b4a19733d35e6362239f035c78092bbe, using merge base 98662215dc0f485cbbedd4b9a17de8c62fe1fb3f, including the shared prerequisite. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-shuffle-pr. Read existing discussion and threads, excluding Copilot.

Exact-head CI: 53 successful checks, 15 skipped, none failed or pending. Run 37222585075 passed all 20 Comet profile jobs and all nine Spark 4.1 SQL shards. Other upstream Spark SQL profiles, macOS and benchmarks were skipped.

Validation: Current-source JVM main/test compilation, suite registration and whitespace checks passed. Using the SHA256-verified native artifact for this head, all 32 focused SQL/planner/mask tests passed. Two additional probes reproduced the known top-K defect. All five top-K probes passed with base planner/serde controls. Separate admission, filter-consumption and bounded timing probes substantiated the findings. Local runtime validation covered Spark 4.1.3 only. Native compilation and broader platform testing relied on CI. Project source remains unchanged.

Comment thread spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala
@sunchao
sunchao force-pushed the codex/nextday-levenshtein-evaluation-oss-20260915 branch from 8f033a6 to 8ec0da7 Compare October 4, 2026 19:55
@sunchao

sunchao commented Oct 4, 2026

Copy link
Copy Markdown
Member Author

Rebased onto the #5533 fused Top-K/AQE fix and carried its Spark 3.5-compatible projection/offset regression into this stack. The current head is 8ec0da7, which also addresses the latest FilterExec admission comment.

Validation: all ten shared evaluation-mask tests and 25 focused planner/next_day/levenshtein tests pass on Spark 4.1.3. Strict Spark 3.5.9 main/test compilation and both new filter regressions pass. Broader CI is pending; the PR description now records the current results and limits.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member Author

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: Batch execution could evaluate unbase64 and ANSI next_day on rows Spark skips. Spark also disagrees between interpreted and generated levenshtein evaluation with nullable thresholds.
  • Design approach: Share evaluation-mask enrollment and planner protection, preserve affected Spark pipelines, and reject nullable-threshold expressions before whole-expression dispatch.
  • Correctness / compatibility analysis: Compared relevant Spark sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The reported filter-input and fused Top-K/AQE regressions are fixed. No additional introduced P1/P2 issues found within this review.
  • Key design decisions: Shared traversal, an early exit, and existing aggregate-buffer restoration reduce duplication. Separating the local Top-K selector from its enclosing projection and offset matches their actual behavior.
  • Implementation sketch: Serde checks identify sensitive expressions. Planner analysis follows limits, deferred inputs and generated predicate ordering, then restores affected native operators while retaining safe scans.
  • Behavioral changes worth calling out: Compared with branch-1.1 at 992c806a7e38c2e88bd018aa5774164b0850e1fa, masked pipelines and nullable-threshold calls intentionally lose acceleration. The previously reported eager-parent fallback concern remains reproducible: with ANSI enabled and Parquet columns containing (DATE '2024-01-01', 'Monday'), date_sub(next_day(d,dow),1), next_day(next_day(d,dow),'Friday'), and array(next_day(d,dow)) use Spark Project, while merge-base planner/serde controls use CometProject. Spark necessarily evaluates these children. This confirms unnecessary loss of native execution without claiming an unmeasured timing regression.
  • Suggested improvements: Finish addressing the existing eager-parent fallback concern. Extend eagerlyEvaluatedChildren in spark/src/main/scala/org/apache/comet/serde/QueryPlanSerde.scala:1237 for these verified eager inputs and retain admission regressions. Request changes for this existing blocker. No duplicate inline finding is added.

Reviewed full SHA 8ec0da7af711442b8c540ed0eb21f95dcc409b9c. Confirmed open and non-draft. Reviewed all eight commits and 23 files, including the prerequisite stack, against requested base 683f89dc993001f094aa2a6f4a2d00152a0698c8 using GitHub's merge base 98662215dc0f485cbbedd4b9a17de8c62fe1fb3f. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-shuffle-pr. Read existing discussion and threads, excluding Copilot.

Exact-head CI at review completion: 30 successful checks, 21 running, 13 skipped, none failed. All nine Spark 4.1 SQL shards and 11 of 20 Comet profile jobs remain running in run 37230096075. Native build and strict Spark 3.5 compilation passed. Other upstream Spark SQL profiles, macOS and benchmarks were skipped.

Validation: Current-source JVM main/test compilation passed. Using the SHA256-verified exact-head CI native library, all 35 focused tests and six additional Top-K/admission probes passed on Spark 4.1.3. The separate merge-base admission control also passed and confirmed the remaining fallback regression. Whitespace and CI-registration checks passed. Local runtime testing covered Spark 4.1.3 only. Native compilation relied on CI, and local lint checks were skipped. Project source remains unchanged.

@sunchao
sunchao force-pushed the codex/nextday-levenshtein-evaluation-oss-20260915 branch from 8ec0da7 to c7c8c73 Compare October 5, 2026 21:28
@sunchao

sunchao commented Oct 5, 2026

Copy link
Copy Markdown
Member Author

Addressed the eager-parent follow-up in c7c8c73 and rebased the conflicting branch onto current main (0ac4dad). DateSub and NextDay now share the binary eager-prefix rule, and CreateArray evaluates all children. Nullable left operands still protect masked right children, and a deferred parent output still keeps its projection in Spark.

The new SQL regression failed on the preceding logic in both whole-stage codegen modes, reporting the three reviewed fallback reasons. With the fix, Spark 4.1.3 passed all 106 focused SQL/planner tests. Spark 3.5.9 passed strict main/test compilation and 105 tests; the existing Spark-4-only map-grouping test was canceled by its version guard. Coverage includes native admission, invalid weekdays under each eager parent, a leading-null array element, and null-masked DateSub/NextDay controls. The native/protobuf inputs match the verified CI library from main 7425780 (run 37346691204). Broader profile/Spark SQL CI is requested by the existing labels and is pending.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member Author

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: Batch execution could evaluate unbase64 and ANSI next_day on rows Spark skips. Spark also disagrees between interpreted and generated levenshtein evaluation with nullable thresholds.
  • Design approach: Share evaluation-mask enrollment and planner protection, preserve affected Spark pipelines, and guard expressions before whole-expression dispatch.
  • Correctness / compatibility analysis: Reviewed the complete nine-commit, 23-file diff, including prerequisites, against 0ac4dadae70838d8675bfd846c2d1c7617c25c78. Compared relevant Spark implementations across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. No additional introduced P1/P2 issues found within this review. The existing eager-parent fallback concern remains reproducible as described below.
  • Key design decisions: Shared traversal and aggregate-buffer restoration reduce duplication. Cached enrollment, an early exit and a supported configuration switch limit the protection’s overhead. The revised filter analysis and distinction between local Top-K selection and final projection address the previously reported cases.
  • Implementation sketch: Serializer guards reject sensitive expression shapes. Planner analysis follows limits and deferred inputs, tags protected operators, and restores affected native ancestors and incompatible aggregate buffers across AQE passes.
  • Behavioral changes worth calling out: Compared with branch-1.1 at 7b7eec69282a7abc33a0e8e687ca594c72f140ec, masked pipelines and nullable-threshold calls intentionally lose acceleration. However, the existing broad eager-parent concern remains for structs. With ANSI enabled, AQE disabled and Parquet input (DATE '2024-01-01', 'Monday'), SELECT struct(next_day(d,dow)) FROM t uses Spark Project on this head but CometProject with the exact-base planner/serde controls. This reproduces with whole-stage codegen both enabled and disabled. Spark’s CreateNamedStruct.eval and doGenCode evaluate every value field, so this fallback protects no skipped child. This confirms unnecessary loss of native execution, without claiming a measured timing regression.
  • Suggested improvements: Finish addressing the existing eager-parent fallback concern. In spark/src/main/scala/org/apache/comet/serde/QueryPlanSerde.scala:1218, recognize CreateNamedStruct.valExprs as eagerly evaluated inputs. Suggest adding native-admission and deferred-output controls alongside the array cases. Request changes for this existing blocker. No duplicate inline finding is added.

Reviewed full SHA: c7c8c7340232d62727fe9c243646ee9ccb7ae681. Confirmed open and non-draft. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-shuffle-pr. Read existing reviews, issue comments, inline comments and threads, excluding Copilot.

Exact-head CI: 52 successful checks, 15 skipped, none failed or pending. Run 37376062781 passed all 20 Comet profile jobs and all nine Spark 4.1 SQL shards. Other upstream Spark SQL profiles, macOS and benchmarks were skipped.

Validation: Current-source JVM main/test compilation and all 134 focused tests passed on Spark 4.1.3 using the digest-verified exact-head CI native library. Separate head/base admission probes confirmed the remaining struct fallback. Whitespace checks passed. Native compilation and other Spark profiles relied on CI. Local lint checks were skipped, and no performance benchmark was run. Project source remains unchanged.

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

Labels

area:expressions Expression evaluation bug Something isn't working run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants