skip to content

asyncio.Queue with no maxsize grows without bound in a streaming job — how do you apply backpressure?

level: seniorimportance: should knowfreq 40%

answer

  1. The default constructor bounds nothing
  2. Memory grows at the rate difference
  3. Make put() a suspension point
  4. The slowest stage should set the pace
  5. asyncio.Queue(maxsize=n); QueueFull to shed

basics

~20 s

An asyncio.Queue defaults to maxsize=0, meaning unbounded, so a fast producer buffers the whole stream in memory. Give it a maxsize: once full, await queue.put(...) suspends the producer until a consumer drains an item, which is backpressure.

solid answer

~40 s

`asyncio.Queue()` is unbounded by default (`maxsize=0`). When the producer is faster than the consumers — reading rows is cheap, processing them is not — the difference accumulates in the queue and memory grows for the length of the run with no error and no slowdown until the process is killed. Constructing it as `asyncio.Queue(maxsize=n)` turns `await put()` into a suspension point: with the buffer full the producer's task parks on a waiter future and yields to the loop, and is resumed only when a `get()` frees a slot. That propagates the consumers' real rate back to the source. Choose `maxsize` as a smoothing buffer, not a reservoir — a small multiple of the worker count is usually right. `put_nowait()` raises `asyncio.QueueFull` instead of waiting, which is how you shed load rather than absorb it.

code

python · 17 lines
python
import asyncio

async def produce(q):
    for row in range(500):
        await q.put(row)      # suspends once the buffer is full
    await q.put(None)

async def consume(q):
    while (row := await q.get()) is not None:
        await asyncio.sleep(0.001)   # the slow stage

async def main():
    q = asyncio.Queue(maxsize=100)
    await asyncio.gather(produce(q), consume(q))
    print("peak buffering capped at", q.maxsize)

asyncio.run(main())

go deeper

for a junior

Remember that asyncio.Queue() is unbounded unless you pass maxsize, and that awaiting put on a full bounded queue makes the producer wait rather than failing. That single argument is the difference between bounded and unbounded memory.

for a middle

Explain the mechanics: a full queue parks the producer on a putter future and yields to the event loop, and each get wakes one waiter. Contrast await put() with put_nowait() raising asyncio.QueueFull.

for a senior

Diagnose it in production: correlate rising RSS with queue depth, decide between backpressure and load-shedding, size maxsize as a jitter buffer, and watch for the feedback-loop deadlock a bound introduces.

for a principal

Own the end-to-end flow-control story: where pressure should surface, what gets dropped when it cannot be absorbed, which SLOs the queueing latency affects, and where an in-process queue must give way to a durable broker.

## The default is the trap `asyncio.Queue()` with no argument is **unbounded**: `maxsize` is 0, and `put()` never waits. That is a reasonable default for a queue used as a mailbox between a handful of tasks, and a serious one for a streaming pipeline. Nothing warns you — there is no error, no log line, and no slowdown. The only symptom is resident memory climbing at a rate equal to (produce rate − consume rate) × item size, for as long as the job runs. ## A concrete shape of the bug Take a nightly payment reconciliation job: - one producer coroutine streams ledger rows from an upstream source into an `asyncio.Queue`; - a handful of consumer tasks compare each row against the counterpart ledger. Reading rows is nearly free; comparing them is not. On the original input the whole run took about 27 minutes and the queue rarely held more than a few hundred rows, so nobody thought about it. Then a second reconciliation pass was added to chase a floating-point rounding drift between the two ledgers, roughly doubling the per-row consumer cost while the producer stayed exactly as fast. The producer now outruns the consumers by a steady margin; the gap goes into the queue; the job's memory grows monotonically across the whole 27-minute window and the container is OOM-killed near the end — late enough that it looks like a problem with the *last* rows rather than with the first ones. ## Why maxsize is the fix With `asyncio.Queue(maxsize=n)`: 1. `await queue.put(item)` on a full queue creates a putter future, appends it to the queue's putter deque and awaits it, yielding to the event loop — the producer task is suspended, not the thread, so the consumers keep running. 2. Each `get()` frees a slot and wakes the first waiting putter. The consumers' throughput therefore *becomes* the producer's throughput, which is the whole idea of **backpressure**: the slowest stage sets the pace, and memory is capped at `maxsize` items plus whatever the workers hold in flight. The bug converts from an unbounded memory leak into a bounded, visible queueing delay — and a delay you can measure with `qsize()` is a problem you can plan capacity for. ## Choosing the number `maxsize` is a **smoothing buffer** that absorbs jitter, not a reservoir that stores the run. A small multiple of the consumer count (say 2–4×) is a good starting point: enough that a briefly slow consumer does not stall the producer, small enough that the memory ceiling is obvious and the queueing latency stays short. A huge `maxsize` is functionally the unbounded case with a later crash. Watch the steady-state `qsize()`: - **pinned at the maximum** means the consumers are the bottleneck and more of them (or faster ones) is the answer; - **hovering near zero** means the producer is, and a bigger buffer will not help. ## Backpressure versus load-shedding - **Waiting** is the right response when the source can be slowed — a file, a cursor, a paged fetch. - When the source cannot be slowed (an inbound event stream that will not wait), absorbing is not a strategy: use `put_nowait()` and catch `asyncio.QueueFull` to drop, sample or divert the item, and count the drops. Choosing which of the two applies is the actual design decision; `maxsize` alone only decides *where* the pressure lands. ## Related pitfalls - **Deadlock:** a bounded queue makes deadlock possible where an unbounded one merely leaked: if a consumer ever puts back into the same queue it reads from, both sides can end up waiting. - **Cancellation** is safe — a producer cancelled while waiting to put removes its own putter and the item is simply not enqueued — but that item is lost unless you handle it. - Pair the bounded queue with **the completion contract**: `task_done()` in a `finally` and `await queue.join()`, so the job knows when it is genuinely finished rather than merely drained. - And remember the **scope**: `asyncio.Queue` is an in-process, single-event-loop object with no thread-safety and no durability. Backpressure across processes or hosts is a different mechanism entirely.

  • How would you have caught the unbounded growth before the job was OOM-killed?
    Sample `qsize()` on a timer and emit it as a metric alongside items produced and consumed. A queue depth that trends upward across a run is the signal; steady-state depth pinned at the maximum on a bounded queue tells you the consumers are the bottleneck. Without the bound there is nothing to pin against, which is why an unbounded queue is also an unobservable one — the number just grows.
  • The upstream source cannot be slowed down. What replaces waiting?
    Load-shedding. Use `put_nowait()` and catch `asyncio.QueueFull`, then drop, sample or divert the item to slower storage, and count every drop as a first-class metric. Blocking a source that will not wait just moves the failure — it stalls whatever feeds it. The design decision is whether the pressure should propagate upstream or be absorbed by discarding work, and only the domain can answer that.
  • What new failure does a bounded queue introduce that an unbounded one cannot have?
    Deadlock. If a consumer ever puts back into the queue it reads from — a retry or a follow-up item — it can block on a full queue while holding the only capacity to drain it. Use a separate queue for the feedback path, or `put_nowait()` with an explicit overflow policy. Bounding converts a memory problem into a scheduling problem, and scheduling problems can stall.
  • Does a bigger maxsize ever fix a throughput problem?
    Only if the mismatch is transient. A buffer absorbs jitter — a consumer that occasionally stalls — and pays for it with latency. Against a sustained rate difference it changes nothing except how long you run before the memory ceiling is hit. If steady-state depth sits at the maximum, the answer is more or faster consumers, or less work per item, never a larger buffer.

An unbounded queue is a conveyor belt with no end stop: the line never halts, boxes just pile up on the floor until the room fills. A maxsize is the end stop that makes the loader wait for the packer.

saying these in an interview costs you the question

  • "asyncio.Queue() already has a sensible default limit"
  • Setting maxsize to a huge number and calling it bounded
  • Believing a full queue blocks the event loop rather than the producer task
  • Treating a growing queue as a consumer bug with no producer implications
  • Never measuring qsize(), so the growth is invisible until the OOM kill
  • Assuming asyncio.Queue backpressure extends across processes or hosts

context