Repository navigation
fix: [branch-1.1] stop the leak of Arrow memory from struct outputs of the codegen dispatcher (#6552) - #6663
Conversation
… 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)
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem:
StructVectoreagerly 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
getFieldguard 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, adjustRenamedStructVector, 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.
Which issue does this PR close?
None. This is the
branch-1.1backport 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_jsonor of a ScalaUDF that returns a case class or a tuple. The leak is 64 KiB per batch forstruct<string, int>and 80 KiB forstruct<bigint, string>.RenamedStructVector's parent constructor builds a writer that allocates a child vector for each field, andinitializeChildrenFromFieldsthen replaces those children without closing them. The buffers come fromCometArrowAllocator, which has no limit and is never closed, so nothing reports the leak, and the executor's off-heap memory just grows.branch-1.1has the sameallocateOutputcode 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.0label, and the commit applies tobranch-1.0without 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.RenamedStructVectorpasses theStructVectorconstructor a field without children, so the writer creates no child vectors and the struct owns all the memory it allocates.getFieldreturns the export field onceinitializeChildrenFromFieldshas run.allocateOutputtakes anallocatorparameter that defaults toCometArrowAllocator. 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:
CometCodegenSuite(106),CometCodegenSourceSuite(60),CometCodegenFuzzSuite(28),CometCodegenHOFSuite(5) andCometSpecializedGettersDispatchSuite(3).CometJsonJvmSuite,CometJsonExpressionSuiteandCometScalaUDFClassLoaderSuite, which coverfrom_jsonand case-class UDF outputs: 15 pass.RenamedStructVectorback to this branch's version and kept the newallocatorparameter, which the test needs. The test then fails withMemory was leaked by query. Memory leaked: (65536)onstruct<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.