Design a Distributed Job Queue, stage 6 of 9: change it
Twenty million jobs waiting
Amazon's Builders' Library describes how backlogs turn one outage into two: the recovery itself overloads the dependency that just recovered.
System so far· 7 parts
Select a component to see what it is responsible for and which state it owns.
- 1Web servers → Enqueue gateway: Enqueue job
- 2Enqueue gateway → Kafka: Append to topic
- 3Relay → Kafka: Read topics
- 4Relay → Redis queues: Push at a controlled rate
- 5Workers → Redis queues: Lease jobs
- 6Workers → Databases and services: Do the work
What you need to know
After an outage, the backlog is a second problem. Workers can usually run far faster than the downstream system can absorb, and that system has just recovered: its caches are cold and its connection pools are refilling.
Draining at full speed points several times the normal load at the weakest system in the room.
Work it out
20 million jobs are waiting. The downstream database can take 4,000 extra jobs a second on top of normal traffic. About how many minutes does a safe drain take?Not every job in a backlog is still worth running. Some lose their value with time: a typing indicator from 40 minutes ago, or a push notification about a message the user has already read. A job can carry an expiry, and a worker that sees an expired job acknowledges it without running it.
The useful measure of a backlog is the age of the oldest job per type, not the count. A million fast jobs can be minutes of work; a hundred stuck ones can be a broken feature.
Check
During the drain, a user sends a message. Its notification job is enqueued behind 40 minutes of old notifications. What should happen?