Skip to content

Monitoring Queue Depth and Item Age

Queue depth is the metric everyone exports and the one that misleads most often. A depth of 500 is fine for a queue drained at 10,000 items per second and a disaster for one drained at 50; a depth of 0 can hide a consumer that stopped five minutes ago if producers stopped too. What users feel is age — how long an item waited before someone started on it. In a simulation with 200 items per second arriving and four consumers keeping up comfortably, depth and oldest-item age were both zero; when the consumers' downstream call slowed from 15 ms to 30 ms at t = 3 s, depth climbed by about 65 items per second and the oldest item's age by about 330 ms per second — 1.67 s after five seconds. Both rose, but only age translated directly into an SLO. This guide instruments a queue for both and turns them into alerts.

Prerequisites

1. Stamp items when they are queued

asyncio.Queue calls _put() and _get() to store and retrieve items — the same hooks PriorityQueue and LifoQueue override. A subclass can wrap every item with its enqueue time without changing callers:

import asyncio
import time


class TimedQueue(asyncio.Queue):
    """A Queue that knows how long its oldest item has waited, and how long each item waited."""

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

    def _get(self):
        enqueued, item = super()._get()
        self.last_wait = time.monotonic() - enqueued
        return item

    def oldest_age(self) -> float:
        return time.monotonic() - self._queue[0][0] if self._queue else 0.0

self._queue is the underlying deque for a FIFO queue; reading its head is O(1). For a PriorityQueue the head is the next item to be served, not the oldest — track the oldest separately if you need it there. last_wait gives the per-item wait as each item is taken, which feeds a histogram.

Verify: put an item, sleep 100 ms, and check oldest_age() reports about 0.1 and last_wait after get() matches.

2. Export depth, age and wait time together

Three numbers, sampled every few seconds and on every get, cover the useful questions:

async def export_queue_metrics(name: str, q: TimedQueue, metrics, every: float = 5.0) -> None:
    while True:
        metrics.gauge("queue_depth", q.qsize(), queue=name)
        metrics.gauge("queue_oldest_age_seconds", q.oldest_age(), queue=name)
        await asyncio.sleep(every)


async def consumer(name: str, q: TimedQueue, metrics) -> None:
    while True:
        item = await q.get()
        metrics.observe("queue_wait_seconds", q.last_wait, queue=name)
        started = time.monotonic()
        try:
            await handle(item)
        finally:
            metrics.observe("queue_service_seconds", time.monotonic() - started, queue=name)
            q.task_done()

Depth is how much is waiting. Oldest age is how late the next item will already be. Wait time, as a histogram, is what each item experienced and gives you percentiles. Service time, measured separately, tells you whether a rise in wait is caused by consumers getting slower or by more arrivals — the distinction that decides between "fix the dependency" and "add consumers".

Verify: under steady load, oldest age stays below one service time and the wait histogram's p99 stays flat.

What depth and age did when consumers slowed down 3 lanes over time. What depth and age did when consumers slowed down consumers 15 ms per item, keeping up 30 ms per item, falling behind depth +65 per second, 328 at t=8 s oldest age 0 ms +330 ms per second, 1.67 s at t=8 s 0 to 8 s → Both grow once capacity drops below arrivals; age is the one expressed in the units of an SLO.

3. Alert on age, not on depth

Depth thresholds need a different number for every queue and every traffic level. Age thresholds come straight from the latency budget: if a job must start within 30 seconds, alert when the oldest item is 15 seconds old.

# Prometheus alert rules
- alert: QueueFallingBehind
  expr: queue_oldest_age_seconds > 15
  for: 2m
  labels: {severity: page}

- alert: QueueAgeGrowing
  expr: deriv(queue_oldest_age_seconds[5m]) > 0.5     # gaining 0.5 s of age per second
  for: 5m
  labels: {severity: ticket}

The derivative alert catches a slow divergence early: in the simulation, age grew at 0.33 s per second the moment capacity fell below arrival rate, long before any absolute threshold would fire. A queue whose age grows steadily is a queue whose consumers can no longer keep up, regardless of its current depth.

Verify: replay the slowdown scenario in staging; the growth alert fires within minutes, before the absolute threshold.

4. Read depth through Little's law

Depth is still useful once divided by throughput. Little's law says average items waiting = arrival rate × average wait, so depth divided by throughput is an estimate of wait time you can compute from two counters even without per-item timestamps:

def estimated_wait(depth: int, throughput_per_s: float) -> float:
    return depth / throughput_per_s if throughput_per_s > 0 else float("inf")

In the simulation, at t = 8 s the queue held 328 items and consumers were completing about 133 per second (4 consumers at 30 ms), so the estimate was about 2.5 s for an item joining the back of the queue — consistent with an oldest-item age of 1.67 s that was still rising. When the estimate and the measured age disagree wildly, something is off with one of the counters: often a consumer that is "running" but stuck, which oldest_age exposes and throughput hides.

Verify: compute the estimate from your exported counters and compare it with measured wait p50; they should be within a small factor.

Queue metrics and what each one answers A grid of 5 rows by 3 columns. Queue metrics and what each one answers metric answers alert on depth how much is waiting rarely: needs per-queue tuning oldest-item age how late the next item already is threshold from the SLO wait histogram what items experienced p99 versus budget service time are consumers slower rise with flat arrivals throughput how fast it drains drop to zero with depth > 0 Age and wait are in the units of the user's experience; depth only makes sense next to throughput.

5. Catch the stalled-consumer case explicitly

The worst failure is a queue whose consumers all hang — waiting on a dependency with no timeout — while depth looks moderate because producers have also slowed. Depth may even be zero if the producers block on a full bounded queue. Two signals catch it:

async def watchdog(name: str, q: TimedQueue, done_counter, stall_after: float = 60.0) -> None:
    last_done, last_change = done_counter(), time.monotonic()
    while True:
        await asyncio.sleep(5)
        now_done = done_counter()
        if now_done != last_done:
            last_done, last_change = now_done, time.monotonic()
        elif q.qsize() > 0 and time.monotonic() - last_change > stall_after:
            log.error("queue %s: %d items waiting, no completions for %.0fs",
                      name, q.qsize(), time.monotonic() - last_change)

"Items waiting and nothing completing" is unambiguous: the consumers are stuck. When it fires, dump the consumer tasks' await chains to see where — the tools are in dumping stacks of a hung asyncio program.

Verify: make the downstream call hang; the watchdog fires within stall_after seconds.

Reading a rising queue age A decision on What is throughput doing while age rises with 3 outcomes. Reading a rising queue age What is throughput doing while age rises? flat arrivals exceed capacity add consumers or shed load falling a dependency got slower check service time zero, items waiting consumers stalled dump their await chains Age says something is wrong; throughput says which kind of wrong.

Verification

Queue monitoring is complete when:

  • Every queue exports depth, oldest age, wait histogram and service time.
  • Alerts use age, with thresholds derived from latency budgets, plus a growth alert.
  • Stalled consumers are detected by completions stopping while items wait.
  • Depth is always read next to throughput, never alone.

Diagnostic Hook: put oldest age, wait p99 and throughput for each queue on one dashboard row. During an incident the shape tells the story at a glance: age rising with flat throughput is a capacity problem; age rising with throughput falling is a slow dependency; age rising with throughput at zero is a stall.

Pitfalls & edge cases

  • Measuring wait at put() time. The wait ends at get(); stamping and measuring on the same side reports zero.
  • Priority queues. The head is the highest-priority item, not the oldest; low-priority items can starve without moving oldest_age().
  • Sampling only depth every minute. Short bursts disappear; export wait as a histogram observed per item.
  • Unbounded label sets. One metric series per tenant or key explodes cardinality; label by queue, not by item.

Frequently Asked Questions

What is the most useful metric for an asyncio queue?

The age of the oldest waiting item, alongside a histogram of per-item wait time. They are measured in seconds, so alert thresholds come straight from latency budgets, unlike depth.

How do I measure how long items wait in an asyncio.Queue?

Subclass asyncio.Queue and override _put to store each item with time.monotonic() and _get to compute the wait when it is taken. Report the head item's age as a gauge and each wait as a histogram observation.

Why is queue depth a poor alert signal?

The right depth threshold depends on throughput, which changes with traffic and differs per queue. The same depth can mean milliseconds or minutes of delay. Divide depth by throughput, or alert on age directly.

How do I detect stuck queue consumers?

Alert when items are waiting but no task_done or completion has been recorded for a period. That combination means consumers are hung, regardless of what the depth is.