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
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
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.
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 sWorker AWorker BStorage- 0 sWorker ATakes the lease with fencing token 33
- 2 sWorker AStalls (GC pause) for 12 s
- 10 sLockA's lease expires
- 10 sWorker BTakes the lease with fencing token 34
- 11 sWorker BWrites the job's result
- 11 sStorageAccepts B's write (token 34)
- 14 sWorker AResumes, still believes it holds the lease, and writes
- 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.
Check
A worker finishes a job and crashes just before acknowledging it. What happens?So every handler must be safe to run twice: idempotent. The usual way is to record the job ID alongside the effect ("email for job 812 sent") and skip work that has already been recorded. See Idempotency.
Retries need care too. If a job failed because a database is struggling, retrying immediately adds load to that database. Exponential backoff waits 1 s, 2 s, 4 s… and random jitter spreads retries so they don't arrive in waves.
Work it out
A job retries with backoff of 1 s, then 2 s, 4 s, 8 s, 16 s, 32 s and 64 s. About how many seconds pass across those seven waits?