Skip to content

Diagnosing Unbounded Queue Memory Growth

The most common asyncio "memory leak" is not a leak: it is a queue being filled faster than it is emptied. Every item stays referenced until a consumer takes it, so memory grows linearly with the rate difference, and so does the time items wait. Measured with a producer offering about 1,000 items of 10 KB per second to a consumer that handled about 460: an unbounded asyncio.Queue grew by about 464 items and 4.5 MiB every second — 2,318 items and 22.5 MiB after five seconds — and the last items had waited 2.4 s, a number that only goes up. The same pipeline with maxsize=100 stayed at 100 items and 1.0 MiB, items waited 220 ms, and the producer simply slowed to the consumer's pace. A tracemalloc snapshot diff pointed straight at the producer's allocation line, +12.7 MiB in 3 s. This guide confirms a queue is the cause, finds which one, and fixes it.

Prerequisites

1. Recognise the shape of queue growth

Queue growth has a signature that distinguishes it from other leaks: memory rises linearly while load is high and stops — or falls — when load drops, and latency rises in step with it:

q: asyncio.Queue = asyncio.Queue()           # unbounded: maxsize=0


async def producer():
    while True:
        await q.put(make_item())             # never waits: the queue is never "full"


async def consumer():
    while True:
        item = await q.get()
        await handle(item)                   # slower than make_item() arrives
        q.task_done()

Measured per second: 464, 927, 1,392, 1,855, 2,318 items queued; 4.5, 9.0, 13.5, 18.0, 22.5 MiB. Each item's wait is the queue depth divided by the consumer's rate, so latency grows linearly too — 2.4 s after five seconds. A true leak keeps growing at constant load and does not shrink when traffic stops; a queue drains when the producer slows. Watch both memory and the age of items as they are dequeued.

Verify: memory growth correlates with request rate and reverses at quiet times — a queue — or does not — a leak.

Memory after 5 s, producer at ~1,000/s, consumer at ~460/s 2 horizontal bars comparing unbounded queue with the others. Memory after 5 s, producer at ~1,000/s, consumer at ~460/s unbounded queue 22.5 MiB, 2,318 items, 2.4 s wait Queue(maxsize=100) 1.0 MiB, 100 items, 0.22 s wait 10 KB items; tracemalloc totals, sampled every second; the unbounded queue grew 4.5 MiB/s. The bounded queue slowed the producer instead of storing the difference.

2. Measure every queue's depth and the age of its items

You cannot fix the right queue without knowing which one is growing. Wrap queues so they report depth and how long items waited:

import time


class InstrumentedQueue(asyncio.Queue):
    def __init__(self, name: str, maxsize: int = 0) -> None:
        super().__init__(maxsize)
        self.name = name

    async def put(self, item) -> None:
        await super().put((time.monotonic(), item))

    def put_nowait(self, item) -> None:
        super().put_nowait((time.monotonic(), item))

    async def get(self):
        enqueued, item = await super().get()
        QUEUE_WAIT.labels(queue=self.name).observe(time.monotonic() - enqueued)
        QUEUE_DEPTH.labels(queue=self.name).set(self.qsize())
        return item

Depth tells you how much is stored; wait time tells you how stale it is — and since get records the wait of items actually consumed, a rising wait is the earliest warning, visible before memory is large. Export a depth gauge from a periodic task as well, so a queue whose consumer has died (and therefore never calls get) is still visible. The pattern of separating wait time from service time is covered in measuring queue wait and service time separately.

Verify: a dashboard shows depth and p99 wait per named queue.

3. Confirm with a tracemalloc snapshot diff

When metrics are missing, two snapshots taken some seconds apart show which line keeps allocating memory that is not freed:

import tracemalloc

tracemalloc.start(10)                        # 10 frames of traceback per allocation
first = tracemalloc.take_snapshot()
await asyncio.sleep(30)                      # under normal load
second = tracemalloc.take_snapshot()
for stat in second.compare_to(first, "lineno")[:10]:
    print(stat)

# measured: producer.py:7: size=17.0 MiB (+12.7 MiB), count=3549 (+2652), average=5027 B

Measured: over 3 seconds, the top entry was the producer's allocation line, growing by 12.7 MiB and 2,652 objects; the queue's own internal storage (asyncio/queues.py) grew by only 10.8 KiB, because the deque holds references, not copies. So the snapshot points at what is put into the queue, which is the line to look at; the queue that holds those objects is the one fed by that line. Use "traceback" grouping to see the full call path when the allocation happens in a shared helper.

Verify: the top growth line in the diff matches the producer of the queue your metrics flagged.

Diagnosing queue-driven memory growth A flow of 5 stages. Diagnosing queue-driven memory growth growth tracks load? queue, not leak per-queue depth + wait which queue tracemalloc diff producer line grows bound it maxsize + backpressure or add capacity more consumers Measured: the diff showed +12.7 MiB on the producer line, +10.8 KiB in the queue.

4. Bound the queue and let backpressure work

The fix for a producer that outruns its consumer is to make the producer wait. maxsize turns put into a point where the producer pauses:

work: asyncio.Queue = asyncio.Queue(maxsize=100)

async def producer():
    while True:
        await work.put(make_item())          # waits while 100 items are queued

Measured: depth stayed at 100, memory at 1.0 MiB, and items waited 220 ms — the queue size divided by the consumer's rate — while the producer produced only as fast as the consumer consumed (2,427 items in five seconds instead of 4,637). Backpressure has to reach something that can slow down: a network read that stops reading, a request handler that waits, a message consumer with a prefetch limit. If the producer is a request handler that must not wait, use put_nowait and reject when full — the bound then becomes a load-shedding point, as in load shedding when the event loop is overloaded.

Verify: under sustained overload, depth is pinned at maxsize and memory is flat.

5. Size the bound from latency, then add capacity

A bound is a promise about how long items wait: maxsize / consumer_rate. Choose it from the waiting time you can accept, then decide whether the consumer needs more capacity:

def maxsize_for(max_wait_s: float, consumer_rate_per_s: float) -> int:
    return max(1, int(max_wait_s * consumer_rate_per_s))


maxsize_for(0.25, 460)          # -> 115: about a quarter-second of buffering

# If the producer's rate is consistently above the consumer's, no bound fixes throughput:
consumers = [asyncio.create_task(consumer()) for _ in range(4)]   # add consumers for I/O-bound work

A queue absorbs bursts; it cannot absorb a sustained rate difference. If depth sits at the bound most of the time, the system is under-provisioned: add consumers for I/O-bound work, move CPU-bound work to processes, or reduce what the producer generates. Very large bounds look safe and are not — they turn overload into minutes of latency and gigabytes of memory before anything visibly fails.

Verify: at peak load, average depth is well below maxsize, and the bound is only reached during short bursts.

What should this queue do when the consumer falls behind? A decision on Can the producer wait with 4 outcomes. What should this queue do when the consumer falls behind? Can the producer wait? yes (pipeline stage, reader) maxsize + await put backpressure no (request handler) put_nowait, reject when full shed load depth always at the bound add consumers / reduce input capacity problem any queue export depth + wait see it early Unbounded queues store the overload; bounded ones push it somewhere visible.

Verification

Queue growth is under control when:

  • Every queue is bounded, with the bound derived from acceptable wait.
  • Depth and wait time are exported per queue.
  • Backpressure reaches the producer, or excess is rejected explicitly.
  • Sustained overload leads to more capacity, not a larger buffer.

Diagnostic Hook: alert on p99 queue wait per named queue, not just process memory. Wait time rises at the start of an overload, while memory is still small; by the time memory alerts fire, an unbounded queue has usually been accumulating stale work for minutes.

Pitfalls & edge cases

  • asyncio.Queue() with no maxsize. Measured: +4.5 MiB/s under a modest rate difference.
  • Treating queue growth as a leak. It reverses when load drops; the fix is backpressure.
  • Huge bounds. Memory and latency become large before anything fails.
  • Bounds the producer ignores. put_nowait in a loop with retries re-creates an unbounded queue.

Frequently Asked Questions

Why does my asyncio service's memory keep growing under load?

Often an unbounded asyncio.Queue fed faster than it is consumed. In testing, a producer at about 1,000 items/s and a consumer at about 460 grew memory by 4.5 MiB per second.

How do I find which asyncio queue is growing?

Export each queue's depth and the wait time of dequeued items with a name label, or compare two tracemalloc snapshots: the top-growing allocation line is the producer of the queue that is filling.

What maxsize should an asyncio.Queue have?

The acceptable waiting time multiplied by the consumer's rate. With maxsize=100 and a consumer at about 460 items/s, items waited about 220 ms in testing.

What happens to the producer when a bounded asyncio.Queue is full?

await put() waits until a consumer takes an item, slowing the producer to the consumer's pace; put_nowait() raises QueueFull, which you can turn into a rejection.