Bounding the Number of Live Tasks¶
A service that starts a task per incoming message — per queue item, per webhook, per stream event — has as many live tasks as messages arrive faster than they finish. Bounding concurrency is not the same as bounding tasks: the common pattern of a semaphore inside each task limits how many do work at once, but every message still gets a task that sits waiting for the semaphore. Measured on Python 3.14 with 20,000 messages from a fast source, each needing 50 ms of I/O: starting a task per message with no limit finished in 0.64 s with up to 3,603 tasks — and so 3,603 concurrent calls to whatever was downstream. A Semaphore(100) inside each task limited the calls to 100 but kept up to 19,502 tasks alive and 35.3 MiB of traced memory. Acquiring the same semaphore before create_task, so the producer waits for a slot, kept at most 102 tasks and 0.2 MiB, in 10.38 s — the time 100 concurrent calls need for 20,000 messages. A pool of 100 workers reading a bounded queue matched it: 102 tasks, 0.2 MiB, 10.16 s. This guide shows how to bound tasks, not just work.
Prerequisites¶
- Python 3.11+.
- Keeping tasks referenced, from preventing task garbage collection with strong references.
- The topic overview, Task Scheduling & Lifecycle.
1. Measure live tasks, not just concurrency¶
The number of live tasks is visible directly. Sample it alongside whatever concurrency limit you have, because the two can diverge by orders of magnitude:
async def watch_tasks(gauge, interval=1.0):
while True:
gauge.set(len(asyncio.all_tasks()))
await asyncio.sleep(interval)
Measured with 20,000 messages: unbounded spawning peaked at 3,603 live tasks, all doing I/O at once; the semaphore-inside version peaked at 19,502 live tasks, of which 100 were working and the rest waiting. Each waiting task holds its coroutine frame, its message, and a future on the semaphore's wait queue — 35.3 MiB of traced memory here for small messages, far more for real payloads. When the source is faster than the work for long enough, the count grows without bound until memory runs out, as described in tracking task growth in long-running services.
Verify: a live-task gauge exists, and under sustained load it stays near your concurrency limit.
2. Acquire the slot before creating the task¶
Move the wait from inside the task to the producer. The producer acquires a slot, then creates the task; the task releases the slot when it finishes:
async def consume(source, handle, limit: int = 100):
slots = asyncio.Semaphore(limit)
tasks: set[asyncio.Task] = set()
async def run(msg):
try:
await handle(msg)
finally:
slots.release()
async for msg in source:
await slots.acquire() # wait here, before a task exists
task = asyncio.create_task(run(msg))
tasks.add(task)
task.add_done_callback(tasks.discard)
await asyncio.gather(*tasks)
Measured: at most 102 live tasks — 100 workers plus the producer and the sampler — and 0.2 MiB, against 19,502 tasks and 35.3 MiB with the semaphore inside. The producer stops reading the source while all slots are busy, so the bound propagates upstream: a message queue keeps its unread messages, a stream applies TCP flow control, and nothing piles up in the process. The finally guarantees the slot is returned whether handle succeeds, fails or is cancelled.
Verify: while all slots are busy, the producer is waiting on acquire() and the live-task count stays at the limit.
3. Or use a fixed pool of workers¶
The other way to bound tasks is to create them once: a fixed number of workers, each reading items from a bounded queue:
async def consume(source, handle, workers: int = 100):
queue: asyncio.Queue = asyncio.Queue(maxsize=workers)
async def worker():
while (msg := await queue.get()) is not None:
await handle(msg)
async with asyncio.TaskGroup() as tg:
for _ in range(workers):
tg.create_task(worker())
async for msg in source:
await queue.put(msg) # waits when workers are behind
for _ in range(workers):
await queue.put(None)
Measured: 10.16 s, 102 live tasks, 0.2 MiB — the same bound as step 2. The worker pool reuses tasks instead of creating 20,000 of them, which saves the per-task cost — a few microseconds each — and gives each worker a place to keep per-worker state, such as a connection. Step 2's pattern gives each message its own task, which keeps per-message cancellation, naming and error isolation simple. Per-item error handling inside workers is essential, as described in handling per-item errors in async pipelines.
Verify: the pool's queue never exceeds its maxsize, and a failure in one item does not stop a worker.
4. Decide whether unbounded is ever acceptable¶
Unbounded spawning was the fastest run here — 0.64 s against about 10 s — because nothing limited the 50 ms calls, and every one of the 3,603 concurrent calls was served instantly by a simulated downstream. Real downstreams are not like that: a database pool, an API's rate limit or a CPU-bound handler would have queued or rejected most of those calls, and the queueing would have happened in the tasks, in memory. Unbounded spawning is acceptable only when both of these hold:
# acceptable: the source itself is bounded and small
async with asyncio.TaskGroup() as tg:
for user_id in request.user_ids[:50]: # capped by validation
tg.create_task(load_user(user_id))
the number of items is bounded by something you control — a validated request, a fixed list — and the downstream can take that many at once. Anything driven by an external stream, queue or client fails the first test, as covered in bounding user-controlled fan-out.
Verify: every unbounded create_task loop iterates over something with a known, small maximum size.
5. Shut down a bounded consumer cleanly¶
A bounded consumer holds at most limit in-flight items when asked to stop. Stop reading the source first, then let in-flight tasks finish within a deadline, and cancel what remains:
async def drain(tasks: set[asyncio.Task], deadline: float) -> None:
if not tasks:
return
done, pending = await asyncio.wait(tasks, timeout=deadline)
for task in pending:
task.cancel()
await asyncio.gather(*pending, return_exceptions=True)
log.info("drained %d tasks, cancelled %d", len(done), len(pending))
With the slot pattern, the number of tasks to drain is never more than the limit, so the shutdown deadline can be sized from the limit and the per-item time — at most 100 items of 50 ms here. With unbounded spawning, the drain could face thousands of tasks, most of which never started their work. Messages that were taken from a durable source but not completed must be returned to it — the SQS visibility reset, a Kafka offset not committed — as described in Graceful Shutdown & Signals.
Verify: a shutdown during load finishes within its deadline and leaves no unacknowledged work lost.
Verification¶
Live tasks are bounded when:
- A live-task gauge stays near the concurrency limit under sustained load.
- Slots are acquired before tasks are created, or a fixed worker pool reads a bounded queue.
- Unbounded spawning is limited to small, validated inputs.
- Shutdown drains at most
limittasks within a deadline and returns unfinished work to its source.
Diagnostic Hook: compare the live-task count with the number of items actually being worked on. A large and growing difference means tasks are waiting for a limit inside themselves — the semaphore-inside pattern — and memory will follow the backlog until the source slows down or the process is killed.
Pitfalls & edge cases¶
- Semaphore inside the task. Measured: 19,502 live tasks for a limit of 100.
- Unbounded spawning from a stream. Measured: 3,603 concurrent downstream calls.
- Forgetting to release on failure. Release in
finally, or slots leak away. - Draining an unbounded backlog at shutdown. Thousands of tasks, most never started.
Frequently Asked Questions¶
How do I limit the number of asyncio tasks?
Acquire a semaphore before calling create_task and release it in the task's finally block, or use a fixed pool of workers reading a bounded queue. Both kept 20,000 messages to about 100 live tasks in testing.
Doesn't a semaphore inside each task limit concurrency?
It limits how many run the protected code at once, not how many tasks exist: 19,502 tasks were alive at once, waiting, using 35.3 MiB.
Is it faster not to limit tasks?
Only if the downstream can absorb everything: unbounded spawning finished in 0.64 s by making 3,603 concurrent calls, which a real database or API would queue or reject.
How do I see how many tasks are running?
len(asyncio.all_tasks()) gives the live count; export it as a gauge and compare it with your concurrency limit.
Related¶
- Task Scheduling & Lifecycle — up to the topic overview.
- Scheduling coroutines at wall-clock times — the other kind of task lifecycle control.
- Asyncio Fundamentals & Event Loop Architecture — the section overview.