Skip to content

Async Data Pipelines in Python

A data pipeline moves records through stages — read, fetch, transform, enrich, write — and asyncio is a natural fit when most of those stages wait on I/O: an HTTP API, a database, object storage, a message broker. The concurrency that makes asyncio fast also makes a naive pipeline dangerous: a producer that can generate items faster than the slowest stage consumes them will simply fill memory. Measured on a two-stage pipeline with a slow consumer, an unbounded queue grew to 998 items and 9.8 MB of traced memory; the same pipeline with maxsize=50 peaked at 50 items and 0.5 MB — and finished in exactly the same 2.10 s, because the bottleneck was the consumer either way. Bounding queues costs no throughput and buys a hard memory ceiling. Batching costs almost nothing either and buys a great deal: inserting 50,000 rows into PostgreSQL went from 8,648 rows/s one row at a time to 400,527 with executemany chunks of 10,000, and 2.8 million with COPY chunks across four connections.

This section covers the patterns that turn a sequence of coroutines into a production pipeline: bounded stages, chunked work, fan-in of several sources, checkpoints that make long jobs resumable, and ordering guarantees across concurrent stages. The parent section, Concurrent Execution & Worker Patterns, covers the worker pools and queues each stage is built from.

Scope of this section:

  • Staged pipelines with bounded queues and per-stage concurrency.
  • Chunking inputs and outputs for batch-friendly sinks.
  • Merging several async sources into one stream, with clean shutdown.
  • Checkpointing progress so a crash costs minutes, not hours.
  • Preserving input order when stages run concurrently.

Architectural principles

  • Every queue is bounded. A pipeline's memory use must be a function of its configuration, not of its input size or the speed difference between stages.
  • Concurrency is set per stage. Each stage gets as many workers as its own latency and its downstream's capacity justify; one global concurrency number is wrong for every stage.
  • Batch at the boundaries. Sinks and sources with per-call overhead — databases, object stores, APIs with bulk endpoints — get chunks, not single items.
  • Progress is durable. A long job records where it is, so a restart resumes rather than repeats; the work after the last checkpoint must be safe to redo.
  • Failure has one owner. A stage failure cancels the pipeline through a TaskGroup and surfaces one error, instead of leaving stages blocked on queues forever.
A staged pipeline with bounded queues A flow of 6 stages. A staged pipeline with bounded queues source reads records bounded queue maxsize 100 fetch stage 16 workers, I/O bounded queue maxsize 100 transform 2 workers, CPU chunked sink batches of 1,000 Each queue caps memory between two stages; each stage's worker count matches its own latency.

Execution model: backpressure through awaits

In an asyncio pipeline, backpressure is not a protocol or a library feature; it is the side effect of await queue.put() on a full bounded queue. When the transform stage falls behind, its input queue fills, the fetch workers block on put, the fetch stage's input queue fills, and the source blocks on its put — the entire pipeline slows to the pace of its slowest stage, automatically, holding at most maxsize items per queue.

An unbounded queue breaks that chain. put never blocks, so the producer runs at its own speed and the difference accumulates in memory. The measurement made this concrete: with a fast producer and a 2 ms consumer, the unbounded queue held almost all 1,000 items at its peak (9.8 MB of 10 KB records), the bounded one never more than 50 (0.5 MB). Total time was 2.10 s in both cases, because the consumer set the pace regardless. In the multi-stage test — 16 fetch workers on 20 ms I/O feeding 4 writers — throughput was 766 items per second unbounded and 767 bounded, and the first queue's peak dropped from 1,984 items to 100.

One subtlety of the event loop matters here: a producer that puts into a queue in a tight loop without other awaits never yields until the queue is full, so bounded queues also improve fairness — they force the producer to let other stages run. The queue mechanics are covered in bounded asyncio queue with backpressure under load.

Same pipeline, slow consumer, bounded versus unbounded 2 horizontal bars comparing unbounded queue with the others. Same pipeline, slow consumer, bounded versus unbounded unbounded queue 9,816 KiB, 998 queued maxsize=50 503 KiB, 50 queued 1,000 records of 10 KB; consumer takes 2 ms each; both runs took 2.10 s. The bound changed memory twentyfold and throughput not at all.

Pattern catalogue

Staged pipeline with bounded queues

When to use: any multi-step processing where steps have different latencies.

async def pipeline(source, fetch, write, *, fetchers=16, writers=4, qsize=100):
    q_in: asyncio.Queue = asyncio.Queue(qsize)
    q_out: asyncio.Queue = asyncio.Queue(qsize)

    async def produce():
        async for item in source:
            await q_in.put(item)

    async def fetch_worker():
        while True:
            item = await q_in.get()
            try:
                await q_out.put(await fetch(item))
            finally:
                q_in.task_done()

    async def write_worker():
        while True:
            item = await q_out.get()
            try:
                await write(item)
            finally:
                q_out.task_done()

    async with asyncio.TaskGroup() as tg:
        workers = [tg.create_task(fetch_worker()) for _ in range(fetchers)]
        workers += [tg.create_task(write_worker()) for _ in range(writers)]
        await produce()
        await q_in.join()
        await q_out.join()
        for w in workers:
            w.cancel()

Trade-off: simple and fast; failure handling and shutdown need the care described in building a staged async pipeline with bounded queues.

Chunked sink

When to use: the sink has per-call overhead, as almost every database and storage API does.

async def chunked_writer(q: asyncio.Queue, conn, size: int = 1000) -> None:
    batch = []
    while True:
        batch.append(await q.get())
        if len(batch) >= size or q.empty():
            await conn.copy_records_to_table("items", records=batch)
            for _ in batch:
                q.task_done()
            batch = []

Trade-off: large chunks maximise throughput and delay visibility; measured, chunk size moved insert throughput from 8,648 to 400,527 rows per second. Details in chunking large inputs for async batch processing.

Fan-in of several sources

When to use: one consumer for several feeds — several API cursors, several queues, several files.

async def merge(*sources, maxsize: int = 100):
    q: asyncio.Queue = asyncio.Queue(maxsize)
    DONE = object()

    async def pump(src):
        async with aclosing(src) as it:
            async for item in it:
                await q.put(item)

    async def run_all():
        try:
            async with asyncio.TaskGroup() as tg:
                for s in sources:
                    tg.create_task(pump(s))
        finally:
            await q.put(DONE)

    runner = asyncio.create_task(run_all())
    try:
        while (item := await q.get()) is not DONE:
            yield item
        await runner
    finally:
        runner.cancel()
        await asyncio.gather(runner, return_exceptions=True)

Trade-off: arrival order rather than any source's order; verified to close every source on early exit and to surface a failing source's error. See merging multiple async iterators into one stream.

Checkpointed long job

When to use: jobs that run for minutes or hours over an ordered input.

async def run_resumable(items_from, process, store, every: int = 500):
    start = await store.load()
    since = 0
    async for pos, item in items_from(start):
        await process(item)
        since += 1
        if since >= every:
            await store.save(pos + 1)
            since = 0

Trade-off: work after the last checkpoint is redone after a crash — 200 items in the test with a checkpoint every 500 — so processing must be idempotent. See checkpointing progress in long-running async jobs.

Order-preserving concurrent stage

When to use: a stage must run concurrently but its output must keep input order.

Tag each item with a sequence number and re-sequence after the concurrent stage with a bounded reorder buffer. The trade-off — buffer memory against head-of-line stalls — is covered in preserving order across concurrent pipeline stages.

Inserting rows into PostgreSQL by chunk size and method 5 horizontal bars comparing COPY 5,000 x 4 connections with the others. Inserting rows into PostgreSQL by chunk size and method COPY 5,000 x 4 connections 2.83M rows/s COPY chunks of 10,000 1.54M rows/s executemany chunks of 10,000 401k rows/s executemany chunks of 100 197k rows/s one row per call 8.6k rows/s asyncpg 0.31, PostgreSQL 17 in a local container, unlogged table, 3 columns. Batching is the largest single lever in a pipeline's sink, by two orders of magnitude.

Queues or generators between stages?

Not every stage boundary needs a queue. Async generators chained together — write(transform(fetch(read()))) — are the simplest pipeline of all: each stage pulls from the one before it, memory is bounded by construction because nothing runs ahead, and there are no worker tasks to manage. The cost is that a generator chain is sequential: while the transform stage handles item 7, the fetch stage is not fetching item 8. For a chain whose slow stage is I/O, that wastes almost all of the concurrency asyncio offers.

The working rule is to use generators where a stage is cheap or must be sequential, and a queue with workers where a stage is slow and parallelisable:

async def records(api):                         # generator: paging is inherently sequential
    cursor = None
    while True:
        page = await api.page(cursor)
        for rec in page.items:
            yield rec
        if not (cursor := page.next):
            return


async def main(api, pool):
    q: asyncio.Queue = asyncio.Queue(200)       # queue: fan out the slow per-record calls
    async with asyncio.TaskGroup() as tg:
        workers = [tg.create_task(enrich_worker(q, pool)) for _ in range(16)]
        async for rec in records(api):
            await q.put(rec)
        await q.join()
        for w in workers:
            w.cancel()

Paging through an API cursor cannot be parallelised — each page depends on the previous one — so a generator is the natural shape; enriching each record is independent and slow, so it gets sixteen workers behind a bounded queue. The longer comparison, with measurements, is in async generators vs queues for streaming pipelines. Mixing the two is normal; most production pipelines are generator chains with one or two queue-and-worker stages where the latency is.

Generator or queue at this stage boundary? A decision on What is the next stage like with 3 outcomes. Generator or queue at this stage boundary? What is the next stage like? cheap or sequential async generator nothing runs ahead slow, items independent bounded queue + workers concurrency here a batch-friendly sink chunked writer amortise per-call cost Spend concurrency only where the latency is; everywhere else, keep the chain simple.

Resource boundaries

Resource What consumes it Bound it with
Memory between stages queued items × item size Queue(maxsize) per stage
Downstream concurrency workers per stage a worker count from the downstream's capacity
Database connections writer workers writers ≤ pool size
Batch memory chunk size × item size chunk by count and bytes
Redo after crash items since last checkpoint checkpoint interval
Reorder buffer throughput × straggler time a window, accepting stalls

The worker-count row is where pipelines most often go wrong. Sixteen fetch workers against an API that tolerates eight concurrent calls do not double throughput; they produce rate-limit errors and retries. Size each stage from its downstream's limits — the same reasoning as in optimizing worker pool sizes for mixed I/O and CPU workloads — and let bounded queues absorb the speed differences.

Integrated production example

A three-stage pipeline — read pages from an API, enrich each record, bulk-load into Postgres — with bounded queues, per-stage concurrency, chunked writes, a checkpoint per committed batch, and one failure path:

import asyncio
import logging

import asyncpg
import httpx

log = logging.getLogger("pipeline")


async def run(api: httpx.AsyncClient, pool: asyncpg.Pool, checkpoint, *, enrichers: int = 8,
              batch: int = 1000, qsize: int = 200) -> int:
    raw: asyncio.Queue = asyncio.Queue(qsize)
    ready: asyncio.Queue = asyncio.Queue(qsize)
    loaded = 0

    async def read_pages():
        cursor = await checkpoint.load()                       # resume where we stopped
        while cursor is not None:
            r = await api.get("/records", params={"cursor": cursor, "limit": 500})
            r.raise_for_status()
            body = r.json()
            for rec in body["items"]:
                await raw.put((body["cursor"], rec))           # blocks when enrichers lag
            cursor = body.get("next")

    async def enrich():
        while True:
            page_cursor, rec = await raw.get()
            try:
                async with asyncio.timeout(10):
                    extra = await api.get(f"/details/{rec['id']}")
                rec["details"] = extra.json()
                await ready.put((page_cursor, rec))
            finally:
                raw.task_done()

    async def load():
        nonlocal loaded
        buf: list = []
        while True:
            buf.append(await ready.get())
            if len(buf) >= batch or ready.empty():
                rows = [(r["id"], r["name"], r["details"]["score"]) for _, r in buf]
                async with pool.acquire() as conn, conn.transaction():
                    await conn.copy_records_to_table("records", records=rows)
                await checkpoint.save(buf[-1][0])              # after the commit, never before
                loaded += len(buf)
                for _ in buf:
                    ready.task_done()
                buf = []

    async with asyncio.TaskGroup() as tg:
        workers = [tg.create_task(enrich()) for _ in range(enrichers)]
        workers.append(tg.create_task(load()))
        await read_pages()
        await raw.join()
        await ready.join()
        for w in workers:
            w.cancel()
    log.info("pipeline loaded %d records", loaded)
    return loaded

Diagnostic Hook: export each queue's depth and each stage's items per second and busy workers. The stage whose input queue is full and whose output queue is empty is the bottleneck — add workers there if its downstream can take them, or make its per-item work cheaper. A checkpoint age that keeps growing while throughput is non-zero means the loader's commits are failing and retrying.

Diagnostic hook callout

Watch three things on every pipeline run:

  • Queue depth per stage. Full input and empty output identifies the bottleneck stage instantly.
  • Items per second per stage, in and out. A stage whose out-rate is below its in-rate is either failing items or accumulating them.
  • Checkpoint age. The time since the last durable checkpoint is the amount of work a crash would redo; alert if it exceeds what you are willing to repeat.

Thresholds: a queue at maxsize for more than a minute is a sustained bottleneck, not a burst; an error rate above a fraction of a percent per stage usually means a dependency problem that retries will only amplify.

Failure modes

Failure mode Root cause Detection Fix
Memory grows until OOM unbounded queue before a slow stage queue depth rises with input size maxsize on every queue
Pipeline hangs at the end a stage crashed, others block on queues stalled depths, no errors logged TaskGroup so a failure cancels all
Downstream rate-limited too many workers in a stage 429s, retries workers from downstream capacity
Slow sink one row per write low rows/s, high call count chunked writes, COPY
Full restart after crash no checkpoint hours of rework checkpoint after each committed batch
Duplicates after resume non-idempotent writes after checkpoint duplicate keys idempotent upserts
Output out of order concurrent stage without re-sequencing consumer sees gaps sequence numbers and a reorder buffer

Frequently Asked Questions

How do I build a data pipeline with asyncio?

Connect stages with bounded asyncio.Queue objects, give each stage its own pool of worker tasks sized for its latency and its downstream's limits, run everything in a TaskGroup so a failure stops the whole pipeline, and batch writes to the sink.

Why should asyncio pipeline queues be bounded?

An unbounded queue lets a fast producer accumulate the difference between stage speeds in memory. In testing, bounding a queue at 50 cut peak memory from 9.8 MB to 0.5 MB with no change in total run time.

How much does batching help an async pipeline?

At the sink, enormously: inserting rows into PostgreSQL went from 8,648 rows per second one at a time to about 400,000 with executemany chunks of 10,000 and 2.8 million with COPY across four connections.

How do I make a long-running async job resumable?

Save a checkpoint — the position after the last item whose effects are durably committed — every few hundred items, and resume from it on start. Make processing idempotent, because items after the last checkpoint are redone.

How do I merge several async iterators?

Run one pump task per source that puts items into a shared bounded queue, yield from the queue until all sources finish, and cancel the pumps in a finally block so breaking out early closes every source.