Monitoring Background Job Queues¶
A job queue fails quietly. Workers fall behind and jobs wait longer, a worker dies holding jobs, a bug makes every job fail and retry — and the web app keeps enqueueing happily throughout. Monitoring it means reading the queue's own state, and the obvious number is often the wrong one. Measured with taskiq 0.13.0's Redis stream broker on Redis 8.10: after 1,000 jobs had all been processed, XLEN reported 1,000 — the stream keeps acknowledged entries — while the consumer group's lag and pending counts were both 0. With 200 jobs per second offered to a worker that could do about 100, lag grew from 100 to 701 over eight seconds and the age of the oldest undelivered job from 0.5 s to 3.5 s. The worker, running at most 5 jobs at a time, nevertheless held 91–100 jobs as pending, because it read ahead. When it was killed, those 90 jobs stayed pending with no live consumer, idle past 5 seconds, and would be reclaimed only after the broker's default 10-minute idle timeout. This guide derives the numbers that describe a queue's health and the alerts that catch each failure.
Prerequisites¶
- A Redis-backed job queue — taskiq's stream broker here; the ideas carry over to arq, Celery and SQS.
- Queue-age thinking, from monitoring queue depth and item age.
- The topic overview, Background Jobs & Task Queues.
1. Read lag and pending, not the stream length¶
A Redis stream is a log: entries stay after consumers acknowledge them, until the stream is trimmed. Its length counts history, not work. The consumer group tracks what matters — entries not yet delivered (lag) and entries delivered but not acknowledged (pending):
async def queue_state(r, stream: str) -> dict:
group = (await r.xinfo_groups(stream))[0]
return {
"xlen": await r.xlen(stream), # history, not backlog
"lag": group["lag"], # waiting to be delivered
"pending": group["pending"], # delivered, not yet acknowledged
"last_delivered": group["last-delivered-id"],
}
Measured after 1,000 jobs had all completed: xlen 1,000, lag 0, pending 0. A dashboard showing XLEN would have reported a thousand-job backlog on an idle queue, and kept growing with every job ever enqueued. The lag field needs Redis 7 or later; on older versions, derive it from the stream's last entry ID and the group's last-delivered ID. Other brokers have their equivalents — ApproximateNumberOfMessages and ...NotVisible on SQS, list length for list-based brokers — and the same split between waiting and in progress.
Verify: your queue-depth metric is zero when the queue is idle, however many jobs have passed through it.
2. Alert on the age of the oldest waiting job¶
Backlog size alone does not say whether anyone is suffering: 700 waiting jobs is a few seconds of work for a big pool and an hour for a small one. The age of the oldest undelivered job is the delay users actually experience. Stream entry IDs begin with a millisecond timestamp, so it can be read straight from Redis:
async def oldest_undelivered_age(r, stream: str, last_delivered: str) -> float | None:
entries = await r.xrange(stream, min="(" + last_delivered, count=1) # first entry after it
if not entries:
return None
millis = int(entries[0][0].split("-")[0])
return time.time() - millis / 1000
Measured with jobs offered at 200 per second and a worker capacity of about 100: the oldest undelivered job's age rose 0.5, 1.5, 2.5, 3.5 s at two-second intervals — growing at half the elapsed time, exactly the shortfall between arrivals and capacity. Alert when age exceeds what the job's purpose tolerates — seconds for a password-reset email, hours for a nightly report — rather than on a count. Age keeps rising whenever capacity is short, whatever the cause, which makes it the single most useful queue alert.
Verify: an alert fires when the oldest waiting job is older than the job type's latency objective.
3. Account for prefetch in "pending"¶
pending counts every entry delivered to a consumer and not yet acknowledged — including entries the worker has read ahead but not started. In the backlog test, the worker ran at most 5 jobs at a time, yet pending stayed at 91–100: the stream broker reads entries in batches, and those were waiting inside the worker process. That has two consequences. Pending is not "jobs running now"; for that, have workers export their own in-progress count. And everything prefetched by a worker is lost with it until reclaimed:
taskiq worker tasks:broker --max-async-tasks 5 --max-prefetch 0
taskiq's --max-prefetch bounds deliveries beyond execution capacity at the worker level; check what your broker's own read batch size is as well, since it decides how many jobs a dead worker strands. Smaller batches mean more round trips to Redis and less stranded work.
Verify: you know how many jobs a single worker can hold at once, running and prefetched together.
4. Detect jobs held by dead workers¶
When a worker dies, its pending entries stay assigned to its consumer name, with an idle time that keeps growing. Nothing else will run them until something claims them — in taskiq's stream broker, after idle_timeout, which defaults to 10 minutes:
async def stuck_pending(r, stream: str, group: str, idle_ms: int) -> int:
entries = await r.xpending_range(stream, group, min="-", max="+", count=1000)
return sum(1 for e in entries if e["time_since_delivered"] > idle_ms)
Measured six seconds after the worker was killed: 90 pending entries had been idle for more than 5 seconds, lag was 534 and the oldest undelivered job was 10.75 s old — the replacement capacity was not there, and 90 jobs were waiting on a timer. Alert when entries stay pending for much longer than the slowest job takes; it means a consumer has died or hung. Set the broker's reclaim idle time just above the longest legitimate job duration, so stranded jobs come back in minutes rather than at a default chosen for someone else's workload, as discussed in running taskiq with FastAPI.
Verify: killing a worker raises a stuck-pending alert within a few multiples of the longest job's duration.
5. Measure throughput and failures per job type¶
Queue state shows backlog; worker-side metrics show why. Count completions, failures and retries per job type, and record durations, from a middleware so every task is covered:
from taskiq import TaskiqMiddleware
class MetricsMiddleware(TaskiqMiddleware):
async def pre_execute(self, message):
message.labels["_started"] = time.perf_counter()
return message
async def post_execute(self, message, result):
task = message.task_name
JOB_DURATION.labels(task).observe(time.perf_counter() - message.labels["_started"])
(JOB_FAILED if result.is_err else JOB_DONE).labels(task).inc()
broker.add_middlewares(MetricsMiddleware())
A rising failure rate with steady completions points at bad input or a downstream outage; completions falling while the queue's age rises points at capacity or a stuck worker; durations that grow mean the downstream is slowing and capacity will follow. Put enqueue counts beside completion counts per job type: the difference, integrated over time, should match the backlog, and when it does not, jobs are being lost somewhere between the two — the failure that the broker comparison in running taskiq with FastAPI measured.
Verify: per-job-type completions, failures and durations are exported, and enqueues minus completions tracks the backlog.
Verification¶
A job queue is monitored when:
- Depth means waiting work: lag and pending, never stream length.
- The oldest waiting job's age is alerted on, against each job type's latency objective.
- Stranded jobs are detected: pending entries idle longer than the slowest job.
- Workers export completions, failures and durations per job type, and enqueues minus completions matches the backlog.
Diagnostic Hook: chart lag, pending and oldest age on one panel. Lag and age rising together is a capacity problem; pending flat at a high value while lag rises is a dead worker holding jobs; everything at zero while users report missing results points at jobs never enqueued — or acknowledged and lost — which only the enqueue-versus-completion comparison reveals.
Pitfalls & edge cases¶
- Alerting on
XLEN. Measured: 1,000 on an idle queue. - Reading
pendingas "running". Measured: 91–100 pending with 5 running. - Default reclaim timeouts. Stranded jobs waited for a 10-minute default.
- Count-based backlog alerts. Age says whether anyone is waiting too long.
Frequently Asked Questions¶
How do I measure the backlog of a Redis stream job queue?
Use the consumer group's lag (undelivered) and pending (delivered, unacknowledged) from XINFO GROUPS. XLEN counts every entry still in the stream: it read 1,000 after all 1,000 jobs had finished.
What is the best alert for a background job queue?
The age of the oldest waiting job against the job's latency objective: in a backlog it rose from 0.5 s to 3.5 s in 8 s while counts gave no sense of urgency.
How do I detect a dead worker holding jobs?
Look for pending entries idle longer than the slowest job (XPENDING); after a worker was killed, 90 entries sat idle until the broker's reclaim timeout.
Why is pending higher than my worker's concurrency?
Workers prefetch: a worker running 5 jobs at a time held 91 to 100 pending entries in testing. Limit prefetch to strand fewer jobs when a worker dies.
Related¶
- Background Jobs & Task Queues — up to the topic overview.
- Chaining and grouping background jobs — groups that also need watching.
- Concurrent Execution & Worker Patterns — the section overview.