Skip to content

fix: return a native plan's memory before its Spark task ends - #6261

Merged
andygrove merged 4 commits into
apache:mainfrom
andygrove:release-plan-stop-tokio-task
Sep 29, 2026
Merged

andygrove merged 4 commits into
apache:mainfrom
andygrove:release-plan-stop-tokio-task

Conversation

@andygrove

@andygrove andygrove commented Sep 26, 2026 •

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #2453.
Closes #5504.

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 plan's stream where releasePlan can reach it, and the producer task locks it only while polling it. Before it flushes the final metrics, releasePlan closes the channel, aborts the task, takes the stream and drops it on the Spark task thread, as it already does with a JVM-fed plan's stream. Taking the stream waits only for a poll the task is already in. Waiting for the task to be cancelled instead would wait for a free Tokio worker, and every worker can be tied up, for instance waiting in Spark's acquireMemory for the memory the stream holds. A panic while dropping the stream comes back as an error.
  • 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 stream, 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 as soon as a worker runs them, which the next change keeps possible while Spark blocks an acquire.
  • SparkMemory calls Spark's acquireMemory inside tokio::task::block_in_place. Spark blocks that call when the task has to wait for other tasks to release memory, and on a Tokio worker block_in_place hands the worker's other tasks to another thread meanwhile, so they keep running. That includes aborted tasks that only need to be cancelled to release the memory being waited for. On the Spark task thread, where JVM-fed plans run, it only steps out of the runtime context. In a micro-benchmark with an acquire that takes 1 µs, it added about 0.1 to 0.3 µs per call for one task, and up to about 2.4 µs per call with eight tasks acquiring at once on a 28-core machine.
  • Paragraphs in the development guide's threading section and its rules for native code.

Because the stream is dropped before the final metrics are flushed, this also takes care of the ordering problem in #5504. #5505 aborts the producer and waits at most 100 ms for it. For memory, a bounded wait for the cancellation is not enough, since a producer that is evaluating an expensive expression or merging spill files can take longer and its reservations would outlive the task, and an unbounded one can deadlock, since the cancellation needs a free worker.

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, one whose stream panics when dropped, and one whose runtime's only worker is blocked until the stopped plan's memory comes back, 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, and check that an acquire which Spark blocks until another task runs completes on a one-worker runtime.
  • CometExecSuite, CometNativeShuffleSuite, CometTaskMetricsSuite and CometExecIteratorLifecycleSuite pass locally, and none of their plans hit the new one-second limit. With block_in_place they took 2 min 34 s, against 2 min 57 s without it.

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.

A plan with no JVM input runs on a Tokio task that sends its batches to
the Spark task thread. releasePlan only dropped the receiving end of the
channel, so when the JVM consumer stopped early or threw, the task kept
the plan's stream, and the memory the stream had reserved, after the
Spark task ended. Spark handed that memory to other tasks while native
code still used it, and logged "release called on N bytes but task only
has 0 bytes" when the stream was finally dropped.

releasePlan now aborts the task and waits for it to finish before
flushing the final metrics. Tasks that DataFusion operators spawn, such
as the ones a sort's merge reads its sorted runs through, are only
aborted when the plan is dropped, so releasePlan also waits, for up to a
second, until every reservation the plan made has been returned.
PlanMemoryPool counts the bytes each plan holds.

Closes apache#2453
@andygrove andygrove added run-iceberg-tests run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue labels Sep 26, 2026
@github-actions github-actions Bot added bug Something isn't working area:memory Memory pools, reservations, OOM handling labels Sep 26, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Native producer tasks and operator-spawned tasks could retain reservations after Spark reclaimed the task's memory.
  • Design approach: Retain the producer's JoinHandle and track outstanding reservations through PlanMemoryPool.
  • Correctness / compatibility analysis: Found one introduced P1 shutdown deadlock. Spark's synchronous memory-acquisition wait can occupy the Tokio workers needed to complete the new join. Verified the relevant Spark behavior in 3.4.3, 3.5.9, 4.0.4, 4.1.3 and branch-4.2 sources.
  • Key design decisions: Per-plan accounting preserves the shared pool's allocation policy. The localized wrappers are straightforward, but an unconditional join requires cancellation to remain runnable under memory pressure.
  • Implementation sketch: releasePlan stops the producer, flushes metrics, drops the context and waits up to one second for remaining reservations.
  • Behavioral changes worth calling out: Teardown now waits for background work. The reproduced performance impact is an indefinite task stall before the reservation timeout is reached.
  • Suggested improvements: Make blocking Spark memory acquisition compatible with Tokio scheduling before joining the producer, and add the saturated-worker regression described below.

Reviewed the entire five-file diff from bc4be39964cbe9cdb5f2a949740a8164e6b5755b to 18ddc8f0912be10963e29a64168edd921e5a79d9. The PR is not a draft. Read AGENTS.md and applied review-comet-pr, review-comet-ffi-pr and review-comet-memory-pr. Snapshot and live discussion contained no existing reviews or comments.

Exact-head CI: Linux build, Rust tests, Comet suites and TPC checks passed. The execution group passed 1,078 tests, including both new lifecycle tests. Spark 4.1 SQL and Iceberg 1.11 extensions jobs remained running. Two label-run aggregate checks failed because upstream jobs were cancelled.

Validation: All seven new native unit tests passed locally. A disposable harness using the exact producer implementation reproduced the deadlock and verified the base and blocking-aware controls. This was a bounded scheduler reproduction, not an end-to-end Spark deadlock test. Scala suites were verified through CI rather than rerun locally. Project source remains unchanged.

Comment thread native/core/src/execution/jni_api.rs Outdated
} = self;
drop(batches);
task.abort();
match runtime.block_on(task) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P1] Keep producer cancellation runnable before joining it. With COMET_WORKER_THREADS=1 and two Spark tasks, task A can stop early while its native producer still holds off-heap reservations. Task B can occupy the sole Tokio worker inside Spark's ExecutionMemoryPool.acquireMemory, waiting for A to release memory. task.abort() only schedules cancellation, so this join waits for a worker that cannot become available until A releases its reservations. A also cannot return to Spark's final task-memory cleanup. Both tasks therefore hang, and the later one-second timeout is never reached. The base teardown returned and allowed Spark cleanup to unblock B. Please make potentially blocking JNI memory acquisition blocking-aware, or otherwise guarantee cancellation can execute while preserving the memory-release ordering.

Evidence: Compiled a disposable harness with byte-for-byte current BatchProducer code and the pinned Tokio 1.53.1 dependencies. On a one-worker runtime, producer A yields one batch then remains pending while retaining memory. Another task occupies the worker with a synchronous wait for that memory. Calling stop() from a separate thread remained blocked through the 300 ms observation window and completed only after externally releasing the wait. The base receiver-drop/detach control returned immediately. Wrapping acquisition in tokio::task::block_in_place also allowed the head implementation to finish. Reproduction: timeout 30 /tmp/comet-6261-current-review/lifecycle-tests review_stop_with_saturated_workers_and_memory_wait --nocapture. Production source confirms JniMemoryManager::acquire calls Spark synchronously without a blocking-aware wrapper, Spark's acquisition uses lock.wait(), and executor memory cleanup follows task completion.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reproduced: stopping_a_batch_producer_does_not_need_a_free_worker holds the only worker of a one-worker runtime in a synchronous wait for the stopped plan's memory, and on 18ddc8f stop() timed out waiting for a free worker.

9d405f3 removes the dependence on a worker instead of making the acquire blocking-aware. The task locks the stream only while polling it, and stop() takes the stream and drops it on the calling thread, as releasePlan already does for a JVM-fed plan. That waits only for a poll the task is already in, which is already running on a worker. The task is still aborted, but no longer joined. The test passes now.

Tasks that DataFusion operators spawn, such as the sort merge's per-run tasks, still need a worker to be cancelled. If every worker is blocked waiting on the memory they hold, the PlanMemoryPool wait gives up after one second, the task ends, and Spark's cleanup wakes the waiters, so that case falls back to the late release instead of deadlocking. Wrapping the JNI acquire in block_in_place would cover it too, but it hands the worker off to another thread on every acquire, so I'd rather do that separately if the fallback turns out to matter.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Since your approval, 5d24232 also makes the acquire blocking-aware, as you first suggested. SparkMemory now calls Spark's acquireMemory inside block_in_place, so a worker blocked there hands its other tasks to another thread, and aborted DataFusion tasks holding the memory it waits for can still be cancelled. a_blocked_acquire_leaves_its_worker_running_other_tasks covers it: on a one-worker runtime, an acquire that waits for another task to run times out without the change and completes with it.

A micro-benchmark with an acquire that takes 1 µs measured about 0.1 to 0.3 µs added per call for one task, and up to about 2.4 µs with eight tasks acquiring at once. The Comet suites I ran took no longer. Stopping the producer still drops the stream on the calling thread, so releasing a plan does not depend on the scheduler either way.

Joining the producer after aborting it needs a free Tokio worker to run
the cancellation. Spark's acquireMemory blocks the thread that calls it
until other tasks free memory, so when every worker was waiting on memory
the stopped plan's stream held, the join never finished and neither task
could end.

The producer task now locks the stream only while polling it, and
releasePlan takes the stream and drops it on the calling thread, as it
already does for a JVM-fed plan. That waits only for a poll already in
progress. The task is aborted but no longer joined.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Background producers and operator-spawned tasks could retain native reservations after Spark reclaimed task memory.
  • Design approach: BatchProducer makes the stream accessible to teardown. PlanMemoryPool tracks each plan’s outstanding reservations while preserving the underlying allocation policy.
  • Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Checked Spark memory-acquisition and task-completion behavior in 3.4.3, 3.5.9, 4.0.4, 4.1.3 and branch-4.2. The previously reported worker-starvation deadlock is addressed, and its regression test passes.
  • Key design decisions: The two localized wrappers separate producer lifetime from reservation accounting. Pool operations retain the existing task-shared allocation and spilling semantics.
  • Implementation sketch: releasePlan closes the receiver, aborts the producer, takes and drops its stream, flushes metrics, drops the context, then waits for remaining reservations.
  • Behavioral changes worth calling out: Teardown waits for an active stream poll and then up to one second for outstanding reservations. The latter timeout logs a warning and allows cleanup to proceed. Synchronization is added around stream polls and reservation accounting. Throughput was not benchmarked.
  • Suggested improvements: None meeting the P1/P2 evidence bar.

Reviewed all five changed files and both commits from base bc4be39964cbe9cdb5f2a949740a8164e6b5755b through full head 9d405f39ef379d27cb7460cb20ab3faf18b839b5. The PR remains non-draft. Read AGENTS.md, existing reviews, comments and the review thread. Applied review-comet-pr, review-comet-ffi-pr and review-comet-memory-pr.

Exact-head CI: 64 checks passed and 10 were skipped, with no failures. Passing coverage includes Linux/Rust tests, Comet suites, Spark 4.1 SQL and Iceberg suites. The execution group passed 1,099 tests, including both new lifecycle tests, with no reservation-timeout warnings. macOS and Spark 3.4/3.5/4.0 SQL checks were skipped.

Validation limits: All eight added Rust tests passed locally in a disposable harness using the exact changed implementations and pinned DataFusion/Tokio versions. This validates native lifecycle behavior without the JNI integration. Full native/JNI builds and Scala suites were not rerun locally. Project files remain unchanged, and nothing was published.

Spark's acquireMemory blocks the calling thread until other tasks
release memory, and Comet calls it on Tokio workers. A worker blocked
there ran none of its other tasks, which could include the aborted tasks
of a released plan whose memory the acquire was waiting for. The release
then waited out its one-second limit, and the memory went back to Spark
after its task had ended.

SparkMemory now makes the call inside tokio::task::block_in_place, which
on a worker hands the worker's other tasks to another thread while the
call blocks. On the Spark task thread, where JVM-fed plans run, it only
steps out of the runtime context.
@andygrove

Copy link
Copy Markdown
Member Author

I also checked this at the JVM level. With COMET_WORKER_THREADS=1, local[4] and 96m or 128m of off-heap memory, a native scan feeding sortWithinPartitions hung on main in 4 runs out of 4. The only worker was parked in ExecutionMemoryPool.acquireMemory waiting for 1/2N of the pool, with all four task threads in executePlan. With this PR merged onto main it passed 11 runs out of 11, in about 5 s. One of those runs hit the same wait and got through it.

A side effect worth knowing about: because each acquire can hand the worker's core to another thread, a one-worker runtime now runs several spawned plans at once. The same query without memory pressure took 13.0 s on main with one worker and 7.7 to 8.6 s with this PR, against 5.5 s with four workers either way.

This PR doesn't change one related thing, and I've filed it separately as #6294. A producer dropped by the runtime shutting down still looks like end of stream to executePlan, so the task finishes with whatever it had.

@andygrove
andygrove added this pull request to the merge queue Sep 28, 2026
@github-merge-queue
github-merge-queue Bot removed this pull request from the merge queue due to a conflict with the base branch Sep 28, 2026
…io-task

# Conflicts:
#	spark/src/test/scala/org/apache/spark/CometExecIteratorLifecycleSuite.scala
@andygrove
andygrove enabled auto-merge September 28, 2026 17:45
@andygrove
andygrove added this pull request to the merge queue Sep 28, 2026
Merged via the queue into apache:main with commit 8805e3f Sep 29, 2026
73 checks passed
andygrove added a commit to andygrove/datafusion-comet that referenced this pull request Sep 29, 2026
Resolve the conflict with apache#6261 on the memory_pools import in jni_api.rs by
importing both overcommit and PlanMemoryPool. The registry still holds the
pool that create_memory_pool returns rather than the PlanMemoryPool that
wraps it, so the memory usage log still reads each pool's overcommit through
the task-shared and tracking wrappers.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:memory Memory pools, reservations, OOM handling bug Something isn't working run-iceberg-tests run-spark-4.1-tests Run the Spark 4.1 SQL tests on this pull request instead of waiting for the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Cancel background batch producers before collecting final plan metrics ExecutionMemoryPool errors releasing more memory than allocated

2 participants