Load Testing Queue Throughput

HTTP load tools measure request latency; a queue load test has to measure something different — how fast work is completed and how long it waits — because the enqueue call returns in a millisecond no matter how far behind the workers are. This guide designs and runs such a test as part of Capacity Planning for Job Queues in Observability & Monitoring for Job Queues.

Problem Statement

Before a product launch, a team wants to know the maximum sustained throughput of its Celery pipeline (API enqueue → Redis → 20 worker pods → Postgres) and where it breaks. A previous test with an HTTP load tool reported "10,000 requests per second, p99 12 ms" and everyone relaxed; in production, the queue fell behind at 900 jobs per second. The HTTP tool had measured only how fast jobs could be enqueued. You need a test that drives realistic jobs at controlled rates, reports completion throughput and queue wait from the system's own metrics, finds the saturation point, and names the bottleneck — workers, broker, database, or an external API.

Prerequisites

  • A staging environment with production-like worker counts, broker, and database sizes (or a known scaling factor).
  • Queue metrics: enqueue rate, completion rate, queue wait histogram, job duration histogram, busy slots — see Prometheus Metrics for Workers.
  • Resource metrics for the broker and database (CPU, connections, latency).
  • External dependencies stubbed with realistic latency, or sandbox endpoints that can take the load.

Step 1 — Define What "Throughput" Means for This Test

Write down the questions and the measurements that answer them before running anything:

# loadtest.yaml — the contract for the test
questions:
  - max sustained completion rate with p95 queue wait < 10 s
  - which resource saturates first
  - backlog drain rate after a 10-minute burst at 2x max
measure_from: prometheus            # system metrics, not the generator
signals:
  arrival:     sum(rate(jobs_enqueued_total{queue="events"}[1m]))
  completion:  sum(rate(jobs_completed_total{queue="events"}[1m]))
  queue_wait:  histogram_quantile(0.95, sum by (le)(rate(job_queue_wait_seconds_bucket{queue="events"}[1m])))
  backlog:     sum(queue_depth{queue="events"})
  busy_ratio:  sum(worker_busy_slots) / sum(worker_total_slots)
job_mix:                             # match production proportions
  ingest_event: 0.70
  enrich_event: 0.25
  rebuild_rollup: 0.05

The key rule: the system is saturated when completion rate stops following arrival rate and the backlog grows, not when the generator reports errors. Enqueue will keep succeeding long after workers have fallen behind.

Saturation: completions stop following arrivals During a stepped load test, arrival rate increases in steps. Completion rate tracks it closely up to about 1,100 jobs per second. Beyond that, completions plateau while arrivals keep rising, and the backlog line starts to climb. The plateau is the sustained throughput ceiling, even though enqueue calls continue to succeed. Stepped load test arrivals (steps) completions plateau ~1,100/s backlog grows

Step 2 — Build a Generator That Enqueues Real Jobs

Generate load by enqueuing through the same code path production uses, with realistic arguments. A dedicated generator process, rate-controlled with a token bucket, is more precise than an HTTP tool and avoids measuring the API layer by accident.

# loadgen.py — enqueue a job mix at a controlled rate, in steps
import random, time, argparse
from app.tasks import ingest_event, enrich_event, rebuild_rollup
from fixtures import sample_event_ids, sample_tenants

MIX = [(ingest_event, 0.70), (enrich_event, 0.25), (rebuild_rollup, 0.05)]

def pick():
    r, acc = random.random(), 0.0
    for task, w in MIX:
        acc += w
        if r <= acc:
            return task
    return MIX[-1][0]

def run(rate: float, seconds: int, batch: int = 50) -> None:
    interval = batch / rate
    deadline, next_t = time.monotonic() + seconds, time.monotonic()
    while time.monotonic() < deadline:
        for _ in range(batch):
            t = pick()
            t.apply_async(kwargs={"event_id": random.choice(sample_event_ids),
                                  "tenant_id": random.choice(sample_tenants)},
                          headers={"loadtest": "run-2026-09-18"})    # tag for filtering
        next_t += interval
        time.sleep(max(0.0, next_t - time.monotonic()))

if __name__ == "__main__":
    p = argparse.ArgumentParser()
    p.add_argument("--steps", default="200,400,600,800,1000,1200,1400")
    p.add_argument("--step-seconds", type=int, default=600)
    a = p.parse_args()
    for r in map(float, a.steps.split(",")):
        print(f"step {r}/s"); run(r, a.step_seconds)

Use argument samples drawn from production-shaped fixtures (tenant sizes, event sizes), because a test where every job touches the same small tenant measures a warm cache, not the system. Run several generator processes if one cannot reach the target rate, and confirm the arrival metric — not the generator's own count — matches the intended rate.

Step 3 — Run Stepped Load and Hold Each Step

Increase load in steps and hold each for long enough to reach steady state — at least 5–10 minutes, longer if caches, autoscalers, or connection pools take time to settle. Ramps that move continuously never show where the system stabilises.

# Freeze autoscaling for the capacity test so you measure a fixed fleet
kubectl scale deployment events-worker --replicas=20
kubectl patch scaledobject events-worker --type merge -p '{"spec":{"minReplicaCount":20,"maxReplicaCount":20}}'

python loadgen.py --steps 200,400,600,800,1000,1200,1400 --step-seconds 600

Test the fixed fleet first to find its ceiling; test autoscaling separately (Step 6), or its reaction time will blur the saturation point.

Step 4 — Read the Results from System Metrics

For each step, record the steady-state values over the last few minutes of the step. A small script against the Prometheus API keeps it reproducible.

# report.py — steady-state values per step
import requests

PROM = "http://prometheus.staging:9090/api/v1/query"
Q = {
    "arrival":    'sum(rate(jobs_enqueued_total{queue="events"}[3m]))',
    "completion": 'sum(rate(jobs_completed_total{queue="events"}[3m]))',
    "wait_p95":   'histogram_quantile(0.95, sum by (le)(rate(job_queue_wait_seconds_bucket{queue="events"}[3m])))',
    "dur_p50":    'histogram_quantile(0.5, sum by (le)(rate(job_duration_seconds_bucket{queue="events"}[3m])))',
    "busy":       'sum(worker_busy_slots{queue="events"}) / sum(worker_total_slots{queue="events"})',
    "db_cpu":     'avg(rate(node_cpu_seconds_total{instance=~"pg-.*",mode!="idle"}[3m]))',
    "redis_cpu":  'rate(redis_cpu_sys_seconds_total[3m]) + rate(redis_cpu_user_seconds_total[3m])',
}

def snapshot(at_unix: float) -> dict:
    return {k: float(requests.get(PROM, params={"query": q, "time": at_unix}).json()
                     ["data"]["result"][0]["value"][1]) for k, q in Q.items()}

A typical result table from this kind of test:

Arrival Completion Wait p95 Duration p50 Busy slots DB CPU Redis CPU
400/s 400/s 0.3 s 45 ms 23% 31% 12%
800/s 800/s 0.6 s 52 ms 48% 63% 24%
1,000/s 998/s 2.1 s 71 ms 68% 82% 30%
1,200/s 1,090/s growing 118 ms 97% 96% 33%

The ceiling is about 1,100 jobs per second. The bottleneck is visible in the columns: duration p50 more than doubles as database CPU approaches 100%, while Redis sits at a third of a core. Adding workers would not help — they would wait on the same saturated database, making each job slower. The capacity arithmetic behind reading this table is in calculating worker count with Little's Law.

Which resource hit the wall? At the 1,200 jobs per second step, database CPU is 96 percent and worker busy slots are 97 percent, while Redis CPU is 33 percent. Job duration doubled as the database saturated, so workers are busy waiting on the database. The database is the bottleneck; more workers would not raise throughput. Utilisation at 1,200 jobs/s offered Postgres CPU 96% worker busy slots 97%, waiting on DB Redis CPU 33%: headroom Busy workers are a symptom here, not the cause: their jobs doubled in duration.

Step 5 — Test a Burst and Measure the Drain

Sustained throughput is one question; recovery from a burst is another. Run at 70% of the ceiling, burst to twice the ceiling for ten minutes, return to 70%, and measure how long the backlog takes to clear.

python loadgen.py --steps 770,2200,770 --step-seconds 600   # baseline, burst, recovery
# Backlog and its drain rate during recovery
sum(queue_depth{queue="events"})
-deriv(sum(queue_depth{queue="events"})[5m:30s])            # jobs/s drained

Compare the measured drain time with the prediction backlog / (ceiling − arrival). A drain much slower than predicted usually means a secondary effect: backlog jobs are older and miss caches, or retries from burst-time failures are adding load. The live-incident version of this calculation is in forecasting backlog drain time.

Step 6 — Test Autoscaling Separately

With the fixed-fleet ceiling known, re-enable autoscaling and run the burst again. Measure the time from burst start to the first new pod processing jobs, and the peak backlog.

kubectl patch scaledobject events-worker --type merge -p '{"spec":{"minReplicaCount":6,"maxReplicaCount":40}}'
python loadgen.py --steps 300,2200,300 --step-seconds 600
kubectl get events --field-selector involvedObject.kind=HorizontalPodAutoscaler -w

If scale-up takes three minutes and the burst adds 1,100 jobs per second of excess, that is 200,000 jobs queued before new capacity arrives — often the most important number the test produces.

The cost of scale-up lag The burst starts at minute zero. The autoscaler notices the backlog after about thirty seconds, requests pods, and they take another two and a half minutes to schedule, pull images, and start. During those three minutes the backlog grows by about 200,000 jobs. Once the new pods are processing, the backlog stops growing and begins to drain. Burst begins, capacity arrives three minutes later detect schedule, pull image, start: ~2.5 min new pods processing, backlog drains backlog +200k Pre-pulled images, a warm node pool, and a higher floor all shorten the red segment.

Tuning is covered in scaling workers with KEDA on queue length.

Verification

A trustworthy test run has three properties you can check:

[x] arrival metric matched the intended rate within 2% at every step
[x] each step held steady state for >= 5 minutes (completion flat, backlog flat or growing linearly)
[x] the bottleneck named in the report was confirmed by a follow-up (e.g., DB scaled up -> ceiling moved)

The third is the real test of the analysis: change the suspected bottleneck and run the step again. If the ceiling moves, the diagnosis was right.

Gotchas & Edge Cases

Measuring the generator. Generator-side latency and error counts describe enqueue, not processing. Always report from queue and worker metrics.

Unrealistic data. Jobs that all hit one tenant, one small file, or one cache key measure a hot path. Sample arguments from production-shaped fixtures.

Stubbed dependencies that are too fast. A payment API stub returning in 1 ms makes jobs look five times faster than production. Add realistic latency distributions to stubs.

Leftover load. A backlog from one step bleeds into the next. Let the queue drain between runs, and tag test jobs (as with the header above) so they can be purged.

FAQ

Can I use k6 or Locust for this? Yes, as the enqueue driver — k6 can call your enqueue API or push to the broker via an extension, and Locust can call task code directly. Keep the measurement in Prometheus either way.

How long should a full test take? A stepped capacity test with 7 steps of 10 minutes, a burst test, and an autoscaling test is about 2–3 hours. Automate it so it can run before every major launch.

Should I test in production? Shadow or replay traffic in production is valuable but risky. Staging with a known scaling factor answers most questions; if you do test in production, cap load well below the staging ceiling and have a kill switch.

Related