Repository navigation
fix: give each nondeterministic dispatched expression its own kernel - #6712
dwsmith1983 wants to merge 12 commits into
Conversation
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.
sunchao
left a comment
There was a problem hiding this comment.
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:
DispatchOccurrenceadds a unique identifier to serialized subtrees containing CatalystNondeterministicnodes. - Correctness / compatibility analysis: Focused Spark 4.1.3 probes covered monotonic IDs, seeded random values, nested expressions and
CodegenFallbackchildren 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 CatalystNondeterministicnodes 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.
|
@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
left a comment
There was a problem hiding this comment.
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.
Added in 1bfca75 as |
…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
left a comment
There was a problem hiding this comment.
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") { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
Could the fixture take over the cross-batch part with a
Configline settingspark.comet.batchSize=2, asif_expr.sqldoes, 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.
…op the duplicate UDF test
sunchao
left a comment
There was a problem hiding this comment.
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
Nondeterministicnodes could share a cached kernel, causing one occurrence to advance another’s counter or random sequence. - Design approach:
DispatchOccurrenceadds 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
randandrandn,round(rand(1), 2), nested expressions, andCodegenFallbackchildren. 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
taganduntagoperations 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.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1, repeated nondeterministic dispatched occurrences now maintain independent state. This is an intended correctness fix. Stateful user functions without CatalystNondeterministicdescendants 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
|
@andygrove the |
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
Nondeterministicnode such asmonotonically_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 returneda = 0..7, b = 8..15per batch where Spark returnsa = 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.emitJvmCodegenDispatchwraps the bound subtree in it before serializing when the tree contains aNondeterministicnode, 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.lookupOrCompilestrips 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
ScalaUDFmarked nondeterministic,Invoke,StaticInvoke) are not CatalystNondeterministicnodes and are not wrapped. #5526 covers them separately.How are these changes tested?
expressions/math/round_nondeterministic_child.sqlrunsSELECT id, round(rand(1), 2) AS a, round(rand(1), 2) AS bover an 8-row table written as one file. ANSI is off, soRoundcarries no query context and the two copies serialize to the same bytes.spark.comet.batchSize=2splits 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 thatroundandrandran in the dispatcher. Without the occurrence tag it fails, sincebcontinuesa'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 occurrenceinCometCodegenSuite: 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 batchesinCometCodegenSuite: two ids give two cache entries, a repeated id gives one hit, and the counter continues across batches for one id.The fixture and
CometCodegenSuitepass on the Spark 4.1 and 3.5 profiles.