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]
