Repository navigation
fix: stop the leak of Arrow memory from struct outputs of the codegen dispatcher - #6552
Conversation
… 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.
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.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Struct outputs leaked Arrow buffers because
StructVectoreagerly created children thatinitializeChildrenFromFieldssubsequently 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
CometArrowAllocatoras 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.
sunchao
left a comment
There was a problem hiding this comment.
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
StructVectorwithout 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:
childrenInitializedprevents constructor-time access to the export children. Allocator injection enables leak testing while preservingCometArrowAllocatoras 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.
|
@sunchao @andygrove Thanks! |
… 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>
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_jsonand a ScalaUDF that returns a case class or a tuple. The dispatcher is on by default: the default value ofspark.comet.exec.scalaUDF.codegen.enabledistrue.Arrow Java's
StructVectorcreates its writer in a field initializer:The
NullableStructWriterconstructor readsgetField().getChildren(). For each child, it callsaddOrGetand thenallocateNewSafe().allocateOutputcreates the struct withnew RenamedStructVector(field, allocator), andfieldcontains 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
allocateOutputcallsinitializeChildrenFromFields. This method adds the new children throughAbstractStructVector.addandputVector. With the defaultCONFLICT_REPLACEpolicy,putVectorremoves 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>. Forstruct<_1:bigint,_2:string>, the leak is 81,920 bytes for each batch. Onmain, the output vector gets its memory fromCometArrowAllocator. 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.evaluateArrow 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:
ListVectorandMapVectorcreate a writer only ingetWriter().What changes are included in this PR?
RenamedStructVectornow gives theStructVectorconstructor a field without children. Thus, the writer creates no child vectors, and the struct owns all of the memory that it allocates.getFieldreturns the export field only afterinitializeChildrenFromFields.allocateOutputdoes not close vectors that it did not create.allocateOutputhas a newallocatorparameter. The default value isCometArrowAllocator, so the current callers do not change. The new test uses this parameter to measure the memory that stays allocated afterclose().How are these changes tested?
A new test in
CometCodegenSuitedoes these steps for each output type:CometArrowAllocator.allocateOutput.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:close()struct<name:string,age:int>struct<_1:bigint,_2:string>struct<inner:struct<name:string,age:int>,tags:array<string>,attrs:map<string,int>>array<struct<name:string,age:int>>map<string,struct<name:string,age:int>>stringWith the fix, the value is 0 for all types, and these suites pass:
CometCodegenSuite, 105 tests.CometCodegenSourceSuite,CometCodegenFuzzSuite,CometCodegenHOFSuite,CometJsonJvmSuite,CometJsonExpressionSuiteandCometScalaUDFClassLoaderSuite.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.0andbranch-1.1have the same problem:allocateOutputcode.Thus, it is possible that these branches need a backport of this fix. The labels for this are
backport-1.0andbackport-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: theStructVectorconstructor gets a field without children. Thus, #5603 also stops the leak. This PR is a small fix formainand the release branches. If #5603 merges first, this PR will only add the regression test.