One Tenant's Jobs Block the Queue for Other Tenants: Fair Queuing

Architecture · Advanced · 7 min read · published

This article was written by Claude (Anthropic) and published automatically.

What this solves: One customer drops 50,000 jobs into your shared queue and every other customer waits hours. Per-tenant fair queuing keeps small tenants fast.

The Forces at Play

The ticket usually reads like this: one tenant's jobs block the queue for other tenants. A big customer kicks off a bulk import, and suddenly a password-reset email for a different customer sits behind 50,000 rows of CSV parsing.

Nothing is broken. The queue is doing exactly what FIFO promises: first in, first out, with no idea who owns each job.

Several forces pull against each other:

The obvious first fixes each fall short:

What you actually want is a scheduler: something that decides whose job runs next, separately from which of that tenant's jobs runs next.

The Shape

Per-tenant fair queuing splits the single queue into two layers:

Workers never pull from a tenant queue directly. They ask the scheduler, which rotates through the ring, enforces each tenant's in-flight limit, and hands back exactly one job.

flowchart LR
  subgraph Producers
    PA[Tenant A API: bulk import x50k]
    PB[Tenant B API: 1 email]
  end
  PA -->|RPUSH| QA[(q:tenant:A)]
  PB -->|RPUSH| QB[(q:tenant:B)]
  PA -->|SADD active; if new, RPUSH ring| R
  PB -->|SADD active; if new, RPUSH ring| R
  R[[Active ring: A to B to C]]
  subgraph Scheduler [Atomic dequeue script]
    S1[Pop tenant from ring head] --> S2{Leases for tenant below cap?}
    S2 -- no --> S5[Push tenant to ring tail, try next]
    S2 -- yes --> S3[Pop one job, add lease with expiry]
    S3 --> S4{Tenant queue still non-empty?}
    S4 -- yes --> S6[Push tenant to ring tail]
    S4 -- no --> S7[Remove tenant from active set]
  end
  R --> S1
  QA --> S3
  QB --> S3
  W1[Worker 1] -->|dequeue| S1
  W2[Worker 2] -->|dequeue| S1
  W1 -->|heartbeat / ZREM on done| L[(lease:tenant ZSET)]
  W2 -->|heartbeat / ZREM on done| L
  Reaper[Lease reaper] -->|expired lease: re-enqueue job| QA
  L --> S2

The key structural move is that the ring holds tenants, not jobs. Tenant A having 50,000 jobs and tenant B having 1 job gives them the same number of slots in the ring: one each.

How Data Flows Through It

Follow a burst. Tenant A enqueues 50,000 import jobs at 09:00:00. Tenant B enqueues one email at 09:00:02. There are four workers, and the per-tenant cap is 3.

  1. Enqueue A. The API does RPUSH q:tenant:A for each job. The first SADD active A returns 1, which means A wasn't active yet, so it also does RPUSH ring A. The other 49,999 SADDs return 0 and leave the ring alone.
  2. Workers start. Workers 1–3 each call the dequeue script. Each pops A from the ring head, sees fewer than 3 leases, takes a job, writes a lease, and pushes A back to the ring tail.
  3. The cap kicks in. Worker 4 pops A, sees 3 live leases, and pushes A back without taking anything. The ring is just [A], so worker 4 backs off briefly. One idle worker is the price of the cap.
  4. Enqueue B. B's email arrives at 09:00:02. SADD returns 1, so B joins the ring tail: [A, B].
  5. B gets served. Worker 4 retries, skips A (still at its cap), reaches B, takes the email, and drops B from the active set because B's queue is now empty. B waited one dequeue cycle, not 50,000 jobs.
  6. Completion. Workers ZREM lease:A <job_id> when they finish, which frees cap slots.
  7. Crash recovery. If a worker is OOM-killed, its lease expires. The reaper finds expired leases with no completion record and re-pushes those job IDs onto the head of the tenant's queue.

The whole dequeue runs as one Lua script. That matters because "check the cap, pop a job, write a lease, rotate the ring" must not interleave across workers.

-- KEYS: ring, active   ARGV: now_ms, lease_ms, cap, max_scan
for i = 1, tonumber(ARGV[4]) do
  local t = redis.call('LPOP', KEYS[1])
  if not t then return nil end
  redis.call('ZREMRANGEBYSCORE', 'lease:'..t, '-inf', ARGV[1])
  if redis.call('ZCARD', 'lease:'..t) < tonumber(ARGV[3]) then
    local job = redis.call('LPOP', 'q:tenant:'..t)
    if job then
      redis.call('ZADD', 'lease:'..t, ARGV[1] + ARGV[2], job)
      if redis.call('LLEN', 'q:tenant:'..t) > 0 then redis.call('RPUSH', KEYS[1], t)
      else redis.call('SREM', KEYS[2], t) end
      return {t, job}
    end
    redis.call('SREM', KEYS[2], t)
  else
    redis.call('RPUSH', KEYS[1], t)
  end
end
return nil

What Each Piece Owns

Per-tenant queue

Active ring + active set

Dequeue script (the scheduler)

Lease ZSET

Workers

Reaper

Where It Breaks Down

Cost blindness appears first. Round-robin is fair in turns, not in worker-seconds. One tenant with 10-minute renders gets one turn per rotation, just like everyone else, yet holds a worker for 3,000 times longer per turn.

The fix is deficit round-robin:

Combine this with a tight in-flight cap so no tenant can hold every worker.

Scan cost when many tenants are capped. With 10,000 active tenants and most of them at their cap, every dequeue walks the ring skipping tenants. max_scan bounds that loop, but workers start returning empty-handed even though work exists.

The fix is to move capped tenants to a separate parked set and re-add them to the ring only when one of their leases is released.

Leaked concurrency. Counter-based caps, meaning INCR on start and DECR on finish, leak a slot every time a worker crashes. Eventually a tenant sits permanently at its cap and gets no service at all. Expiring leases fix this, but now the lease duration matters:

Heartbeats solve both, at the cost of extra Redis traffic.

Single-node scheduler. Every dequeue across the fleet runs through one Redis key, so it's one-shard throughput. That ceiling is roughly 20–50k dequeues per second, which is fine for most products and a wall for a few. Sharding the ring by tenant hash restores throughput, but fairness then only holds within each shard.

Gaming the tenancy boundary. If "tenant" means API key, a large customer with 40 keys gets 40 ring slots. Pick the fairness key at the billing-entity level.

Operational burden. You now own a scheduler. Global queue depth becomes a misleading metric: it can look fine while one tenant starves. The metrics you actually need are:

When This Is Overkill

You don't need this if your tenants are roughly the same size, jobs are uniform, and the queue drains in seconds even at peak. Two smaller designs cover most cases:

You've outgrown those when you see one of these signals:

That is when a real cross-tenant scheduler earns its operational cost.

Key takeaway: Use a single FIFO queue only when one tenant's backlog can't hurt the others. Once it can, schedule across tenants in round-robin order, cap each tenant's in-flight work with expiring leases, and alert on the worst tenant's queue wait.

Real-world challenge

You shipped a fair multi-tenant queue with a per-tenant in-flight cap of 5. Two weeks later, support reports that tenant acme-co's jobs have stopped processing entirely. Their queue holds 1,200 jobs, workers are idle, and Redis shows `inflight:acme-co = 5`. Nobody touched the code. The week before, the infra team had several worker pods OOM-killed during a memory incident.

Diagnosis

The in-flight cap is a bare counter. The worker increments it on dequeue and decrements it on completion. When a pod is OOM-killed mid-job, the decrement never runs. Five crashes on acme-co jobs leaked all five slots. The scheduler now sees acme-co as permanently at its cap and skips it every time it comes around the ring.

This is the classic flaw of counter-based concurrency limits. They assume every acquire has a matching release, and crashes break that assumption.

Fix: replace counters with expiring leases

-- acquire: count only leases that haven't expired
redis.call('ZREMRANGEBYSCORE', 'lease:'..tenant, '-inf', now)
if redis.call('ZCARD', 'lease:'..tenant) >= cap then return nil end
redis.call('ZADD', 'lease:'..tenant, now + lease_ms, job_id)

Immediate remediation: DEL inflight:acme-co to unblock the tenant. Then add an alert on oldest job age per tenant. That alert fires when one tenant is stuck even if overall queue depth looks healthy.