ranflarion opened a new pull request, #25763: URL: https://github.com/apache/datafusion/pull/25763
## Which issue does this PR close? - Closes #25762. ## Rationale for this change Dropping a plan's stream only aborts the tasks that drive its inputs, so what they own (streams, memory reservations and the pool behind them) is released whenever the runtime next gets to them, and a caller that needs those resources back, DataFusion's own cancellation tests included, can only poll for them. The issue has the details. ## What changes are included in this PR? `datafusion-common-runtime` gets opt-in reclaimable spawns: `SpawnedTask::spawn_reclaimable`, and `JoinSet::spawn_reclaimable` and `spawn_reclaimable_on`. The task's future lives in a slot the worker locks while polling it. Dropping the handle, or calling `abort_all` or `shutdown` on the set, takes the future out and drops it on the spot when no worker is polling it, then aborts the task. Otherwise it aborts as today, and the future is dropped when that poll returns. - The in-place drop runs the future's destructor inside the task's runtime but outside the task, where `tokio::task::try_id()` is `None`. That is why it is opt-in, and the constructors document it. `spawn` and `spawn_on` are unchanged. - Reclaimable spawns hand out no `AbortHandle`, and the `JoinSet` constructors return the task's `Id` instead. So apart from runtime shutdown, only the owner cancels them, and it finishes the drop before aborting. The handle holds the slot only while it takes the future out, so no worker ever waits on a destructor running elsewhere. - A panic from the destructor is raised again when tokio drops the task, so joining reports `JoinError::Panic`. The exception is a runtime shutdown during `abort_all`, and it is documented. - An injected `JoinSetTracer` wraps the task rather than the reclaimed future, so tokio still drops the tracer inside the task. `JoinSet::join_all` joins through the wrapper, so dropping it reclaims before aborting. In `datafusion-physical-plan`, `ReceiverStreamBuilder` and `RecordBatchReceiverStreamBuilder` get `spawn_reclaimable` and `spawn_reclaimable_on`, and their existing `spawn` and `spawn_on` are unchanged. The tasks that drive plan streams use the new methods: `spawn_buffered`, the builder's internal input runner behind `CoalescePartitionsExec` and `AnalyzeExec`, and `RepartitionExec`'s input tasks. ## Are these changes tested? Yes. New unit tests in `datafusion-common-runtime` cover: - an idle task's future is gone by the time `drop` returns - a task that is mid-poll is released when that poll returns - the in-place destructor runs in the runtime but outside the task, while plain `spawn` keeps its task context - destructor panics are contained on drop and reported through `abort_all` and `join_next` - runtime shutdown does not wait for a destructor running on another thread - dropping `join_all` reclaims before aborting - joined tasks are forgotten and detached tasks keep running An integration test in its own binary, since the tracer is process-wide, checks that an injected tracer is still dropped inside its task. The existing cancellation tests exercise the operator opt-ins: with this change, all 11 that use `assert_strong_count_converges_to_zero` see no references left at their first check, where 4 of them needed its retry before. Benchmarks: `dfbench` built with `--release` from main (6a792c671) and from this branch. TPC-H SF10 parquet on local NVMe, on a 32 vCPU i4i.8xlarge with default target partitions, 5 iterations per run, with runs ordered main, branch, branch, main. `compare.py` on each adjacent pair: | | main | this PR | `compare.py` | |---|---|---|---| | `tpch` SF10, round 1 | 6701.9 ms | 6672.1 ms | 21 no change, Q19 1.06x slower | | `tpch` SF10, round 2 | 6702.5 ms | 6642.0 ms | 21 no change, Q9 1.09x faster | | `sort_tpch` SF10, round 1 | 52790.4 ms | 52775.9 ms | 11 no change | | `sort_tpch` SF10, round 2 | 52808.3 ms | 52664.6 ms | 11 no change | The same binary varied by up to 4.6% per query between its two runs, and the two flagged queries differ between rounds and point in opposite directions, so I read them as noise. ## Are there any user-facing changes? New public APIs: `SpawnedTask::spawn_reclaimable`, `JoinSet::spawn_reclaimable` and `spawn_reclaimable_on`, and `spawn_reclaimable` and `spawn_reclaimable_on` on both stream builders. Existing spawn methods keep their behavior. The tasks DataFusion spawns to drive streams under `spawn_buffered`, `CoalescePartitionsExec`, `AnalyzeExec` and `RepartitionExec` now drop their futures on the thread that drops the stream when idle. So destructors of streams under those operators may run outside a tokio task, as they already do whenever a stream is not spawned. -- 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]
