Describe the bug
All the native plans in a Spark task share one memory pool, and that pool keeps the CometTaskMemoryManager of whichever plan created it. acquire_task_shared_pool calls its create closure only when the task has no live pool (task_shared.rs#L100-L118), so every later plan's manager is dropped unused (mod.rs#L39-L72). The first plan's getUsed therefore reports the whole task's native reservations, and every other plan's reports 0.
CometExecIterator.close() warns when its manager's usage isn't zero (CometExecIterator.scala#L352-L355). When the first plan closes while another plan still holds memory, it logs a leak that isn't one. A real leak in any later plan is never reported.
Steps to reproduce
A JVM sort-merge join over two native sorts puts two native plans in each task. Set spark.comet.exec.sortMergeJoin.enabled=false and disable broadcast joins. Then join a table with keys 0 to 999 against one with keys 0 to 2,000,000. In 7 of 10 tasks, the left plan logged a warning like this one when it closed while the right side's sort still held its batches:
WARN CometExecIterator: CometExecIterator closed with non-zero memory usage : 33576604
Expected behavior
The warning fires only when the task's native memory outlives all of its plans, and it reports that amount.
Additional context
CometNativeWriteExec and CometIcebergWriteExec over a Comet child also put two plans in a task, and the child plan is created first. Once #6247 reserves the writer's buffers, every such write will log this warning when the child finishes. This was the third "Not investigated" item in #5212. One way to fix it would be to keep one CometTaskMemoryManager per task attempt on the JVM side, and check for non-zero usage when the task's last plan closes.
Describe the bug
All the native plans in a Spark task share one memory pool, and that pool keeps the
CometTaskMemoryManagerof whichever plan created it.acquire_task_shared_poolcalls itscreateclosure only when the task has no live pool (task_shared.rs#L100-L118), so every later plan's manager is dropped unused (mod.rs#L39-L72). The first plan'sgetUsedtherefore reports the whole task's native reservations, and every other plan's reports 0.CometExecIterator.close()warns when its manager's usage isn't zero (CometExecIterator.scala#L352-L355). When the first plan closes while another plan still holds memory, it logs a leak that isn't one. A real leak in any later plan is never reported.Steps to reproduce
A JVM sort-merge join over two native sorts puts two native plans in each task. Set
spark.comet.exec.sortMergeJoin.enabled=falseand disable broadcast joins. Then join a table with keys 0 to 999 against one with keys 0 to 2,000,000. In 7 of 10 tasks, the left plan logged a warning like this one when it closed while the right side's sort still held its batches:Expected behavior
The warning fires only when the task's native memory outlives all of its plans, and it reports that amount.
Additional context
CometNativeWriteExecandCometIcebergWriteExecover a Comet child also put two plans in a task, and the child plan is created first. Once #6247 reserves the writer's buffers, every such write will log this warning when the child finishes. This was the third "Not investigated" item in #5212. One way to fix it would be to keep oneCometTaskMemoryManagerper task attempt on the JVM side, and check for non-zero usage when the task's last plan closes.