Skip to content

Building a Staged Async Pipeline with Bounded Queues

A staged pipeline splits work into steps connected by queues, each step with its own pool of workers. It is the shape almost every async data job ends up in: read items, fetch something for each over the network, write results. The two decisions that matter are how many workers each stage gets and how big each queue may grow. Measured on a pipeline of 2,000 items — 16 fetch workers doing 20 ms of I/O, 4 writers doing 2 ms each — unbounded queues gave 766 items per second with the first queue peaking at 1,984 items; queues bounded at 100 gave 767 items per second with that queue peaking at 100, and traced memory falling from 106 KiB to 25 KiB. The bound is free; leaving it out is a memory leak waiting for a large input. This guide builds the pipeline end to end, including the shutdown and failure handling that most examples leave out.

Prerequisites

1. Define stages as worker functions over queues

Each stage reads from an input queue, does its work, and writes to an output queue. Keep the stage functions ignorant of the pipeline around them:

import asyncio
from collections.abc import Awaitable, Callable


def stage(fn: Callable[[object], Awaitable[object]],
          inbox: asyncio.Queue, outbox: asyncio.Queue | None):
    async def worker() -> None:
        while True:
            item = await inbox.get()
            try:
                result = await fn(item)
                if outbox is not None and result is not None:
                    await outbox.put(result)           # blocks when the next stage is behind
            finally:
                inbox.task_done()
    return worker

task_done() runs in finally, so a failing item still counts as handled and join() cannot hang on it. The await outbox.put() is where backpressure happens: when the next stage's queue is full, this worker waits, which leaves its own input queue to fill, and so on back to the source. Returning None from a stage function drops the item, which is a convenient filter.

Verify: a stage whose function raises for one item still processes the rest, and its queue's join() returns.

Two stages, two bounded queues A flow of 5 stages. Two stages, two bounded queues source await put() queue 1 maxsize 100 16 fetch workers 20 ms I/O each queue 2 maxsize 100 4 writers 2 ms each Backpressure flows right to left through every await put().

2. Size each stage from its own latency

The worker count per stage should keep the stage's throughput at or above the pipeline's target, without exceeding what its downstream tolerates. Little's law gives the first half:

import math


def workers_needed(target_per_s: float, latency_s: float, headroom: float = 1.25) -> int:
    return max(1, math.ceil(target_per_s * latency_s * headroom))


workers_needed(750, 0.020)   # fetch stage: 19
workers_needed(750, 0.002)   # write stage: 2

With 16 fetch workers at 20 ms, the fetch stage tops out near 800 items per second, which is why the measured pipeline ran at 766–767: the fetch stage was the bottleneck, and the four writers mostly waited (their queue peaked at 16). Adding writers would not have helped; adding fetchers would, if the service being fetched from could absorb more concurrent calls. Check that second half against the downstream's limits — rate limits, connection pools — as described in limiting concurrent requests with asyncio.Semaphore.

Verify: under load, the bottleneck stage's input queue sits near full and its output queue near empty; other stages show the reverse.

3. Bound every queue

q_in: asyncio.Queue = asyncio.Queue(maxsize=100)
q_out: asyncio.Queue = asyncio.Queue(maxsize=100)

Measured with and without bounds: throughput 766 versus 767 items per second — identical within noise — while the first queue's peak went from 1,984 to 100 and traced memory from 106 KiB to 25 KiB for small items. The effect scales with item size: in a slow-consumer test with 10 KB records, the unbounded queue reached 9.8 MB and the bounded one 0.5 MB, again with identical run time. A bound only costs throughput if it is so small that a fast stage idles between bursts; a few times the downstream stage's worker count is usually plenty.

Verify: feed the pipeline ten times as many items; peak memory stays the same.

Peak depth of the first queue, 2,000 items 2 horizontal bars comparing unbounded with the others. Peak depth of the first queue, 2,000 items unbounded 1,984 items, 766/s maxsize=100 100 items, 767/s 16 fetch workers at 20 ms, 4 writers at 2 ms; traced memory 106 KiB vs 25 KiB. The bound changed memory, not speed: the slowest stage sets the pace either way.

4. Run everything in a TaskGroup and shut down in order

Workers loop forever, so the pipeline needs an explicit end: feed all input, wait for each queue to drain in order, then cancel the workers. A TaskGroup ties the workers' lifetime to the pipeline's and turns any worker crash into a pipeline failure:

async def run_pipeline(source, fetch, write, *, fetchers=16, writers=4, qsize=100) -> None:
    q_in: asyncio.Queue = asyncio.Queue(qsize)
    q_out: asyncio.Queue = asyncio.Queue(qsize)
    async with asyncio.TaskGroup() as tg:
        workers = [tg.create_task(stage(fetch, q_in, q_out)()) for _ in range(fetchers)]
        workers += [tg.create_task(stage(write, q_out, None)()) for _ in range(writers)]
        async for item in source:
            await q_in.put(item)
        await q_in.join()                  # everything fetched
        await q_out.join()                 # everything written
        for w in workers:
            w.cancel()                     # idle workers blocked in get(): stop them

Join in stage order: q_in.join() returning means every fetched result has been put into q_out, so q_out.join() afterwards sees all of them. On Python 3.13+, queue.shutdown() is an alternative to cancelling idle workers, as in shutting down queues with Queue.shutdown.

Verify: the function returns only after the last item is written, and no worker tasks remain afterwards.

5. Decide what an item failure does to the pipeline

The stage helper above catches nothing; an exception in fn propagates out of the worker, the TaskGroup cancels everything, and the pipeline fails fast. That is right when any failure means the output is unusable. When items are independent, handle failures per item and keep going:

def tolerant(fn, dead_letters: asyncio.Queue, max_failures: int):
    failures = 0

    async def wrapped(item):
        nonlocal failures
        try:
            return await fn(item)
        except Exception as exc:
            failures += 1
            await dead_letters.put((item, repr(exc)))
            if failures > max_failures:
                raise RuntimeError(f"too many failures ({failures})") from exc
            return None                                     # drop this item, continue
    return wrapped

A failure budget combines both behaviours: isolated bad items go to a dead-letter queue, but a systematic problem — the downstream is down, every item fails — still stops the pipeline instead of dead-lettering the whole input. The dead-letter pattern itself is in implementing a dead letter queue with asyncio.

Verify: one bad item is dead-lettered and the run completes; a dependency outage stops the run after max_failures.

What should one failed item do to the pipeline? A decision on Is the output usable without this item with 3 outcomes. What should one failed item do to the pipeline? Is the output usable without this item? no let it raise TaskGroup stops all yes, items independent dead-letter, continue per-item handling yes, but watch for outages failure budget stop after N Isolated bad records should not stop a run; a dead dependency should.

Verification

The staged pipeline is correct when:

  • Every queue is bounded, and peak memory does not grow with input size.
  • Each stage's worker count matches its latency and its downstream's limits.
  • Shutdown joins queues in stage order and leaves no workers running.
  • Item failures follow a deliberate policy: fail fast, or dead-letter with a budget.

Diagnostic Hook: export depth for each queue and items processed per second for each stage. The bottleneck is the stage with a full inbox and an empty outbox; if it is an I/O stage with headroom downstream, add workers; if its downstream is saturated, the pipeline is at capacity and adding workers will only add errors.

Pitfalls & edge cases

  • Unbounded queues. Memory grows with the input instead of with the configuration.
  • task_done() outside finally. One failing item hangs join() forever.
  • Joining the second queue first. Items still in flight from stage one are missed.
  • Too many workers for the downstream. Throughput stays flat while errors rise.

Frequently Asked Questions

How do I connect asyncio pipeline stages?

With bounded asyncio.Queue objects: each stage's workers get from an input queue and put into an output queue. A full output queue makes put wait, which applies backpressure back through every earlier stage.

Does bounding asyncio queues slow a pipeline down?

Not when the bound is a few times the next stage's worker count. In testing, a two-stage pipeline ran at 766 items per second unbounded and 767 bounded, while the first queue's peak fell from 1,984 to 100.

How many workers should each pipeline stage have?

Enough to meet the target throughput — target rate times per-item latency, plus headroom — and no more than the stage's downstream dependency can handle concurrently.

How do I shut down a pipeline of worker tasks cleanly?

After the source is exhausted, join each queue in stage order so all in-flight items finish, then cancel the idle workers, all inside a TaskGroup so a crashed worker fails the whole pipeline.