Skip to content

Design a Distributed Job Queue, stage 5 of 9: break it

What 'at least once' commits you to

Workers lease a job (it becomes invisible to other workers until the lease expires), run it, then acknowledge. Some jobs fail; some workers crash mid-job.

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. A lease (a "visibility timeout" in SQS) hides a job from other workers for a fixed time while one worker runs it. If the worker acknowledges in time, the job is deleted. If not, the lease expires and the job becomes visible again for another worker.

    That's what makes crashes safe. It's also how the same job ends up running twice.

  2. Here the lease belongs to worker A. Freeze A for longer than its lease and watch what happens to the job. Leave the fencing option off for now.

    A lease that runs out while its holder is asleep. The lease lasts 10 s and A renews it every 3 s, but A freezes at 2 s (a GC pause, a VM migration, a slow disk). Change how long A is frozen, and whether storage checks fencing tokens.
    12 s
    Worker AWorker BStorage
    1. 0 sWorker ATakes the lease with fencing token 33
    2. 2 sWorker AStalls (GC pause) for 12 s
    3. 10 sLockA's lease expires
    4. 10 sWorker BTakes the lease with fencing token 34
    5. 11 sWorker BWrites the job's result
    6. 11 sStorageAccepts B's write (token 34)
    7. 14 sWorker AResumes, still believes it holds the lease, and writes
    8. 14 sStorageAccepts A's write and overwrites B's result

    Two workers acted as the owner. A's lease expired during the pause, B took over, and storage accepted A's late write anyway. The job's result is now whatever A wrote last.

  3. Check

    A worker finishes a job and crashes just before acknowledging it. What happens?