Skip to content

[prototype] asyncio control plane for BlockAllocationTaskScheduler - #1073

Draft
jan-janssen wants to merge 4 commits into
mainfrom
asyncio-blockallocation-prototype
Draft

jan-janssen wants to merge 4 commits into
mainfrom
asyncio-blockallocation-prototype

Conversation

@jan-janssen

Copy link
Copy Markdown
Member

Experimental prototype to evaluate whether zmq.asyncio can cut the administrative thread overhead of the interactive block-allocation backend. Opt-in only via EXECUTORLIB_ASYNCIO=1; the default threaded path is unchanged.

Motivation

BlockAllocationTaskScheduler starts one Python thread per worker. Each thread spends its life blocked, either in future_queue.get() while idle or in SocketInterface.receive_dict() while a task runs. N workers therefore cost N administrative threads.

Approach

  • One private event loop per executor, running in its own background thread. Each worker is a coroutine on that loop instead of a thread.
  • One dispatcher thread blocks on the existing queue.Queue and hands an item to a worker coroutine only once that worker is idle. Tasks therefore stay in the queue until picked up, so cancel_futures, the shutdown sentinels and max_workers resizing keep their current behaviour.
  • AsyncSocketInterface is a zmq.asyncio subclass of SocketInterface with the same wire protocol, so the worker side is unchanged. All interfaces of one executor share one ZMQ context.
  • Jupyter-safe: no asyncio.run, no run_until_complete on the caller's loop, no nest_asyncio. The caller's event loop is never touched.
  • Unchanged API: users still get concurrent.futures.Future objects.

Expected: 2 administrative threads per executor, independent of max_workers.

Changes

  • standalone/interactive/communication.py: AsyncSocketInterface, plus an optional context argument for interface_bootup.
  • task_scheduler/interactive/shared.py: execute_task_dict_async, with the same caching and error semantics as execute_task_dict.
  • task_scheduler/interactive/blockallocation_async.py (new): AsyncWorkerPool, a thread-like AsyncWorker handle, and the worker coroutine.
  • task_scheduler/interactive/blockallocation.py: switches between threads and async workers based on the opt-in flag, and closes the pool on shutdown.
  • tests/unit/task_scheduler/interactive/test_blockallocation_async.py (new): basic use, concurrent tasks, exceptions, shutdown modes, repeated creation, resizing, use inside a running event loop, and thread-count scaling.

Status: not yet verified

  • A local smoke test (1/4/16 workers) worked: results, exceptions and shutdown behaved correctly, the thread count stayed constant (+3 including the dependency-scheduler thread) and returned to baseline after shutdown.
  • The new test module did not finish locally within 5 minutes, so something in it may hang. It needs investigating before this leaves draft.
  • No full-suite run in async mode yet, i.e. with EXECUTORLIB_ASYNCIO=1 set globally.
  • Still to do: a benchmark script (thread count, startup and shutdown time, throughput) and an assessment of whether the approach should also cover OneProcessTaskScheduler.

Known limitations

  • Spawner bootup() and shutdown(wait=True) are still synchronous and run on the event loop. That's fine for MpiExecSpawner, but the pysqa spawner's polling sleep would block all workers while one boots.
  • Future done-callbacks now run on the event-loop thread rather than on per-worker threads.
  • queue_join_on_shutdown=True isn't supported in async mode. The scheduler always sets it to False.

🤖 Generated with Claude Code

Experimental, opt-in via EXECUTORLIB_ASYNCIO=1. Replaces the one thread
per worker of the BlockAllocationTaskScheduler with worker coroutines on
a single private asyncio event loop (one background thread owned by the
executor) plus one dispatcher thread reading the existing queue.Queue.
ZMQ communication uses zmq.asyncio with one shared context per executor.

The public API, concurrent.futures.Future objects, the worker-side
protocol and the threaded default path are unchanged. The calling
thread's event loop (e.g. Jupyter) is never used or modified.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@coderabbitai

coderabbitai Bot commented Oct 1, 2026

Copy link
Copy Markdown
Contributor

Important

Draft PR not reviewed

Draft PRs are not automatically reviewed by default.

  • Trigger a manual review

To automatically review draft PRs, update your CodeRabbit configuration:

reviews:
  auto_review:
    drafts: true
  • Autopilot · Keep fixing CodeRabbit findings and required CI, and resolving merge conflicts

Autopilot is currently an internal CodeRabbit preview.


Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@jan-janssen jan-janssen linked an issue Oct 1, 2026 that may be closed by this pull request
Co-authored-by: jan-janssen <3854739+jan-janssen@users.noreply.github.com>
@codecov

codecov Bot commented Oct 1, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 77.44681% with 53 lines in your changes missing coverage. Please review.
✅ Project coverage is 92.65%. Comparing base (cc75cf9) to head (59420a5).

Files with missing lines Patch % Lines
...ask_scheduler/interactive/blockallocation_async.py 77.53% 31 Missing ⚠️
...c/executorlib/task_scheduler/interactive/shared.py 45.71% 19 Missing ⚠️
...xecutorlib/standalone/interactive/communication.py 93.87% 3 Missing ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main    #1073      +/-   ##
==========================================
- Coverage   94.29%   92.65%   -1.65%     
==========================================
  Files          39       40       +1     
  Lines        2192     2423     +231     
==========================================
+ Hits         2067     2245     +178     
- Misses        125      178      +53     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Fearture] Use asyncio internally to reduce number of threads

2 participants