Skip to content

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
123456SERVICEWeb serversSERVICEEnqueue gatewayLOG / STREAMKafkaWORKERRelayQUEUERedis queuesWORKERWorkersDATABASEDatabasesand services

Select a component to see what it is responsible for and which state it owns.

  1. 1Web servers → Enqueue gateway: Enqueue job
  2. 2Enqueue gateway → Kafka: Append to topic
  3. 3Relay → Kafka: Read topics
  4. 4Relay → Redis queues: Push at a controlled rate
  5. 5Workers → Redis queues: Lease jobs
  6. 6Workers → Databases and services: Do the work

What you need to know

0 of 2 checks done
  1. 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.

  2. 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?
    minutes