andygrove opened a new pull request, #6261: URL: https://github.com/apache/datafusion-comet/pull/6261
## Which issue does this PR close? Closes #2453. ## Rationale for this change A native plan with no JVM input, such as a native Parquet scan feeding a native sort, runs on a background Tokio task that sends its batches to the Spark task thread. `executePlan` discarded that task's `JoinHandle`, and `releasePlan` only dropped the receiving end of the channel. When the JVM consumer stopped early or threw, the task kept running until it next tried to send, and until then it kept the plan's stream and every memory reservation the stream held, after the Spark task had ended. Spark frees whatever a task still holds when the task ends and can hand it to other tasks, so that memory counted as free while native code was still using it. When the task finally dropped the stream, Spark logged `release called on N bytes but task only has 0 bytes of memory from the off-heap execution pool`, the warning in this issue. Tasks that DataFusion operators spawn leave memory behind the same way, on both the native-only and the JVM-fed paths. A sort's in-memory merge reads each sorted run through a task started by `spawn_buffered`, and each of those tasks holds the reservation for its run. Dropping the plan only aborts them, and they give the memory back the next time they yield, on a Tokio thread. The background producer arrived in #3553, after this issue was first filed, so it is a later cause of the same warning rather than the one first reported. ## What changes are included in this PR? - `BatchProducer` keeps the producer task's `JoinHandle` together with its channel. `releasePlan` drops the receiver, aborts the task and waits for it to finish before it flushes the final metrics. A task waiting on its input is cancelled straight away, and one that is producing a batch stops when that poll returns. A panic while dropping the stream comes back as an error, as it does when `releasePlan` drops a JVM-fed plan's stream itself. - `PlanMemoryPool` wraps the pool each plan reserves through and counts the bytes the plan holds. After dropping the plan, `releasePlan` waits for that count to reach zero, and logs a warning if it has not after one second. Unlike the producer, what still holds those reservations is not known at that point, and a reservation that is never released would otherwise hang the task. The aborted tasks give theirs back the next time they yield. - A paragraph in the threading section of the development guide. Because the producer is stopped before the final metrics are flushed, this also takes care of the ordering problem in #5504. #5505 stops the producer the same way but waits at most 100 ms, and for memory a bounded wait is not enough: a producer that is evaluating an expensive expression or merging spill files can take longer, and its reservations would still outlive the task. ## How are these changes tested? - Two tests in `CometExecIteratorLifecycleSuite` run a native sort over a native scan, read one row, stop, and check that the task holds no memory, from a task completion listener that runs after the one closing the plan. A UDF in the native projection above the sort stalls on the sort's second batch, so the producer is busy when the plan is released. One test keeps the sort in memory and the other makes it spill. Without the fix they fail with 1660594 and 139296 bytes still held, and Spark logs `Managed memory leak detected` for each task. - Native unit tests stop a producer that is waiting on its input, one that is blocked in the middle of a poll, and one whose stream panics when dropped, and check that a task the plan spawned keeps its reservation after the producer has stopped until `PlanMemoryPool` sees it returned. They also cover the pool's accounting and its deadline. - `CometExecSuite`, `CometNativeShuffleSuite`, `CometTaskMetricsSuite` and `CometExecIteratorLifecycleSuite` pass locally, and none of their plans hit the new one-second limit. `CometExecIterator.close()` still logs `closed with non-zero memory usage` in tests where one task runs several native plans, 89 times across `CometExecSuite` and `CometNativeShuffleSuite` both with and without this change. That has a separate cause: the task-shared pool charges every plan's reservations to the first plan's `CometTaskMemoryManager`, so when that plan closes first it reports memory the task's other plans still hold. -- 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]
