Skip to content

fix: give each nondeterministic dispatched expression its own kernel - #6712

Open
dwsmith1983 wants to merge 12 commits into
apache:mainfrom
dwsmith1983:fix/6711-dispatcher-kernel-per-occurrence
Open

dwsmith1983 wants to merge 12 commits into
apache:mainfrom
dwsmith1983:fix/6711-dispatcher-kernel-per-occurrence

Conversation

@dwsmith1983

@dwsmith1983 dwsmith1983 commented Oct 6, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #6711.

Rationale for this change

The JVM codegen dispatcher keeps one compiled kernel per task for each distinct pair of serialized expression bytes and column specs. The kernel holds the deserialized expression, and a Catalyst Nondeterministic node such as monotonically_increasing_id() keeps its counter inside it. Every dispatched subtree is bound on its own, so two occurrences of the same subtree in one projection serialize to identical bytes and share one kernel: the second occurrence continues the first one's counter. On a single-partition table with 8-row batches, javaId(monotonically_increasing_id()) in two columns returned a = 0..7, b = 8..15 per batch where Spark returns a = b = id. Spark gives every occurrence its own state.

What changes are included in this PR?

  • DispatchOccurrence, a marker expression that carries an occurrence id. CometScalaUDF.emitJvmCodegenDispatch wraps the bound subtree in it before serializing when the tree contains a Nondeterministic node, taking the id from a JVM-wide counter on the driver. The id travels in the native plan bytes, so every batch and every task of a plan carries the same id for one occurrence, and two occurrences never share one. Deterministic subtrees are not wrapped, so identical ones keep sharing a compiled kernel.
  • CometScalaUDFCodegen.lookupOrCompile strips the wrapper right after deserializing, before compiling. The cache key is unchanged; the id changes the bytes.

Two identical nondeterministic projections now serialize to different native plan bytes, so AQE may no longer treat them as equal for exchange reuse. That only loses a reuse Spark would not perform either, since Spark gives the two projections distinct state.

Out of scope: user functions that keep their own state in one object Spark shares between calls (a ScalaUDF marked nondeterministic, Invoke, StaticInvoke) are not Catalyst Nondeterministic nodes and are not wrapped. #5526 covers them separately.

How are these changes tested?

  • expressions/math/round_nondeterministic_child.sql runs SELECT id, round(rand(1), 2) AS a, round(rand(1), 2) AS b over an 8-row table written as one file. ANSI is off, so Round carries no query context and the two copies serialize to the same bytes. spark.comet.batchSize=2 splits the scan into four batches, so each copy has to carry its generator state from batch to batch. The test compares with Spark and checks that round and rand ran in the dispatcher. Without the occurrence tag it fails, since b continues a's random sequence, and with a kernel rebuilt for every batch it fails too, since rows 2 to 7 repeat rows 0 and 1.
  • non-deterministic dispatched subtrees serialize per occurrence in CometCodegenSuite: two dispatches of a nondeterministic tree produce different payload bytes; a deterministic tree is not wrapped and its payload equals the plain serialization of the bound tree.
  • dispatcher keeps one kernel per occurrence id and its state across batches in CometCodegenSuite: two ids give two cache entries, a repeated id gives one hit, and the counter continues across batches for one id.

The fixture and CometCodegenSuite pass on the Spark 4.1 and 3.5 profiles.

The codegen dispatcher caches one kernel per task for each distinct
serialized expression, and a Nondeterministic node keeps its state in
that kernel. Two occurrences of the same subtree bind to the same
ordinals and serialize identically, so they shared one kernel and the
second occurrence continued the first one's counter. Wrap a subtree
that contains a Nondeterministic node in a marker carrying an id from
a driver-wide counter before serializing it, and strip the marker on
the executor before compiling, so each occurrence gets its own entry
while deterministic subtrees keep sharing compiled kernels.
@github-actions github-actions Bot added bug Something isn't working area:expressions Expression evaluation area:udf labels Oct 6, 2026

@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: Identical dispatched nondeterministic expressions shared a kernel, causing one occurrence to advance another’s counter or random generator.
  • Design approach: DispatchOccurrence adds a unique identifier to serialized subtrees containing Catalyst Nondeterministic nodes.
  • Correctness / compatibility analysis: Focused Spark 4.1.3 probes covered monotonic IDs, seeded random values, nested expressions and CodegenFallback children across eight batches. All matched Spark at this head. The same probes reproduced incorrect results on the base. Relevant Spark initialization semantics were checked across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0.
  • Key design decisions: Deterministic expressions retain cache sharing. Stateful kernels remain task-local. The marker is removed before code generation, adding no per-row evaluation work or protobuf changes.
  • Implementation sketch: Tag during serialization, retain the existing byte-based cache key, then unwrap before compiling. Three added tests cover serialization, query results and cache continuity.
  • Behavioral changes worth calling out: Compared with branch-1.1, repeated nondeterministic occurrences now maintain independent state while preserving continuity across batches. This is an intended correctness fix. User-function state without Catalyst Nondeterministic nodes remains outside this change.
  • Suggested improvements: None at the requested P1/P2 threshold. No introduced P1/P2 issues found within this review.

Reviewed all five changed files in the full diff from a8471cd92e5371fa8454cbc12e9cfdfc97beec44 to a34896d0363aae665741b421c2799c593699e0f7. Routed skills: review-comet-pr and review-comet-expression-pr. The PR is not a draft, and no existing reviews, comments or review threads were present.

Exact-head CI: labeling passed. Comet CI and CodeQL report action_required, providing no build/test verdict.

Validation limits: The isolated probes compiled current production sources against cached Spark 4.1.3 dependencies. Full native/JNI execution, the complete Comet and Spark SQL suites, and runtime tests on other Spark versions were not run.

@dwsmith1983

Copy link
Copy Markdown
Contributor Author

@sunchao I have resolve the merge conflicts which will kill the ci run. Sorry.

…ernel-per-occurrence

# Conflicts:
#	spark/src/main/scala/org/apache/comet/udf/codegen/CometScalaUDFCodegen.scala
#	spark/src/test/scala/org/apache/comet/CometCodegenSuite.scala

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

All of the new tests go through a Java or Scala UDF, but users can hit this without one, through a built-in that dispatches via CodegenDispatchFallback. With spark.sql.ansi.enabled=false, SELECT id, round(rand(1), 2) AS a, round(rand(1), 2) AS b FROM t returns different a and b on main, and this PR fixes it. With ANSI on, main happens to get it right because Round serializes a SQLQueryContext with its position in the SQL text, so the two copies already ship different bytes. ANSI off is the default on Spark 3.4 and 3.5, though. In 1.0 this query fell back to Spark because round on double had no dispatch path, so it's a 1.1.0 regression.

Could we add a SQL file test for it, something like expressions/math/round_nondeterministic_child.sql with -- Config: spark.sql.ansi.enabled=false, an 8-row table and that query? I tried that locally on 4.1. It fails on main and passes with this PR.

With ANSI off, Round carries no query context, so two copies of
round(rand(1), 2) in one projection serialize to the same bytes. Before
the per-occurrence tag they shared one dispatched kernel and the second
column continued the first one's random sequence. The SQL file test reads
both columns from one seeded partition and checks that round and rand
both go through the dispatcher.
@dwsmith1983

dwsmith1983 commented Oct 7, 2026 •

Copy link
Copy Markdown
Contributor Author

Could we add a SQL file test for it, something like expressions/math/round_nondeterministic_child.sql with -- Config: spark.sql.ansi.enabled=false, an 8-row table and that query?

Added in 1bfca75 as expressions/math/round_nondeterministic_child.sql, with ANSI off, an 8-row single-file table and your query. It uses expect_dispatch(round, rand), so it also fails if round over a double or the rand under it stops going through the dispatcher. On main it fails the way you saw: a matches Spark and b continues a's sequence (0.64 then 0.39 in the first row, where Spark has 0.64 twice). It passes on this branch on Spark 4.1 and 3.5.

…ernel-per-occurrence

# Conflicts:
#	spark/src/main/scala/org/apache/comet/udf/codegen/CometScalaUDFCodegen.scala
#	spark/src/test/scala/org/apache/comet/CometCodegenSuite.scala

@comphead comphead left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Thanks @dwsmith1983. The fix matches Spark's per-occurrence state for Catalyst Nondeterministic nodes, and from reading the code I did not find a regression. The thread notes that round(rand(1), 2) is a 1.1.0 regression, so this looks like a backport-1.1 candidate: branch-1.1 also keys its kernel cache on the serialized bytes, so the same tag and untag should apply there. One inline comment on folding the UDF test into the new SQL fixture.

}
}

test("identical non-deterministic dispatched expressions keep their own state") {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

With round_nondeterministic_child.sql in place, this test checks the same fix a second time. The javaId pair reaches the same tag in emitJvmCodegenDispatch as the fixture's round, and per the PR description the idPassthrough pair already passes without the fix. Separate Scala UDF kernels keeping their own counters across batches are already covered by Nondeterministic state persists across two ScalaUDFs in one task. Could the fixture take over the cross-batch part with a Config line setting spark.comet.batchSize=2, as if_expr.sql does, and this test be dropped? I have not run it, but from reading CometExecIterator that setting becomes the native session's batch size.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Could the fixture take over the cross-batch part with a Config line setting spark.comet.batchSize=2, as if_expr.sql does, and this test be dropped?

Done. The fixture now reads its 8 rows in four batches of 2, and the Scala test is gone. A kernel rebuilt for every batch now fails it (rows 2 to 7 repeat rows 0 and 1), which the fixture missed at the default batch size.

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

Reviewed the full six-file diff from 00a4b422f0ed6560ef76dce004b94c3697613108 to bdb059dcddbce6393be8b08e51f87b2313a157ce. The PR is not a draft. No introduced P1/P2 issues found within this review. Existing reviews, conversation comments, inline comments, and threads were checked. The requested SQL fixture and cross-batch coverage are present, with no substantiated existing blocker remaining.

Routed skills: review-comet-pr, review-comet-expression-pr, and review-comet-ffi-pr.

Summary

  • Prior state and problem: Identical dispatched subtrees containing Catalyst Nondeterministic nodes could share a cached kernel, causing one occurrence to advance another’s counter or random sequence.
  • Design approach: DispatchOccurrence adds a unique driver-generated identifier to the serialized subtree. Different occurrences consequently receive different cache keys.
  • Correctness: Focused probes using current production dispatcher and codegen sources matched Spark for monotonic IDs, seeded rand and randn, round(rand(1), 2), nested expressions, and CodegenFallback children. They covered four two-row batches in partitions 0 and 3. The base implementation reproduced the shared-state errors. Plan isolation and cleanup also passed.
  • Compatibility analysis: Verified the relevant Spark sources against versions 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. Their nondeterministic initialization and fallback contracts agree with this implementation. The change preserves return types, support levels, configuration defaults, and Arrow ownership rules.
  • Key design decisions: Identity is assigned during planning and travels with the expression bytes. Deterministic expressions remain unwrapped and retain cache sharing. Each nondeterministic occurrence preserves its state across batches.
  • Implementation sketch: Tag before closure serialization, compute the existing digest from those tagged bytes, and remove the marker after deserialization before compiling and initializing the kernel. The added tests cover serialization identity, cache continuity, and built-in SQL behavior.
  • Performance: The change adds a planning-time tree traversal and separate kernel state where correctness requires it. Unwrapping occurs only on cache misses and introduces no per-row marker evaluation. Deterministic cache sharing passed the probe. No P1/P2 performance regression was identified, and no timing benchmark was run.
  • Design: The change fits the existing task-local cache and plan cleanup lifecycle. It fixes occurrence identity without introducing another registry or changing the protobuf schema.
  • Abstraction & complexity: The private marker’s tag and untag operations keep the mechanism small and localized. Removing it before code generation keeps occurrence bookkeeping out of expression evaluation.
  • Behavioral changes worth calling out: Compared with branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1, repeated nondeterministic dispatched occurrences now maintain independent state. This is an intended correctness fix. Stateful user functions without Catalyst Nondeterministic descendants remain outside its scope.
  • Suggested improvements: No additional code change meets the P1/P2 reporting threshold. Exact-head CI and broader integration validation remain outstanding.

Exact-head CI: labeling passed. Comet CI, CodeQL, and Check PR Title report action_required, providing no build/test verdict.

Validation limits: The standalone probes compiled current production sources against cached Spark 4.1.3 dependencies with minimal allocator/type adapters. They did not exercise native/JNI execution, the complete CometCodegenSuite, the SQL fixture runner, or the Spark SQL matrix. Other supported versions received source comparison rather than runtime testing. The checkout remains unchanged.

…ernel-per-occurrence

# Conflicts:
#	spark/src/test/scala/org/apache/comet/CometCodegenSuite.scala
@dwsmith1983

Copy link
Copy Markdown
Contributor Author

@andygrove the round(rand(1)) fixture is in, and it now covers the cross-batch case with spark.comet.batchSize=2. Could you take another look? CI on aada46823 needs an approval to start.

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

Labels

area:expressions Expression evaluation area:udf bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Identical nondeterministic expressions share one codegen dispatcher kernel and continue each other's state

4 participants