Bounded Queues with Python asyncio

Not every queue needs a broker. Inside one process, asyncio.Queue gives you the producer-consumer pattern with backpressure in a few lines, and this guide builds a production-grade version as part of Producer-Consumer Pattern Design in Queue Fundamentals & Architecture. It is also the pattern that sits inside many workers: a consumer fetches from a broker into a local bounded queue, and a pool of tasks drains it.

Problem Statement

An ingestion service reads records from a paginated partner API and writes them to Postgres. The first version created one task per record with asyncio.gather; on a large account it launched 400,000 coroutines at once, held every record in memory, opened far more database connections than the pool allowed, and crashed with memory errors. A second version processed records sequentially and took eleven hours. You want a pipeline where fetching and writing overlap, concurrency is capped at what the database can take, memory stays bounded no matter how large the account is, errors in one record do not stop the pipeline, and shutdown is clean.

Prerequisites

  • Python 3.11+ (for asyncio.TaskGroup and asyncio.timeout).
  • Async clients for both ends (httpx.AsyncClient, asyncpg or psycopg async).
  • A measured or estimated safe concurrency for the slowest stage — here, the database pool size.
  • Familiarity with the backpressure ideas in backpressure strategies for fast producers.

Step 1 — Size the Queue and the Worker Pool

Two numbers define the pipeline: how many consumers run concurrently (bounded by the slowest downstream) and how many items may wait between stages (bounded by memory and by how far ahead the producer should get).

DB_POOL_SIZE = 20
CONSUMERS = DB_POOL_SIZE          # one write in flight per connection
QUEUE_MAX = CONSUMERS * 4         # enough buffer to hide fetch latency jitter
# Memory bound: QUEUE_MAX items + CONSUMERS in progress + one fetched page
#   = 80 + 20 + 100 records, whatever the account size

A small multiple of the consumer count is usually right for the buffer: large enough that consumers rarely starve while the producer waits on the next page, small enough that memory is trivially bounded. An unbounded queue (maxsize=0) would recreate the original memory problem, only more slowly.

One producer, a bounded buffer, a capped pool A producer task fetches pages of 100 records from the partner API and puts records into an asyncio queue with a maximum size of 80. Twenty consumer tasks take records from the queue and write them to Postgres through a pool of twenty connections. When the queue is full, the producer's put call waits, which slows fetching to the rate the database can absorb. Backpressure by construction producer pages of 100 asyncio.Queue(maxsize=80) put() waits when full 20 consumers DB pool of 20 Memory is bounded by queue size and pool size, not by account size.

Step 2 — Write the Producer with an Awaited put

await queue.put(item) suspends the producer when the queue is full. That single line is the backpressure mechanism: the producer fetches the next page only when consumers have made room.

import asyncio, httpx

async def produce(client: httpx.AsyncClient, account_id: str, q: asyncio.Queue) -> int:
    cursor, count = None, 0
    while True:
        resp = await client.get(f"/accounts/{account_id}/records",
                                params={"cursor": cursor, "limit": 100})
        resp.raise_for_status()
        page = resp.json()
        for record in page["records"]:
            await q.put(record)                  # waits here when consumers fall behind
            count += 1
        cursor = page.get("next_cursor")
        if not cursor:
            return count

Avoid put_nowait in a loop with a retry on QueueFull — it turns backpressure into busy-waiting. If the producer must never block (for example, it serves HTTP requests), use put_nowait and shed load explicitly on QueueFull, as discussed in the backpressure guide linked above.

Step 3 — Write Consumers That Survive Bad Items

Each consumer loops forever, taking an item, processing it, and marking it done. Exceptions for one item are caught and recorded so the consumer keeps going; task_done() runs in a finally so join() is never left waiting on an item that failed.

import logging
log = logging.getLogger("ingest")

async def consume(name: str, q: asyncio.Queue, pool, stats: dict) -> None:
    while True:
        record = await q.get()
        try:
            async with asyncio.timeout(10):                       # bound each item
                async with pool.acquire() as conn:
                    await conn.execute(
                        "INSERT INTO records (id, payload) VALUES ($1, $2) "
                        "ON CONFLICT (id) DO UPDATE SET payload = EXCLUDED.payload",
                        record["id"], json.dumps(record))
            stats["ok"] += 1
        except (asyncio.TimeoutError, asyncpg.PostgresError) as exc:
            stats["failed"] += 1
            await record_failure(record, exc)                     # park for later retry
            log.warning("record failed", extra={"id": record.get("id"), "error": repr(exc)})
        finally:
            q.task_done()

The upsert makes the write idempotent, so re-running an ingestion after a crash does not duplicate records. asyncio.CancelledError is deliberately not caught: cancellation is how shutdown stops consumers (Step 4).

The consumer loop Each consumer waits on queue.get, processes the record with a ten-second timeout, and on success increments the ok counter. On a timeout or database error it records the failure for later and continues. In every case task_done is called in a finally block so that queue.join can complete. Cancellation is not caught, so shutdown can stop the loop. Every item ends in task_done() await q.get() process timeout 10 s ok += 1 park failure, log task_done() CancelledError is not caught: cancelling the task is the shutdown signal.

Step 4 — Coordinate Completion with join and Cancel

The pipeline is done when the producer has finished and every queued item has been processed. q.join() waits for the second condition; then the idle consumers are cancelled. No sentinel values are needed.

async def ingest(account_id: str) -> dict:
    q: asyncio.Queue = asyncio.Queue(maxsize=QUEUE_MAX)
    stats = {"ok": 0, "failed": 0}
    async with httpx.AsyncClient(base_url=PARTNER_URL, timeout=30) as client, \
               asyncpg.create_pool(DSN, min_size=5, max_size=DB_POOL_SIZE) as pool:
        consumers = [asyncio.create_task(consume(f"c{i}", q, pool, stats)) for i in range(CONSUMERS)]
        try:
            produced = await produce(client, account_id, q)   # returns when all pages fetched
            await q.join()                                    # all items processed
        finally:
            for c in consumers:
                c.cancel()                                    # consumers are blocked in q.get()
            await asyncio.gather(*consumers, return_exceptions=True)
    return {"produced": produced, **stats}

The finally also runs when the producer raises (the partner API fails mid-way) or when the whole ingest coroutine is cancelled, so consumers never leak. If a consumer crashed unexpectedly, join() would still complete because remaining consumers drain the queue — but with fewer workers; monitor for that in Step 6.

Step 5 — Add Graceful Shutdown for Long-Running Services

When the pipeline runs inside a long-lived service, a SIGTERM should stop fetching new pages, let buffered items finish within a deadline, and then cancel.

import signal

async def main():
    stop = asyncio.Event()
    loop = asyncio.get_running_loop()
    loop.add_signal_handler(signal.SIGTERM, stop.set)
    ingest_task = asyncio.create_task(ingest_forever(stop))
    await stop.wait()
    try:
        async with asyncio.timeout(25):          # < container grace period
            await ingest_task                    # producer sees stop, returns; join drains
    except TimeoutError:
        ingest_task.cancel()                     # hard stop: in-flight items are retried next run

async def produce_until(stop: asyncio.Event, client, q):
    async for page in paginate(client):
        if stop.is_set():
            return
        for record in page:
            await q.put(record)

Because writes are idempotent upserts and progress is resumable from a cursor, a hard stop costs at most the in-flight items, which the next run repeats harmlessly. The same drain-then-cancel sequence applies to broker-backed workers, as in graceful shutdown for Go workers.

Stop intake, drain, then cancel On SIGTERM the stop event is set and the producer returns before fetching another page. Consumers keep draining the up to 80 buffered records. If the queue drains within 25 seconds the pipeline exits cleanly; otherwise remaining consumers are cancelled and their in-flight records are picked up by the next run. Shutdown within the grace period SIGTERM producer stops drain buffer (max 25 s) cancel rest A bounded buffer is what makes the drain step fast enough to fit the grace period.

Step 6 — Measure Queue Depth and Consumer Utilisation

Two metrics tell you which side is the bottleneck. A queue that sits full means consumers are the limit (and the producer is being held back); a queue that sits empty means the producer is the limit and consumers are idle.

async def report(q: asyncio.Queue, stats: dict, consumers: list[asyncio.Task]):
    while True:
        await asyncio.sleep(10)
        alive = sum(1 for c in consumers if not c.done())
        QUEUE_DEPTH.set(q.qsize())
        CONSUMERS_ALIVE.set(alive)
        log.info("pipeline", extra={"depth": q.qsize(), "max": q.maxsize, "alive": alive, **stats})

In the scenario, depth stayed at 80 of 80 with 20 consumers busy — the database was the limit, as intended — and total time fell from eleven hours to 38 minutes with memory flat at about 150 MB. If depth is always near zero, raise producer concurrency (fetch pages in parallel) before adding consumers.

Verification

async def test_memory_is_bounded_for_huge_accounts(fake_partner_with_records):
    fake_partner_with_records(400_000)
    tracemalloc.start()
    result = await ingest("acct-huge")
    _, peak = tracemalloc.get_traced_memory()
    assert result["ok"] == 400_000
    assert peak < 200 * 1024 * 1024                  # bounded, independent of record count

async def test_failures_do_not_stop_pipeline(fake_partner_with_records, flaky_db):
    fake_partner_with_records(1_000); flaky_db.fail_ids({5, 500})
    result = await ingest("acct-1")
    assert result == {"produced": 1000, "ok": 998, "failed": 2}

Gotchas & Edge Cases

asyncio.gather over everything. Launching one task per item has no bound; it is the problem this pattern replaces. Use a pool of consumers or an asyncio.Semaphore if a queue is overkill.

Forgetting task_done on errors. join() then waits forever. Keep it in finally.

CPU-bound processing. asyncio consumers share one thread; a CPU-heavy transform blocks every consumer. Offload with asyncio.to_thread or a process pool, or move the work to a real job queue.

Priority. asyncio.PriorityQueue gives in-process priority; items must be comparable, so wrap them as (priority, sequence, item) tuples.

FAQ

When should I use a broker instead? When work must survive a process crash, be shared across machines, or be retried later. An in-process queue loses its contents on restart; it is a concurrency tool, not a durability tool.

How do I parallelise the producer too? Run several producer tasks, each fetching a different slice (date ranges, shard ids, or cursor ranges if the API allows), all putting into the same bounded queue. Track producer completion with a counter or a second TaskGroup, and call q.join() only after every producer has returned. The queue bound still caps memory, because every producer waits on the same put when consumers fall behind.

Should failed records be retried in the same run? Usually not inline: a record that timed out may time out again and occupy a consumer slot for another ten seconds. Park failures in a table or a file and retry them in a second pass after the main run, with its own lower concurrency — the same separation between hot path and retry path that job frameworks use.

How is this different from a thread pool? The pattern is the same (queue.Queue plus threads); asyncio suits I/O-bound work with many concurrent operations and low memory per task. Threads suit blocking libraries.

What about anyio or trio? Both offer memory-object streams with bounded buffers and structured cancellation that implement the same pattern with fewer foot-guns around task lifetimes.

Related