System design practice

Durable Job Queue

medium Based on a system published by Slack queues durability availability migration real-world

Slack's job queue: Kafka in front of Redis, so a full Redis loses nothing.

Solve it in your browser Read the lesson first

Slack runs everything too slow for a web request as a background job: posting messages, push notifications, link unfurls, reminders, billing. The web app pushed jobs into Redis, and workers took them from there. Then workers fell behind during a database incident, Redis reached its memory limit and refused new jobs, and the web app could no longer enqueue anything.

The fix: keep Redis and the workers as they are, and put Kafka in front of Redis as a durable buffer.

Functional requirements

Use these use case and scenario names exactly: the traffic, requirements and tests in problem.proschi refer to them.

Scale

Constraints

What is given

problem.proschi declares the web app (four replicas) and the existing workers (eight replicas), and holds the traffic, requirements and tests. Add the gateway, Kafka, the relay, Redis, the connections and the three use cases.

Based on

How your design is checked

You write the design as text in Proschi. Tests run in your browser: a simulation of the traffic above checks latency, availability, cost and what happens when a machine fails. How the simulation works.

Start designing

Lesson · 14 min read

Learn it: Durable Job Queue #

Durable Job Queue: putting a log in front of a queue that can fill up #

Slack runs a lot of work in the background, and all of it went through one Redis-backed job queue. When workers slowed down during a database incident, Redis filled up and refused new jobs. This lesson shows why, and how a durable log in front of Redis fixes it without rewriting the workers.

What you'll learn #

The problem, explained #

Who uses it. The "user" here is Slack's own web application. Every time a request needs slow work done, the PHP web app enqueues a job and moves on. Workers pick jobs up and run them. Nobody outside the company calls this system, but every user feels it when it breaks.

What went wrong. Before the fix, the web app pushed jobs straight into Redis lists and workers popped them off. Redis keeps everything in memory. When workers slowed down, jobs piled up, Redis hit its memory limit, and from then on every enqueue failed. A slowdown in one place (the database the workers used) turned into an outage everywhere (nothing could be enqueued).

Functional requirements.

Non-functional requirements. p99 of Enqueue under 50 ms (99% of enqueues answer faster than that), Enqueue available 99.99% of the time, the job written to a durable store before the web app hears back, any single machine can fail, and the whole thing costs at most $8,000 a month including the existing web app and workers.

What is given. given.proschi fixes the web app (four replicas) and the workers (eight replicas), because they already exist and the point of the exercise is not to rewrite them. It also fixes the traffic: 33k requests per second for each of the three use cases, with 0.1% of relays hitting a full Redis.

What the tests check, in plain words:

Back-of-the-envelope #

Start from the published numbers: 1.4 billion jobs on the busiest days, with a peak of 33,000 a second.

QuantityArithmeticResult
Average rate on a busy day1.4B ÷ 86,400 sabout 16k jobs/s
Peak to average33k ÷ 16kabout 2×
Gateway loadevery enqueue is one HTTP call33k rps
Kafka writesone produce per job33k rps
Relay loadone job delivered per relay step33k rps
Redis load33k pushes from the relay + 33k pops from workers66k rps
Relays that find Redis full0.1% × 33k33 per second
Kafka disk per day (assume ~1 KB per job)1.4B × 1 KBabout 1.4 TB/day

The job size is an assumption, not a published number. But it tells you something important: a day of backlog is more than a terabyte (a few terabytes once Kafka replicates it). That is easy for disks and impossible for the RAM of a single Redis. That is the whole argument for a log in front.

How many replicas? The general formula is:

replicas needed = load ÷ (capacity per replica × target utilisation)

In Proschi a generic service replica ([Go] is a service) handles 2k requests a second. At 100% you would need 33k ÷ 2k = 16.5 replicas just to keep up, and that tier would be saturated. Real systems aim well below 100% because queueing delay explodes near the limit. The model shows a node as "hot" above 70%. So divide by 0.7, round up, and then check one more thing: survive any node failure re-runs the whole analysis with one replica fewer, and the tier must still stay under 100%. Do that arithmetic for both the gateway and the relay; they carry the same 33k rps.

What the simulation does with it. Kafka is a queue node with 50k writes per replica, so three brokers are lightly loaded. Redis is a cache with 100k operations per replica. Unlike a relational database, in the model it takes writes on every replica. The services are the bottleneck. The model uses the M/M/c queueing formula, so a big pool at 70% still answers in about its base 10 ms, while a single replica at 70% would take more than three times that. Enqueue's p99 is the gateway hop plus the Kafka hop with their exponential tails, comfortably inside 50 ms when nothing is hot.

Cost. Each service replica costs $100 a month, a Kafka broker $200 and a Redis replica $150. The web app and workers already take $1,200 of the $8,000. Every gateway or relay replica you add costs $100, so throwing replicas at the problem runs out of budget quickly. That is deliberate: the budget forces you to size, not guess.

Availability. The Enqueue path is the gateway, then Kafka. Each service replica is up 99.5% of the time. But a pool of many replicas is up as long as any one of them is, so in the model the pool is effectively always up. The 99.99% target is easy once nothing on the path is a single replica; with one Kafka broker or one gateway it would not be.

Concepts #

A durable log as a buffer #

A queue in the loose sense is anything that holds work between a producer and a consumer. Two very different things go by that name:

The key property is decoupling under failure: when the consumer side is slow, the producer side keeps working. The backlog lives on disk until consumers catch up.

The trade-offs: a log adds a hop and one more system to operate, and consumers must track their own position (the offset). A log is also a poor fit for per-message features such as delays, priorities or retries acknowledged one message at a time. That is one reason Slack kept Redis behind it rather than pointing workers straight at Kafka.

When not to use it: when the work is cheap and synchronous anyway, or when the volume is small enough that a managed queue with built-in retries (SQS, for example) does the whole job.

At-least-once delivery: commit after the side effect #

A consumer of a log does two things for every message: performs its side effect (here, pushing into Redis) and records that it is done (committing the offset). The order matters:

Losing jobs is worse than running one twice, so job systems pick at-least-once and make the jobs (or the step that consumes them) idempotent, meaning running them twice has the same effect as running them once. The Kafka documentation's "Message Delivery Semantics" section describes exactly this choice.

In Proschi you express the two outcomes as scenarios. The failed side effect uses -x, which models a call that never gets an answer, and the scenario still completes without committing:

title "Buffered hand-off"

app    "App"        [REST API] x2
log    "Job Log"    [Kafka]    x3
mover  "Mover"      [Worker]   x2
target "Downstream" [Database] x2

app   -> log    : produce
log   -> mover  : consume
mover -> target : write

usecase "Move" {
  log -> mover : JobCreated offset 42

  alt "Done" {
    mover   -> target : UPSERT job 42
    target --> mover  : ok
    mover  --> log    : commit offset 42
  } alt "Downstream busy" {
    mover  -x target : UPSERT job 42
    mover --> log    : keep offset, retry with backoff
  }
}

Back pressure belongs in the right place #

Back pressure is how a system tells producers to slow down when it is full. The System Design Primer describes the classic version: bound the queue and answer "busy, try later" when it is full. That is exactly what Redis did, and it was correct behaviour for Redis. The mistake was that the web app received the "busy". The fix moves the pressure boundary: the relay is the only client that ever hits a full Redis. Its answer to "busy" is simply to wait, because the job is safe in Kafka.

A useful rule: put the component that can absorb a backlog between the component that must never stop (accepting work) and the component that can stop (running it).

A stateless gateway in front of a broker #

Kafka clients keep long-lived connections to the brokers, learn which broker leads which partition, and batch messages. A PHP process lives for one request. It would open a fresh connection every time, and the web app would need to know the broker topology. A small stateless service (Slack's Kafkagate) holds the connections and exposes one HTTP endpoint. Stateless means any replica can serve any request, so you scale it by adding replicas and lose nothing when one dies.

The trade-off is one more hop and one more service to deploy. It pays off when clients are many, short-lived or lack a good broker client; it does not when a few long-running services can embed one.

title "Produce over HTTP"

cli "Short-lived Client" [REST API] x2
gw  "Produce Gateway"    [Go]       x2
bus "Event Log"          [Kafka]    x3

cli -> gw  : HTTP
gw  -> bus : produce

usecase "Publish event" {
  cli  -> gw  : POST /events json {"topic": "audit"}
  gw   -> bus : PRODUCE audit
  bus --> gw  : ack
  gw  --> cli : 202
}

Designing it step by step #

1. Scope. Confirm what "enqueue succeeded" must mean: the job is on disk somewhere replicated, not just in a process's memory. Confirm that workers and Redis stay (the migration constraint), the peak rate (33k/s), and that ordering across jobs is not required. Ask what happens today when Redis is full: enqueues fail. That is the failure you are designing away.

2. High-level design. Draw the minimal chain that respects the constraints: web app → something stateless → durable log → something that moves jobs → Redis → workers. Three use cases fall out naturally:

An alternative you should mention and reject: point the workers at Kafka and drop Redis. It removes a component, but it means rewriting every worker at once and losing Redis queue features they rely on. The problem forbids it, and the test Only the relay moves jobs from Kafka to Redis encodes that.

3. Deep dive. Three places deserve time.

4. Wrap-up. Check every requirement against your design: Enqueue p99 is two short hops; the path has no single replica; cost fits with room; every test's flow holds. Then name what you would do next: alert on consumer lag (how far the relay is behind Kafka), cap the relay's push rate so a recovering Redis is not flooded, and add a dead-letter topic for jobs that fail repeatedly.

Common mistakes #

Keep enqueueing into Redis (wrong/enqueue-into-redis). The web app still pushes straight into Redis, even with Kafka sitting next to it. This is the original outage waiting to happen again: the next time workers slow down, Redis fills and the web app's enqueues fail. It fails Enqueuing never touches Redis (the web app is connected to a cache) and The web app enqueues through a stateless gateway.

Write to Kafka, then also to Redis before answering (wrong/enqueue-waits-for-redis). It looks safer: the job is in both places. But the answer now waits for Redis, so a full Redis fails the enqueue exactly as before, and the relay will push a duplicate later. It fails Enqueuing never touches Redis.

A relay that commits even when Redis refused (wrong/relay-ignores-full-redis). The relay gets an out-of-memory error from Redis and commits the offset anyway. The job is now gone from Kafka's point of view and never made it into Redis: silent data loss, the worst kind. It fails A full Redis leaves the job in Kafka, because no scenario shows a failed call to the cache that the use case survives.

The web app talks to Kafka directly (wrong/web-produces-to-kafka). It works on a whiteboard, but in production every short-lived PHP request opens broker connections and must know the cluster layout. It fails The web app enqueues through a stateless gateway.

Rewrite the workers to read Kafka (wrong/workers-read-kafka). Tempting because it deletes a component, but it is a big-bang migration of every job type, and the problem says the workers stay as they are. It fails Only the relay moves jobs from Kafka to Redis.

Other classic mistakes.

In the interview #

Lead with the failure story, because it justifies the whole design: "Accepting work must not depend on the health of the system that runs it. Today it does, because Redis is both the intake and the work queue." Then draw the chain and explain each box in one sentence.

Questions you are likely to get, with short answers:

Further reading #

Now design it

More system design problems