Skip to content

fix: stop the leak of Arrow memory from struct outputs of the codegen dispatcher - #6552

Merged
andygrove merged 5 commits into
apache:mainfrom
peterxcli:fix/codegen-struct-output-child-leak
Oct 4, 2026
Merged

andygrove merged 5 commits into
apache:mainfrom
peterxcli:fix/codegen-struct-output-child-leak

Conversation

@peterxcli

@peterxcli peterxcli commented Oct 2, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #.

Rationale for this change

The JVM codegen dispatcher leaks off-heap Arrow memory for each batch when the output type is a top-level struct. Examples are from_json and a ScalaUDF that returns a case class or a tuple. The dispatcher is on by default: the default value of spark.comet.exec.scalaUDF.codegen.enabled is true.

Arrow Java's StructVector creates its writer in a field initializer:

private final NullableStructWriter writer = new NullableStructWriter(this);

The NullableStructWriter constructor reads getField().getChildren(). For each child, it calls addOrGet and then allocateNewSafe(). allocateOutput creates the struct with new RenamedStructVector(field, allocator), and field contains the children. Thus, when the constructor returns, the struct already has one child vector for each field. Each of these child vectors has buffers with the default capacity.

Then allocateOutput calls initializeChildrenFromFields. This method adds the new children through AbstractStructVector.add and putVector. With the default CONFLICT_REPLACE policy, putVector removes the old children from the struct, but it does not close them. When Comet closes the struct, the struct closes only the new children. Only the writer keeps a reference to the old children, and nothing releases their buffers.

The leak is 65,536 bytes for each batch of struct<name:string,age:int>. For struct<_1:bigint,_2:string>, the leak is 81,920 bytes for each batch. On main, the output vector gets its memory from CometArrowAllocator. This root allocator has no limit, and Comet never closes it. Thus, no error message shows the leak.

The per-task allocator from #5027 made the leak visible. With -Darrow.memory.debug.allocator=true, all of the ledgers that stay open have this allocation stack:

NullableStructWriter.<init> ← StructVector.<init> ← RenamedStructVector.<init> ← CometBatchKernelCodegenOutput.allocateOutput ← CometScalaUDFCodegen.evaluate

Arrow 19.0.0 creates the writer in the same way. Thus, an upgrade to Arrow 19.0.0 does not fix the leak.

The leak does not affect List and Map outputs:

  • ListVector and MapVector create a writer only in getWriter().
  • When they create a nested struct, they use a field type without children. Thus, the writer of the nested struct creates no child vectors.

What changes are included in this PR?

  • RenamedStructVector now gives the StructVector constructor a field without children. Thus, the writer creates no child vectors, and the struct owns all of the memory that it allocates. getField returns the export field only after initializeChildrenFromFields. allocateOutput does not close vectors that it did not create.
  • allocateOutput has a new allocator parameter. The default value is CometArrowAllocator, so the current callers do not change. The new test uses this parameter to measure the memory that stays allocated after close().

How are these changes tested?

A new test in CometCodegenSuite does these steps for each output type:

  1. Create a child allocator of CometArrowAllocator.
  2. Allocate an output vector from this allocator with allocateOutput.
  3. Close the output vector.
  4. Make sure that the allocator has 0 bytes allocated.

Without the fix, the test fails on the first type with this error: Memory was leaked by query. Memory leaked: (65536).

This table shows the bytes that stay allocated after close() without the fix:

Output type Bytes allocated after close()
struct<name:string,age:int> 65,536
struct<_1:bigint,_2:string> 81,920
struct<inner:struct<name:string,age:int>,tags:array<string>,attrs:map<string,int>> 34,304
array<struct<name:string,age:int>> 0
map<string,struct<name:string,age:int>> 0
string 0

With the fix, the value is 0 for all types, and these suites pass:

  • CometCodegenSuite, 105 tests.
  • The other codegen dispatcher suites, 108 tests: CometCodegenSourceSuite, CometCodegenFuzzSuite, CometCodegenHOFSuite, CometJsonJvmSuite, CometJsonExpressionSuite and CometScalaUDFClassLoaderSuite.

The fix also stops the unused allocation. For a 4-row struct<name:string,age:int> output, the peak allocation is 65,634 bytes without the fix. With the fix, it is 98 bytes.

Backport to release branches

branch-1.0 and branch-1.1 have the same problem:

  • They contain the same allocateOutput code.
  • They use Arrow 18.3.0.
  • The dispatcher is on by default.

Thus, it is possible that these branches need a backport of this fix. The labels for this are backport-1.0 and backport-1.1. Refer to Backporting to Release Branches. The commit applies to the two branches without conflicts.

Relation to #5603

#5603 moves this allocation into NativeUtil.createVector. It uses the same approach: the StructVector constructor gets a field without children. Thus, #5603 also stops the leak. This PR is a small fix for main and the release branches. If #5603 merges first, this PR will only add the regression test.

… 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.
@github-actions github-actions Bot added bug Something isn't working area:expressions Expression evaluation labels Oct 2, 2026
@peterxcli peterxcli changed the title fix: stop struct outputs of the codegen dispatcher from leaking Arrow memory fix: stop the leak of Arrow memory from struct outputs of the codegen dispatcher Oct 2, 2026
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.
@peterxcli
peterxcli marked this pull request as ready for review October 4, 2026 13:40

@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: Struct outputs leaked Arrow buffers because StructVector eagerly created children that initializeChildrenFromFields subsequently replaced without closing.
  • Design approach: Construct the struct without children, then expose the cached export schema after child initialization.
  • Correctness / compatibility analysis: The change preserves field names, nullability, nested values, and FFI ownership. Relevant Spark sources were checked across supported 3.4–4.2 versions. No introduced P1/P2 issues found within this review.
  • Key design decisions: The initialization flag keeps the constructor from observing the export children. Allocator injection enables direct leak checks while retaining CometArrowAllocator as the production default.
  • Implementation sketch: Update the allocation helper and forwarding method, adjust RenamedStructVector, and add regression coverage for struct, array, map, and scalar outputs. The implementation remains localized.
  • Behavioral changes worth calling out: The baseline retained 65,536, 81,920, and 34,304 bytes for the three tested struct shapes. The fixed helper released all allocations. It removes unused allocations without adding copies or changing expression evaluation.
  • Suggested improvements: None meeting the P1/P2 reporting threshold.

Reviewed full SHA 945c80dc33d97c136c16846c06566d3fb059c630 against supplied base b0f0de3da544f811b452415438062951271671f2. The full supplied PR patch matches the merge-base-relative diff and includes all five branch commits. The additional 14 files in the direct base-to-head comparison are unchanged since the branch fork and reflect four newer main commits. No existing reviews, issue comments, inline comments, or review threads were present.

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

Exact-head CI: 15 checks succeeded, 9 remained running, and 14 were skipped. Native compilation passed. Spark 4.1 Comet suites, Rust tests, TPC-H/TPC-DS checks, and two lint jobs were still running. Spark SQL and Iceberg suites were skipped. No failures were reported at the final check, but CI was not yet conclusive.

Validation: Compiled the exact output helper against Spark 4.1.3 and Arrow 18.3.0. All 96 bounded allocation/write/Spark-reference/FFI/close cases passed, plus allocation-failure cleanup. Coverage included nulls, empty values, nested collections, and batches through 8,193 rows. Full Maven suites and native Spark integration were not run locally. Other Spark versions received source-level compatibility checks.

@andygrove
andygrove self-requested a review October 4, 2026 14:38

@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: Top-level struct outputs leaked Arrow buffers because eager writer initialization created children that subsequent initialization replaced without closing.
  • Design approach: Construct StructVector without children, then expose the cached export field after child initialization.
  • Correctness / compatibility analysis: Schema, null handling, nested values, and buffer ownership remain consistent in the reviewed paths. Relevant Spark sources were checked across supported 3.4–4.2 versions. No introduced P1/P2 issues found within this review.
  • Key design decisions: childrenInitialized prevents constructor-time access to the export children. Allocator injection enables leak testing while preserving CometArrowAllocator as the production default.
  • Implementation sketch: Update the allocation helper and forwarding method, adjust RenamedStructVector, and add regression coverage for six output shapes. The change remains localized without additional abstraction.
  • Behavioral changes worth calling out: The baseline retained 65,536, 81,920, and 34,304 bytes for the three tested struct shapes. The fixed helper released all allocations. It eliminates unused allocations without adding copies or per-row work.
  • Suggested improvements: None meeting the P1/P2 reporting threshold.

Reviewed full SHA 945c80dc33d97c136c16846c06566d3fb059c630 against supplied base b0f0de3da544f811b452415438062951271671f2. The supplied PR patch exactly matches the full merge-base-relative diff, covering all five branch commits and three changed files. Additional direct base-to-head differences reflect four newer commits on main, not PR changes.

Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr, and review-comet-memory-pr. Read the existing approval and confirmed there are no issue comments, inline comments, or review threads containing unresolved concerns.

Exact-head CI: 25 checks succeeded and 15 were skipped, with no failures. Native build, Rust tests, Spark 4.1 Comet suites, and TPC-H/TPC-DS verification passed. The expressions job log confirms that the new leak regression and existing struct-output test passed. Spark SQL, Iceberg, and macOS suites were skipped.

Validation: Recompiled the exact output helper and baseline in an isolated Spark 4.1.3/Arrow 18.3.0 harness. All 96 allocation/write/Spark-reference/Arrow C Data round-trip/close cases passed, plus allocation-failure cleanup. Cases included nulls, empty structs, nested collections, and batches through 8,193 rows. Full Maven/native integration suites were not rerun locally. Other Spark versions received source-level checks, not local runtime validation.

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

@andygrove
andygrove added this pull request to the merge queue Oct 4, 2026
@andygrove andygrove added backport-1.0 Candidate for backporting to 1.0 release branch backport-1.1 Candidate for backporting to 1.1 release branch labels Oct 4, 2026
Merged via the queue into apache:main with commit 3bc2faa Oct 4, 2026
78 checks passed
@peterxcli
peterxcli deleted the fix/codegen-struct-output-child-leak branch October 4, 2026 16:44
@peterxcli

Copy link
Copy Markdown
Member Author

@sunchao @andygrove Thanks!

andygrove added a commit that referenced this pull request Oct 5, 2026
… dispatcher (#6552) (#6663)

* 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)

Co-authored-by: Peter Lee <peterxcli@gmail.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:expressions Expression evaluation backport-1.0 Candidate for backporting to 1.0 release branch backport-1.1 Candidate for backporting to 1.1 release branch bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants