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¶
- Python 3.11+, stdlib only.
- tracemalloc, from finding memory leaks in asyncio with tracemalloc.
- Bounded pipelines, from Async Data Pipelines.
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.
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.
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.
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 nomaxsize. 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_nowaitin 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.
Related¶
- Memory & Resource Leaks — up to the topic overview.
- Profiling async services with memray — when the growth is not a queue.
- Resilience, Cancellation & Error Handling — the section overview.