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.TaskGroupandasyncio.timeout). - Async clients for both ends (
httpx.AsyncClient,asyncpgorpsycopgasync). - 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.
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).
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.
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
- Producer-Consumer Pattern Design — the pattern and its variants.
- Backpressure Strategies for Fast Producers — blocking, shedding, and buffering.
- Tuning Prefetch and Consumer Concurrency — the broker-side equivalent of the buffer size.
- Competing Consumers vs Pub/Sub Fan-Out — distributing work versus copying it.