ranflarion opened a new issue, #25762:
URL: https://github.com/apache/datafusion/issues/25762

   ### Is your feature request related to a problem or challenge?
   
   `ExecutionPlan::execute` requires that "the Stream that is returned must 
ensure that any allocated resources are freed when the stream itself is 
dropped", and points operators to `SpawnedTask` for spawned work. But 
`SpawnedTask`'s drop, like tokio's `JoinSet`, only calls `abort()`, and tokio 
drops an aborted task's future the next time a worker gets to it: after the 
dropping task yields on a current-thread runtime, or some time later on a 
multi-thread one. Until then the task still owns its streams, its memory 
reservations and through them the `MemoryPool`, the `TaskContext` and any open 
spill files. Nothing tells the caller when they are gone, so a caller that 
needs them back can only poll.
   
   DataFusion's own tests do exactly that. `BlockingExec::refs` notes that 
"tokio might take some time to cancel spawned tasks, so you need to wrap this 
check into a retry loop", and 11 cancellation tests use 
`assert_strong_count_converges_to_zero`, which polls every 10 ms for up to 10 
s. On main, 4 of them (`test_drop_cancel` in `analyze`, `coalesce_partitions` 
and `repartition`, and `record_batch_receiver_stream_drop_cancel`) still see 1 
or 2 references right after the drop and pass only on the check after the first 
10 ms sleep. #22016 measured an 80 ms cancellation delay through 
`CoalescePartitionsExec` and `RepartitionExec`, and #22017 removed the plan 
reference behind it, which leaves the abort-only drop.
   
   ### Describe the solution you'd like
   
   clippy's `disallowed-methods` already routes every spawn through 
`SpawnedTask` and `datafusion_common_runtime::JoinSet` (#6513), so the 
mechanism fits in those wrappers, as opt-in constructors: 
`SpawnedTask::spawn_reclaimable`, and `JoinSet::spawn_reclaimable` and 
`spawn_reclaimable_on`. Such a task's future lives in a slot that the worker 
locks while polling it. When the handle is dropped, or `abort_all` or 
`shutdown` is called on the set, it takes the future out and drops it right 
there if no worker holds the lock, then aborts the task. If a worker is polling 
the task at that moment it aborts as today, and the future is dropped as soon 
as that poll returns.
   
   The in-place drop runs the future's destructor on the dropping thread, 
inside the task's runtime but outside the task, so `tokio::task::try_id()` 
returns `None` there (or the dropping task's ID), and tokio has no public API 
to enter a task's context. That is why it is opt-in: `spawn` and `spawn_on` 
keep their current semantics everywhere, including on 
`RecordBatchReceiverStreamBuilder`, and only the tasks that drive a plan's 
streams opt in, through new `spawn_reclaimable` builder methods used by 
`spawn_buffered` (and so by sorts and `SortPreservingMergeExec`) and by the 
builder's input runner behind `CoalescePartitionsExec` and `AnalyzeExec`, and 
through `SpawnedTask::spawn_reclaimable` for `RepartitionExec`'s input tasks. A 
stream's destructor can't rely on running inside a task today either, since the 
same stream is dropped on whatever thread drops the plan whenever its parent 
does not spawn it, and `spawn_buffered` does not spawn at all on a 
current-thread runtime. No DataFusio
 n code reads the task ID. An injected `JoinSetTracer` wraps the task rather 
than the reclaimed future, so tokio still drops the tracer inside the task.
   
   Reclaimable tasks hand out no `AbortHandle` (the `JoinSet` constructors 
return the task's ID instead), so apart from their runtime shutting down, they 
are cancelled only through their own `SpawnedTask` or `JoinSet`, which drops 
the future before aborting the task. Nothing can observe such a task finish 
while its destructor still runs, and no runtime worker ever waits on a 
destructor running on another thread, since the handle holds the slot only 
while it takes the future out. A panic from the destructor is raised again when 
tokio drops the task, so joining reports `JoinError::Panic` as before, with one 
exception: if the runtime shuts down while `JoinSet::abort_all` is dropping a 
task's future, the set can still be joined afterwards and reports that task as 
cancelled even if its destructor panicked. The in-place drop is skipped while 
the thread is already unwinding. The added cost for opted-in tasks is an 
allocation per spawn and one uncontended mutex lock per poll, and I'll includ
 e benchmark results with the PR.
   
   I have this on a branch with tests. With the stream tasks opted in, all 11 
cancellation tests above see no references left at the first check, and the 
physical-plan (2320 tests), `core_integration` (1189) and sqllogictest suites 
pass.
   
   ### Describe alternatives you've considered
   
   Waiting for the aborted tasks inside `Drop` doesn't work, since a drop can't 
await and blocking there can deadlock when it runs on a worker. Making in-place 
drop the default for every spawn would be simpler, but it changes where 
arbitrary futures' destructors run and what task context they see, which is why 
it is limited to stream tasks. An awaitable signal that a dropped plan's tasks 
are gone would help callers that can await, but it still makes every caller 
wait, and it could be added on top of this for the mid-poll case. Not spawning 
at particular call sites, for example the per-batch sort runs behind 
`spawn_buffered`, removes those tasks for embedders that need it, but gives up 
their parallelism and only covers the sites it touches.
   
   ### Additional context
   
   The case this can't cover is a task that is mid-poll when its handle is 
dropped, since its future can't be taken while a worker is polling it. In 
practice that is CPU work in flight, for example a `spawn_buffered` sort run 
that is sorting its batch when the plan is dropped, which keeps its reservation 
until that sort finishes. A caller that needs release to be synchronous even 
then has to keep that work off separate tasks, which I'd raise separately.
   


-- 
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