Skip to content

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

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.

20,000 messages from a fast source, 50 ms of I/O each A grid of 4 rows by 5 columns. 20,000 messages from a fast source, 50 ms of I/O each pattern time peak live tasks concurrent calls peak traced memory create_task per message, no limit 0.64 s 3,603 3,603 6.1 MiB Semaphore(100) inside each task 10.43 s 19,502 100 35.3 MiB Semaphore(100) before create_task 10.38 s 102 100 0.2 MiB 100 workers + Queue(100) 10.16 s 102 100 0.2 MiB Python 3.14; messages carry a 200-byte payload.

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.

Backpressure from the slot to the source A sequence of 6 messages between 4 participants. Backpressure from the slot to the source source producer slots (100) task message acquire(): slot free create_task(run(msg)) acquire(): all busy, wait finished: release() read next message Unread messages stay in the source, not in memory as tasks.

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.

How should this consumer bound its tasks? A decision on Where do the items come from with 4 outcomes. How should this consumer bound its tasks? Where do the items come from? a small, validated list task per item, no limit capped by validation a stream, task per message acquire slot before create_task 102 tasks a stream, long-lived workers N workers + Queue(N) 102 tasks semaphore inside each task bounds calls, not tasks 19,502 tasks Backpressure has to reach the source, or tasks absorb the excess.

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 limit tasks 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.