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]

Reply via email to