andygrove opened a new pull request, #6374:
URL: https://github.com/apache/datafusion-comet/pull/6374
## Which issue does this PR close?
Closes #6255.
## Rationale for this change
All the native plans in a Spark task reserve memory in one native memory
pool. `acquire_task_shared_pool` (`task_shared.rs`) only runs its `create`
closure when the task has no live pool, so the pool keeps the
`CometTaskMemoryManager` passed with whichever plan created it, and every later
plan's manager goes unused. `CometExecIterator` built a manager per plan and,
in `close()`, warned when that manager's `getUsed` was non-zero. So the task's
first plan reported the memory of every plan in the task and the others always
reported 0. When the first plan closed while a sibling still held memory, it
logged a leak that was not one, and a real leak in any later plan was never
reported.
This happens whenever a task runs more than one native plan, for example a
JVM sort-merge join over two native sorts, or a native Parquet or Iceberg write
over a Comet child. With the repro from the issue, 8 of the join stage's 10
tasks logged `CometExecIterator closed with non-zero memory usage` on current
main.
## What changes are included in this PR?
The fix is on the JVM side only. No native code changes.
- `CometExecIterator.taskMemory` gives every native plan in a task one
shared `CometTaskMemoryManager`, matching the task-shared native pool. The
task's first plan creates it, keyed by the `TaskContext`, and a task completion
listener removes it when the task ends, whether or not the task's plans were
closed.
- Each `CometExecIterator` counts itself as open once its completion
listener is registered, so its `close()` is guaranteed to run. `close()` checks
the shared manager's usage only when it closes the task's last open plan. The
warning now names the task and reports what the task still holds after its last
native plan closed. Memory that outlives one particular plan is still reported
per plan by the native `releasePlan` warning added in #6261.
- The `CometTaskMemoryManager` docs and the task-shared pool section of the
memory management contributor guide describe the shared manager.
- `CometExecIteratorLifecycleSuite.withTaskContext` now completes its
synthetic task, which runs the task completion listeners as the end of a real
task does.
## How are these changes tested?
New tests in `CometExecIteratorLifecycleSuite`:
- `a native plan that closes while another plan in its task holds memory
does not warn` zips two native plans into one Spark task. The second plan's
native sort takes in its input and holds it, and the first plan is then read to
its end, which closes it. The test asserts that the sort held memory at that
point and that no non-zero memory usage warning was logged. Without the fix it
fails with `List("CometExecIterator closed with non-zero memory usage :
1270122") was not empty`.
- `the last native plan in a task to close warns about the memory the task
still holds` opens two plans in one task and acquires 1234 bytes through the
task's manager, standing in for native memory that outlives them. Closing the
first plan logs nothing, and closing the second logs exactly one warning
reporting 1234 bytes.
- `a task's memory manager is released when the task ends, even with a plan
left open` leaves a plan open, ends the task, and checks with a weak reference
that the manager can be collected. With the completion listener's removal
disabled, this test fails.
Run on the default profile (Spark 4.1, Scala 2.13):
- `CometExecIteratorLifecycleSuite` and `CometTaskMemoryManagerSuite`: 21
tests pass.
- `CometJoinSuite` and `CometParquetWriterSuite`: 103 tests pass, and the
run added no `closed with non-zero memory usage` lines to `unit-tests.log`.
- The issue's repro (`spark.comet.exec.sortMergeJoin.enabled=false`,
broadcast joins disabled, keys 0 to 999 joined with 0 to 2,000,000), run as a
scratch test that is not part of this PR: 8 `closed with non-zero memory usage`
warnings in `unit-tests.log` before the fix, 0 after.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]