Skip to content

fix: [branch-1.1] stop the leak of Arrow memory from struct outputs of the codegen dispatcher (#6552) - #6663

Merged
andygrove merged 1 commit into
apache:branch-1.1from
andygrove:backport-6552-branch-1.1
Oct 5, 2026
Merged

andygrove merged 1 commit into
apache:branch-1.1from
andygrove:backport-6552-branch-1.1

Conversation

@andygrove

Copy link
Copy Markdown
Member

Which issue does this PR close?

None. This is the branch-1.1 backport of #6552, which didn't close an issue either.

Rationale for this change

The JVM codegen dispatcher leaks off-heap Arrow memory for every output batch whose type is a top-level struct, such as the output of from_json or of a ScalaUDF that returns a case class or a tuple. The leak is 64 KiB per batch for struct<string, int> and 80 KiB for struct<bigint, string>. RenamedStructVector's parent constructor builds a writer that allocates a child vector for each field, and initializeChildrenFromFields then replaces those children without closing them. The buffers come from CometArrowAllocator, which has no limit and is never closed, so nothing reports the leak, and the executor's off-heap memory just grows.

branch-1.1 has the same allocateOutput code and the same Arrow version (18.3.0), and the dispatcher is on by default here too (spark.comet.exec.scalaUDF.codegen.enabled). #6552 has the details.

#6552 also has the backport-1.0 label, and the commit applies to branch-1.0 without conflicts. That backport isn't open yet. Under the backporting guide's newer-branch rule, it shouldn't merge before this one.

What changes are included in this PR?

A cherry-pick (-x) of #6552's commit. It applied without conflicts, and the patch is line-for-line the same as upstream's.

  • RenamedStructVector passes the StructVector constructor a field without children, so the writer creates no child vectors and the struct owns all the memory it allocates. getField returns the export field once initializeChildrenFromFields has run.
  • allocateOutput takes an allocator parameter that defaults to CometArrowAllocator. Existing callers are unchanged, and the new test uses the parameter to measure what a closed output leaves allocated.

How are these changes tested?

On this branch, with the default Spark 4.1 profile and JDK 17:

  • The codegen dispatcher suites all pass, 202 tests including the new one: CometCodegenSuite (106), CometCodegenSourceSuite (60), CometCodegenFuzzSuite (28), CometCodegenHOFSuite (5) and CometSpecializedGettersDispatchSuite (3).
  • CometJsonJvmSuite, CometJsonExpressionSuite and CometScalaUDFClassLoaderSuite, which cover from_json and case-class UDF outputs: 15 pass.
  • Without the fix, the new test fails on this branch. I put RenamedStructVector back to this branch's version and kept the new allocator parameter, which the test needs. The test then fails with Memory was leaked by query. Memory leaked: (65536) on struct<name: string, age: int>, so this branch has the leak and the test catches it.

I didn't run the other Spark profiles locally. CI runs them on this branch.

… dispatcher (apache#6552)

* fix: stop struct outputs of the codegen dispatcher from leaking Arrow memory

Arrow Java's StructVector builds its NullableStructWriter in a field
initializer, and that writer creates and allocates a child vector for every
child of the struct's field. allocateOutput constructs its struct from the
export field, which carries the children, and then calls
initializeChildrenFromFields. Under the default CONFLICT_REPLACE policy that
drops the writer's children from the struct without closing them, so their
buffers outlive the output vector. Every struct-typed output batch of the
dispatcher leaked them: 64 KiB for struct<string, int> and 80 KiB for
struct<bigint, string>.

Close the writer's children before initializeChildrenFromFields replaces them.
allocateOutput now takes the allocator, defaulting to CometArrowAllocator, so a
test can check that a closed output leaves nothing allocated.

* Write the new comments in Simplified Technical English

* Remove the -ing forms from the new test name and assert message

* Make RenamedStructVector own all of its memory

Give the StructVector constructor a field without children. Thus, the writer
creates no child vectors, and nothing must close vectors that another object
created. getField returns the export field only after
initializeChildrenFromFields. allocateOutput goes back to its original shape.

* Shorten the new comments

(cherry picked from commit 3bc2faa)
@andygrove andygrove added bug Something isn't working area:expressions Expression evaluation labels Oct 5, 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: StructVector eagerly allocated children that subsequent initialization replaced without closing, leaking off-heap memory per output batch.
  • Design approach: Construct the struct without children, initialize them once, then expose the original export field.
  • Correctness / compatibility analysis: Checked Arrow 18.3.0 initialization and ownership against Spark’s struct-writing and UDF semantics across supported 3.4–4.2 versions. Field names, nullability, child ordering, and exported values remain consistent. No introduced P1/P2 issues found within this review.
  • Key design decisions: The constructor-phase getField guard addresses Arrow’s virtual call during construction with a small, localized state change. It adds no row-level work or payload copying.
  • Implementation sketch: Thread an allocator through allocateOutput, adjust RenamedStructVector, and add allocator-accounting regression coverage for scalar and nested outputs.
  • Behavioral changes worth calling out: A disposable probe using the exact-head output implementation passed 96 allocation/write/Spark-reference/FFI/close cases, plus allocation-failure cleanup. Restoring the base constructor reproduced 65,536 bytes retained after close.
  • Suggested improvements: None meeting the P1/P2 evidence bar.

Reviewed the entire three-file diff from 992c806a7e38c2e88bd018aa5774164b0850e1fa to ba7a5ba2ae14557e3a19bbd1d0725cb2f262c018. The PR is non-draft. Snapshot and live discussions contained no existing feedback or substantiated blockers.

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

Exact-head CI: Preflight, title, and labeling checks passed. The current CI run still has build and integration jobs queued. Two earlier label-run aggregate checks failed because their prerequisite jobs were cancelled, as confirmed in their logs. No completed integration-test verdict is available.

Validation limits: Local runtime checks used Spark 4.1.3, Arrow 18.3.0, and JDK 21 in an isolated harness. Full Comet/native builds, Comet suites, Spark SQL suites, other Spark runtime profiles, and throughput benchmarks were not run.

@andygrove
andygrove merged commit 7b7eec6 into apache:branch-1.1 Oct 5, 2026
173 of 181 checks passed
@andygrove
andygrove deleted the backport-6552-branch-1.1 branch October 5, 2026 19:23
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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants