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¶
- Python 3.11+, stdlib only.
- Queue counters, from tracking unfinished work with task_done and join.
- The broader topic, from Async Data Pipelines.
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.
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.
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.
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()outsidefinally. One failing item hangsjoin()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.
Related¶
- Async Data Pipelines — up to the topic overview.
- Chunking large inputs for async batch processing — making the write stage fast.
- Concurrent Execution & Worker Patterns — the section overview.