Skip to content

fix: initialize dispatched kernels with the partition the native plan computes - #6700

Merged
comphead merged 2 commits into
apache:mainfrom
comphead:fix-6570-dispatcher-partition-index
Oct 6, 2026
Merged

comphead merged 2 commits into
apache:mainfrom
comphead:fix-6570-dispatcher-partition-index

Conversation

@comphead

@comphead comphead commented Oct 5, 2026

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #6570.

Rationale for this change

The JVM codegen dispatcher initialized each kernel with TaskContext.partitionId(). Spark and Comet's native expressions use the index of the partition being computed instead. The two differ under a union, a coalesce and a cartesian product, so partition-seeded expressions such as rand, uuid, monotonically_increasing_id and spark_partition_id inside a dispatched expression returned different values from Spark. The dispatcher also kept one kernel per task, so under a coalesce every parent partition after the first reused the first one's kernel and its random state.

Default settings reach this through common expressions. round on a double and regexp_replace dispatch their whole subtree, so round(rand(42) * 100, 2) and regexp_replace(uuid(), '-', '') are affected. The issue has a repro.

What changes are included in this PR?

  • The native planner passes its partition index and plan id (exec_context_id) to JvmScalarUdfExpr, which hands them to CometUdfBridge.evaluate.
  • CometUDF gets an evaluate overload that receives both. Its default calls the existing method, so other CometUDF implementations are unaffected.
  • CometScalaUDFCodegen initializes kernels with that partition index and caches nondeterministic kernels per plan. Deterministic kernels never read the index, so one of them still serves every plan in a task, and a coalesce over many partitions does not recompile them.

How are these changes tested?

A new test in CometCodegenSuite compares a dispatched map(1, spark_partition_id()) and round(rand(42), 6) with Spark under a union, a coalesce and a cross join. Equivalent queries return different answers from Spark on main, see the repro in #6570.

@comphead comphead 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 Oct 5, 2026
@github-actions github-actions Bot added bug Something isn't working area:scan Parquet scan / data reading area:ffi Arrow FFI / JNI boundary 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: Dispatched expressions used the task partition ID and shared random state across parent partitions, producing incorrect results under union, coalesce, and Cartesian execution.
  • Design approach: Pass the native partition index and plan ID through JNI, then isolate nondeterministic kernels by plan.
  • Correctness / compatibility analysis: Partition initialization matches the relevant Spark 3.4–4.2 sources. Focused tests confirm state continuity across batches and isolation between plans.
  • Key design decisions: Deterministic kernels remain shared. The default CometUDF.evaluate overload preserves existing implementations. These are straightforward changes, but plan-specific entries lack cleanup when their plans finish.
  • Implementation sketch: The planner supplies both identifiers to JvmScalarUdfExpr, the bridge forwards them, and the dispatcher uses them for initialization and caching.
  • Behavioral changes worth calling out: Compared with branch-1.1 at 7b7eec69282a7abc33a0e8e687ca594c72f140ec, partition-dependent results are intentionally corrected. Retaining completed kernels introduces the measured heap regression below.
  • Suggested improvements: Release nondeterministic entries when their native plan closes, preserving live plans' state and deterministic sharing. Request changes for this P2.

Reviewed all 10 changed files at 0450646b460fe59074be54db44cf3dbd7153d918 against d16f7bbb7ef5aa566626bc98b6655785724b4da4. The PR is not a draft. Existing discussions and threads contain no duplicate concern. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr, and review-comet-memory-pr.

Exact-head CI: all executed checks passed, including Linux Rust/JVM tests, TPC-H/TPC-DS validation, and Spark 4.1 catalyst/core/hive suites. Other Spark SQL versions, macOS, and Iceberg suites were skipped.

Local validation: 68 focused tests passed using freshly compiled changed JVM sources and the verified exact-head CI native artifact. Additional probes passed for partition IDs, random sequences, UUIDs, monotonic IDs, deterministic sharing, and overload delegation. No clean native rebuild, other-version runtime tests, or heap-exhaustion test was performed.

val entry =
CometScalaUDFCodegen.CacheEntry(compiled, kernel, boundExpr.dataType, outputField)
kernelCache.put(key, entry)
kernelCache.put(if (boundExpr.deterministic) sharedKey else key, entry)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] Release nondeterministic cache entries when their native plan finishes. Every parent partition consumed by coalesce gets a fresh planId, so this insertion retains another serialized expression, deserialized closure, and kernel. CometUdfBridge removes the dispatcher only at task completion, while CometExecIterator.close() does not remove these entries. A bounded 32-partition query retained approximately 32 MiB in cache-key payloads alone after consuming all parents. The base dispatcher retained one approximately 1 MiB payload for the same expression. Consequently, coalescing many partitions now accumulates JVM heap proportional to every completed plan, with additional closure copies beyond those measured bytes. Tie entry cleanup to native-plan lifetime while preserving state for live or interleaved plans.

Evidence: Reproduced on Spark 4.1.3/JDK 21 with the exact-head native artifact. Define a serializable Lookup(val data: Array[Byte]) extends (Int => Int) whose apply(i) returns data(Math.floorMod(i, data.length)).toInt. Evaluate val f = udf(new Lookup(Array.fill[Byte](1024 * 1024)(7))) and spark.range(0, 32, 1, 32).select(f(spark_partition_id()).as("x")).coalesce(1). Inside rdd.mapPartitions, exhaust the rows, then inspect the current task's dispatcher cache before task completion. The plan used CometCoalesce over CometProject and retained 32 entries containing 33,794,624 serialized-key bytes. A separate identical-expression comparison against the dispatcher freshly compiled from base d16f7bb measured head: 32 entries/33,651,968 bytes; base: 1 entry/1,051,624 bytes. Reproduction source and logs: /tmp/6700-pinned-review-4qldjlpo/PR6700ReviewSuite.scala and tests.log.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

I agree with sunchao that this needs fixing here. On this branch, round(rand(42), 6) under coalesce(1) over 16 partitions compiled 16 kernels, and nothing removes them before the task ends. CometExecIterator.close() already has the plan id, because it is the iterator's id. Could it drop that plan's entries right after nativeLib.releasePlan(plan)? For example, a CometUdfBridge.releasePlan(taskAttemptId, id) that calls a hook only the dispatcher implements. A test that runs a coalesce over several partitions and then checks what the cache holds would cover the cleanup. It would also cover the claim that deterministic kernels are shared across plans, which no test checks today. The caching diagram at the top of the class needs updating as well. It still says the kernel cache is keyed on the bound expression and input shapes, for the whole task.

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

Requesting changes for two things in CometScalaUDFCodegen, both inline. The per-plan kernels sunchao flagged stay in the cache until the task ends. And the shared-key decision misses a nondeterministic child under an Invoke, so make_valid_utf8 over spark_partition_id() is still wrong under a coalesce. The rest checks out. The index comes from split.index, plan ids come from newIterId and never repeat within an executor, and the new test fails if I make every kernel shared again.

partitionIndex: Int): CometScalaUDFCodegen.CacheEntry = {
assert(Thread.holdsLock(this), "lookupOrCompile must run under this.synchronized")
kernelCache.get(key) match {
// A deterministic kernel never reads the partition index, so one instance under the planless

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.

The shared-key decision reads boundExpr.deterministic, and that is only the root's flag. Spark doesn't make it transitive everywhere. Invoke.deterministic is isDeterministic && arguments.forall(_.deterministic) and skips targetObject. Spark 4 lowers make_valid_utf8(x) and is_valid_utf8(x) to Invoke(x, "makeValid") and Invoke(x, "isValid"), and CometInvoke dispatches the whole subtree. So a partition-seeded child under one of them still goes into the shared kernel. I ran this on the branch:

spark.range(0, 8, 1, 2)
  .selectExpr(
    "id",
    "make_valid_utf8(cast(spark_partition_id() AS STRING)) AS p",
    "make_valid_utf8(cast(round(rand(42), 6) AS STRING)) AS r")
  .coalesce(1)

The four rows from the second partition come back with p = 0, and with the first partition's rand sequence continued. Spark returns p = 1 for them. The same query under a union matches. Could the check walk the tree instead, for example boundExpr.exists(!_.deterministic)? With that change the query matches Spark locally, and a deterministic expression still compiles once across 16 coalesced plans. Adding this shape to the new test, guarded with isSpark40Plus, would pin it.

val entry =
CometScalaUDFCodegen.CacheEntry(compiled, kernel, boundExpr.dataType, outputField)
kernelCache.put(key, entry)
kernelCache.put(if (boundExpr.deterministic) sharedKey else key, entry)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

I agree with sunchao that this needs fixing here. On this branch, round(rand(42), 6) under coalesce(1) over 16 partitions compiled 16 kernels, and nothing removes them before the task ends. CometExecIterator.close() already has the plan id, because it is the iterator's id. Could it drop that plan's entries right after nativeLib.releasePlan(plan)? For example, a CometUdfBridge.releasePlan(taskAttemptId, id) that calls a hook only the dispatcher implements. A test that runs a coalesce over several partitions and then checks what the cache holds would cover the cleanup. It would also cover the claim that deterministic kernels are shared across plans, which no test checks today. The caching diagram at the top of the class needs updating as well. It still says the kernel cache is keyed on the bound expression and input shapes, for the whole task.

@andygrove andygrove added the backport-1.1 Candidate for backporting to 1.1 release branch label 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: Dispatched kernels used task partition IDs and shared nondeterministic state across parent partitions, producing incorrect results under union, coalesce, and Cartesian execution.
  • Design approach: Pass the native partition index and plan ID through JNI and isolate nondeterministic kernels by plan.
  • Correctness / compatibility analysis: The initialization matches Spark 3.4–4.2 sources. The existing Invoke concern still reproduces: under coalesce, make_valid_utf8(cast(spark_partition_id() AS STRING)) returns 0 for the second partition where Spark returns 1. Random state also continues incorrectly.
  • Key design decisions: Sharing deterministic kernels avoids repeated compilation. The default CometUDF.evaluate overload preserves existing implementations. The design is straightforward, but cache lifetime and descendant nondeterminism need the corrections already requested.
  • Implementation sketch: PhysicalPlanner supplies both identifiers to JvmScalarUdfExpr, the bridge forwards them, and the dispatcher initializes and caches kernels accordingly.
  • Behavioral changes worth calling out: Compared with branch-1.1 at 7b7eec69282a7abc33a0e8e687ca594c72f140ec, correcting partition-dependent results is intended. The existing heap regression is confirmed: a bounded dispatcher comparison retained 32 entries and approximately 32 MiB of serialized keys, versus one entry and approximately 1 MiB at the base.
  • Suggested improvements: Resolve the existing cache cleanup thread and Invoke nondeterminism thread. Release completed plans’ kernels and inspect descendant nondeterminism before sharing.

Reviewed the entire 10-file diff at 0450646b460fe59074be54db44cf3dbd7153d918 against d16f7bbb7ef5aa566626bc98b6655785724b4da4. The PR remains non-draft. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr, and review-comet-memory-pr.

Exact-head CI: 38 successful checks and 31 skipped, with no failures or pending checks. Linux Rust/JVM tests, TPC-H/TPC-DS validation, Spark 4.1 SQL suites, and required checks passed. Other Spark SQL versions, macOS, and Iceberg suites were skipped.

Validation: 68 focused checks and six direct probes passed, including reproductions confirming the existing failures. Changed JVM sources were freshly compiled using cached dependencies and the digest-verified exact-head native CI artifact. Runtime validation used Spark 4.1.3/JDK 17. No clean full Maven/Rust rebuild, other-version runtime matrix, or heap-exhaustion test was performed. Project code remains unchanged.

No additional introduced P1/P2 issues found within this review beyond the existing comments. Existing blockers remain unresolved.

… computes

The JVM codegen dispatcher seeded kernels from `TaskContext.partitionId()`, which differs from
the partition being computed under a union, a coalesce and a cartesian product. It also kept one
kernel per task, so the parent partitions of a coalesce shared random state.

The native planner now passes its partition index and plan id through the UDF bridge.
Nondeterministic kernels are cached per plan and initialized with that index. Deterministic
kernels never read it and stay shared across the plans of a task.
… determinism

A coalesce kept the nondeterministic kernel of every parent partition until the task ended.
`CometExecIterator.close` now calls `CometUdfBridge.releasePlan`, which hands the closed plan to a
new `CometUDF.releasePlan` hook, and the dispatcher drops that plan's kernels there. Kernels
shared by every plan stay until the task ends.

The dispatcher shared a kernel whenever the root expression reported deterministic, but `Invoke`
leaves its target object out of `deterministic`. So `make_valid_utf8` over `spark_partition_id()`
still shared one kernel across the plans of a coalesce. It now checks every node.
@comphead
comphead force-pushed the fix-6570-dispatcher-partition-index branch from 0450646 to bf1f69a Compare October 6, 2026 16:12

@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: Dispatched kernels used task partition IDs and shared nondeterministic state across parent partitions, producing incorrect results under union, coalesce and Cartesian execution.
  • Design approach: Pass the computed partition index and native plan ID through JNI, then isolate nondeterministic kernels by plan.
  • Correctness / compatibility analysis: The initialization matches the relevant Spark 3.4–4.2 sources. Tests confirm that the existing Invoke determinism and completed-plan retention concerns are addressed.
  • Key design decisions: Deterministic kernels remain shared, and default CometUDF methods preserve existing implementations. The lifetime model is straightforward, but the additional cache lookup introduces the measured CPU regression below.
  • Implementation sketch: The planner supplies both identifiers, the bridge forwards them, and CometExecIterator.close() releases the plan’s cached kernels.
  • Behavioral changes worth calling out: Compared with branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1, correcting partition-dependent results is intended. Large nondeterministic closures also incur an unintended per-batch slowdown.
  • Suggested improvements: Reuse the closure hash across shared and per-plan lookups. Request changes for this P2.

Reviewed the entire 11-file diff at bf1f69ab4fe07e0be0b43567d06aca6535109059 against 8c783aa88104616dcf0f7876b4f7a9ba71e2bf31. The PR remains non-draft. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr, and review-comet-memory-pr.

Exact-head CI at the last check: 17 successful, 13 skipped and 8 running, with no failures. Rust tests, JVM suites, TPC validation and Spark 4.1 SQL preparation remain running.

Validation: Full JVM reactor test-compile, 67 focused tests and six direct probes passed. Integration tests used the digest-verified native artifact from this exact head. The 32-parent cleanup reproduction retained zero completed-plan entries. Performance comparisons used identical head builds with the dispatcher implementation varied. Runtime validation used Spark 4.1.3/JDK 17. No local clean Rust rebuild, full Spark SQL suite or other-version runtime matrix was performed.

// key serves every plan in the task. Only a kernel with a nondeterministic node is stored per
// plan.
val sharedKey = key.copy(planId = CometScalaUDFCodegen.NoPlan)
kernelCache.get(sharedKey).orElse(kernelCache.get(key)) match {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

[P2] Avoid hashing large closures twice on every nondeterministic cache hit. For these kernels, get(sharedKey) always misses before get(key) succeeds, and each case-class hash traverses the entire serialized closure through ByteBuffer.hashCode(). This adds another full closure scan on every batch, even within one native plan where the base already returned correct results. A bounded native query with a captured 1 MiB lookup array slowed from a 271 ms median with the base dispatcher to 407 ms at this head, approximately 50%. Reusing the bytes hash in a disposable variant reduced this to 298 ms. Could the two lookups reuse a precomputed closure/schema hash while preserving plan isolation and deterministic sharing?

Evidence: Reproduced with Spark 4.1.3/JDK 17 and the verified exact-head native artifact. Define a serializable Lookup(val data: Array[Byte]) extends (Int => Int) whose apply(i) returns data(Math.floorMod(i, data.length)).toInt. Set val f = udf(new Lookup(Array.fill[Byte](1048576)(7))), then collect spark.range(0, 1048576, 1, 1).select(f(spark_partition_id()).as("x")).agg(sum("x")). After six warmups, nine samples gave medians of 271.107 ms with the dispatcher freshly compiled from base 8c783aa, 407.237 ms with head, and 297.807 ms with the disposable hash-reuse variant. All returned 7,340,032. Only the dispatcher varied between otherwise-identical head builds. A direct 4096-row dispatcher benchmark also measured approximately 1.25 versus 2.28 ms per batch. Sources and logs are under /tmp/pr6700-bf1f69ab-review/, including PR6700EndToEndCostSuite.scala, PR6700CacheCost.scala, and native-cost-*.log.

@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, both are fixed. CometExecIterator.close() now drops the closed plan's kernels through releasePlan, and the shared-key decision walks the whole tree. I checked that the new tests catch each fix. Without the releasePlan call, the cache still holds the kernels of the four coalesced plans after the rows are read. With the root-only check, the make_valid_utf8 query returns p = 0 for the second partition again. If every kernel is stored per plan, the cache test fails as well, so it also pins the sharing of deterministic kernels. Both new tests pass locally on Spark 3.5, 4.0 and 4.1.

@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: Dispatched kernels used task partition IDs and shared nondeterministic state across parent partitions, producing incorrect results under union, coalesce and Cartesian execution.
  • Design approach: Pass the computed partition index and native plan ID through JNI, isolate nondeterministic kernels by plan, and release them when that plan closes.
  • Correctness / compatibility analysis: Initialization matches the relevant Spark 3.4–4.2 sources. Focused tests confirm that the earlier nested Invoke nondeterminism and completed-plan retention concerns are addressed.
  • Key design decisions: Deterministic kernels remain shared. Default CometUDF methods preserve existing implementations. The lifecycle hooks fit the existing design without introducing a separate cache-management framework.
  • Implementation sketch: PhysicalPlanner supplies both identifiers to JvmScalarUdfExpr. The bridge forwards them, and CometExecIterator.close() invokes cache cleanup after native teardown.
  • Behavioral changes worth calling out: Compared with branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1, correcting partition-dependent results is intended. The existing P2 hashing regression remains reproducible: a single-partition query using a UDF with a captured 1 MiB lookup array took a 403 ms median at head versus 275 ms with the base dispatcher. Both returned identical results using the same exact-head native artifact.
  • Suggested improvements: Address the existing redundant closure-hashing concern by reusing the closure/schema hash while preserving plan isolation and deterministic sharing. This remains an unresolved blocker.

Reviewed the entire 11-file, two-commit diff at bf1f69ab4fe07e0be0b43567d06aca6535109059 against 8c783aa88104616dcf0f7876b4f7a9ba71e2bf31. The PR remains non-draft. Existing reviews, comments and threads were read. Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr, and review-comet-memory-pr.

Exact-head CI: 35 successful checks and 15 skipped, with no failures or pending checks. Required checks, Linux Rust/JVM tests, TPC validation, and Spark 4.1 catalyst/core/hive suites passed. Other Spark SQL versions, macOS and Iceberg suites were skipped.

Validation: JVM reactor test-compile, 67 focused tests and six state/lifecycle probes passed on Spark 4.1.3/JDK 17. The 32-parent cleanup probe retained zero completed-plan entries. The native artifact was verified against the exact-head CI archive digest. Benchmark sources and logs are in /tmp/pr6700-review-final-d17ik8bj/. No clean local Rust rebuild or other-version runtime matrix was performed. Full Spark SQL validation relied on CI. Project source remains unchanged.

No additional introduced P1/P2 issues found within this review beyond the existing performance comment.

@comphead
comphead added this pull request to the merge queue Oct 6, 2026
Merged via the queue into apache:main with commit 8d0bb01 Oct 6, 2026
50 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:ffi Arrow FFI / JNI boundary area:scan Parquet scan / data reading area:udf backport-1.1 Candidate for backporting to 1.1 release branch bug Something isn't working 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.

Codegen dispatcher initializes kernels with an incorrect partition index under UNION ALL and coalesce

3 participants