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]

Reply via email to