Skip to content

The finished design, decision by decision

How to design a Distributed Job Queue

Not the one correct diagram, but a design you can defend under these constraints: the finished architecture, then every stage's question with the reasoning that answers it, the tradeoffs it accepts, and where another engineer could land differently.

A job queue that keeps working when workers fall behind

The short answer

7 parts, each with one job. The map below shows how requests and data move between them; the stages after it explain why each part is there.

Web servers
Enqueue jobs during requests and return.
Enqueue gateway
Stateless HTTP service that appends jobs to Kafka; keeps broker connections warm.
Kafka
Durable, disk-backed log of jobs; a topic per job group; days of retention.
Relay
Moves jobs from Kafka into Redis at a configurable rate per job type; one owner per topic.
Redis queues
Per-type queues the existing workers already know how to consume.
Workers
Lease a job, run its handler, acknowledge or schedule a retry.
Databases and services
What the jobs actually change.
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

Why does this design work?

The redesign separates holding a backlog from handing out work. Web servers enqueue through a stateless gateway into Kafka, which keeps days of jobs on disk whatever the workers are doing. A relay feeds each job type into Redis at a rate its workers and downstream systems can take, so Redis only ever holds a small working set and can never fill up and freeze.

Job types get separate queues, worker pools and rate limits, so one slow dependency slows only its own jobs. Workers lease jobs, acknowledge after success, retry with backoff, dead-letter poison jobs and drop expired ones, and handlers deduplicate on the job ID because delivery is at least once. Backlogs are drained deliberately, and the new path was rolled out in shadow mode with heartbeats and stage-by-stage counts before any job depended on it.

Invariants, and where they are enforced

  • Enqueueing succeeds even when workers are slow or Redis is full.

    Jobs land in a disk-backed log with days of retention; the relay only moves them into Redis as fast as Redis and workers can take them. Enforced by Enqueue gateway, Kafka.

  • One job type's backlog cannot starve the others.

    Separate topics, queues, rate limits and worker pools per job type or priority. Enforced by Relay, Redis queues, Workers.

  • Every enqueued job runs at least once.

    The log is durable; the relay retries Redis failures; workers lease jobs and acknowledge only after success. Enforced by Kafka, Relay, Workers.

What does it rely on?

  • Handlers are idempotent on the job ID.
  • Each job type's downstream capacity is known well enough to set rate limits.
  • Kafka retention exceeds the longest outage you expect to drain.
  • Job types declare whether they can expire.

What tradeoffs does it make?

ChoiceGainsCosts
Kafka in front of RedisDurable backlog; unchanged workers; incremental rollout.Two systems and two services to operate.
Per-type isolationFailures stay in one lane.More queues and pools to size and monitor.
At-least-once with idempotent handlersNo lost jobs through crashes.Every handler must deduplicate.

What are the reasonable alternatives?

Managed queue (SQS, Cloud Tasks, Pub/Sub)
Better when building fresh, or the team does not want to operate Kafka and Redis.
Kafka-only with retry topics
Better when consumers can be written for Kafka, and per-partition ordering is useful.
Jobs table in the main database
Better when volume is modest and enqueuing must be transactional with the data change (an outbox).

When does it stop working?

  • A backlog outlasts Kafka's retention.
  • Jobs need strict global ordering across types.
  • Enqueueing must be atomic with a database write (which calls for an outbox).

Every stage, decided and explained

Spoilers, for the whole investigation: each stage's question and its answer, the reasoning behind it, and the tradeoffs it accepts. If you have not worked through the stages yet, you may want to do that first.

Work through the stages

Stage 1 of 9 · Model

What the numbers say

Use 86,400 seconds in a day, and about 2 KB per job.

What you need to know first

A day has 86,400 seconds. Divide a daily total by 86,400 to get the average per second, then compare it with the peak to see how spiky the traffic is.

1.4 billion jobs a day. About how many jobs a second is that on average?

About 16,200 per second.

1,400,000,000 ÷ 86,400 ≈ 16,200 a second. The 33,000 peak is about twice that.

A queue holds the difference between what arrives and what is finished. If workers keep up, the queue stays near empty. If they slow down, the backlog grows by (arrival rate − completion rate) every second, for as long as that lasts.

So the question for any buffer is: how long a bad period can it hold, and what does that cost?

Workers stop completely for 10 minutes at the 33,000-a-second peak. Jobs are about 2 KB. About how many gigabytes of jobs pile up?

About 40 GB.

33,000 × 600 seconds ≈ 20 million jobs. 20 million × 2 KB ≈ 40 GB.

On disk, 40 GB is nothing. In RAM, it's 40 GB of headroom you have to keep free on every normal day, just in case.

Where is it cheapest to hold a backlog that might reach 40 GB a few times a year?

A disk-backed log

Disk costs a small fraction of RAM per gigabyte, and a log written sequentially is fast enough to absorb 33,000 jobs a second. Memory is best kept for the small set of jobs workers are about to take.

What the stage asks

Which statements follow?

  1. Holds

    1.4 billion jobs a day is about 16,000 jobs a second on average, so the peak is about twice the average.

    1,400,000,000 / 86,400 ≈ 16,200. The 33,000 a second peak is about 2×.

  2. Holds

    If workers stop completely for ten minutes at peak, about 20 million jobs pile up, roughly 40 GB.

    33,000 × 600 ≈ 19.8 million jobs; at 2 KB each, about 40 GB. In RAM, that is a lot of headroom to keep permanently free for a bad ten minutes.

  3. Fails

    A queue lets producers outpace consumers indefinitely.

    A queue absorbs bursts. If arrivals stay above completions, the backlog grows until something runs out. See Backpressure and capacity.

  4. Fails

    Since Redis is in memory, it is the best place to hold a large backlog.

    Memory is the most expensive and least elastic place to hold a backlog. A disk-backed log holds days of jobs cheaply; memory is for the small working set workers are about to take.

The reasoning

  1. A backlog grows by arrivals minus completions per second; size buffers for the worst period you expect.
  2. Hold large backlogs on disk; keep memory for the small working set workers are about to take.
  3. A queue absorbs bursts; it can't make producers outpace consumers forever.

A job queue has two jobs that pull in different directions: hand work to workers quickly (which suits memory) and hold a backlog safely when workers fall behind (which suits disk). The original design asked Redis to do both, so its worst day was decided by how much RAM happened to be free.

Stage 2 of 9 · Break it

Why the queue stopped draining

Select the lines that describe causes in the design, not just symptoms.

What you need to know first

Redis has a memory limit (maxmemory). When a Redis used as a queue reaches it, commands that would add data fail with an out-of-memory error. Commands that only remove data still work.

The catch is in the details: a reliable dequeue usually moves the job to a "processing" list (RPOPLPUSH) so it isn't lost if the worker crashes. Moving writes a new entry, and writing needs memory.

Redis is at its memory limit. Workers dequeue with RPOPLPUSH, which writes the job into a processing list. What happens to the queue?

It freezes. Enqueues fail (they add data), and dequeues fail too (they add the job to the processing list). Nothing leaves, so no memory is ever freed.

The only action that would make room is the one Redis refuses.

When a shared pool of workers serves several job types, each worker takes whatever job is next. Fast jobs leave quickly; slow jobs stay. Over time, the workers fill up with whichever type is slowest.

200 workers share a queue. Search jobs used to take 30 ms; now they take 1.1 s. If 10% of arriving jobs are search jobs and the rest take 30 ms, roughly what share of busy worker time goes to search?

About 80 %.

Per 10 jobs: 1 search job × 1,100 ms + 9 others × 30 ms = 1,100 + 270 = 1,370 ms of work. Search is 1,100 ÷ 1,370 ≈ 80% of it.

10% of the jobs take 80% of the workers. Everything else waits for the few workers left.

Every worker holds a connection to every Redis instance. During the outage, operators add 200 workers. What does that do?

Adds connections and polling load to the Redis that is already failing

New workers can't dequeue from a full Redis anyway, and each one adds connections and commands to it. The fix adds pressure to the broken component.

What the stage asks

Select the faulty lines.

TimelineIncident timeline
  1. 114:02 db-cluster-7 p99 latency 40 ms → 900 ms (lock contention)
  2. 214:03 search.index jobs now take 1.1 s instead of 30 ms; all workers busy, most of them on search.index

    All job types share one worker pool, so one slow type occupies every worker and starves the rest. Job types need separate pools or limits.

  3. 314:05 enqueue 31,000/s, dequeue 9,000/s
  4. 414:19 redis-jobs-3 used_memory at 98% of maxmemory
  5. 514:20 web: LPUSH jobs:notify failed: OOM command not allowed

    Enqueue is a synchronous write to the memory-bound queue, so a full queue fails user requests. Enqueue needs a durable buffer that does not depend on consumers keeping up.

  6. 614:20 worker: RPOPLPUSH jobs:search processing:search failed: OOM

    Dequeuing copies the job into a processing list, which needs free memory. A full queue therefore cannot drain: the one action that would free memory is the one that is refused.

  7. 714:21 dequeue rate 0/s; queue frozen
  8. 814:24 ops: +200 workers to drain faster

    Every worker connects to every Redis instance, so adding workers adds load to the component that is already failing, and they cannot dequeue anyway.

  9. 914:58 db-cluster-7 recovered; Redis memory manually freed; queue resumes

What the fix has to do

  • Slow downstream → slow jobs → dequeue slower than enqueue → memory fills.
  • Dequeue needed memory, so a full queue could not drain.
  • Enqueue should land in a durable, disk-backed buffer that does not depend on workers.
  • Job types should not share one worker pool.
  • Adding workers increased load on Redis, because of all-to-all connections.Supporting

The reasoning

  1. Ask what happens when a buffer is completely full; make sure the way out doesn't need the exhausted resource.
  2. In a shared worker pool, the slowest job type ends up occupying most workers.
  3. Adding consumers can add load to the component that is failing.

The outage was not one bug but a chain, and the worst link was a drain that required the resource that was exhausted. Systems that fail this way are common: disks too full to delete files, out-of-memory processes that need memory to shut down cleanly. When designing a buffer, ask what happens when it is completely full, and make sure the way out does not need more of what has run out.

Stage 3 of 9 · Decide

A buffer that can hold a bad day

Thousands of job handlers and the worker fleet consume from Redis. Rewriting all of them at once would be a risky project on the most critical async path in the company.

What you need to know first

A log like Kafka stores messages in order, on disk, and keeps them for a set time (say two days) whether or not anyone has read them. Consumers track how far they have read.

So writing to the log never depends on consumers keeping up. Reading can lag by hours and the producer doesn't notice.

Kafka and a Redis list hand out work differently:

Kafka partitionRedis-style queue
Unit of progressan offset: "everything up to here is done"each job, leased and acknowledged on its own
A slow jobholds up the jobs behind it in its partitionholds up only its own worker
Retrying one jobneeds extra machinery (retry topics)built in: let the lease expire or re-enqueue

A Kafka consumer reads jobs 1 to 10 from a partition. Job 3 fails and must be retried later; 4 to 10 succeed. What offset can it commit?

Up to job 2: committing past job 3 would mark it done

An offset means 'everything before this is finished'. Until job 3 is handled, the consumer can't move past it, unless it copies job 3 somewhere else (a retry topic) first.

When a change touches the most critical path in a system, how you get there matters as much as where you end up. A design that keeps thousands of existing workers unchanged can be rolled out one job type at a time and rolled back at any point. A design that rewrites them all has to work everywhere on the first try.

What the stage asks

How do you make enqueueing safe?

  1. Defensible

    Give the Redis clusters several times more memory

    It moves the cliff further away, at a high cost for capacity that sits idle on normal days. The next long downstream incident finds the new edge.

  2. Defensible

    Replace Redis with Kafka and rewrite workers to consume from Kafka directly

    Kafka can hold the backlog, but it is a different consumption model: per-partition order with offsets, where a slow job blocks those behind it in the partition and per-job retries need extra machinery. Rewriting every worker at once is a big-bang change to a critical system.

  3. Sound

    Put Kafka in front of Redis: web servers enqueue to Kafka through a small gateway, and a relay moves jobs into Redis at a controlled rate; workers stay unchanged

    Enqueues land on disk, with days of retention, whatever the workers are doing. Redis now holds only what workers are about to take, and the relay decides how fast it fills. No handler changes. This is what Slack built: a Go gateway (Kafkagate) and a relay (JQRelay), with one relay owning each topic.

  4. Flawed

    Make web servers block and retry when Redis is full

    It moves the backlog into user requests: pages hang while web servers wait for the queue, and the outage spreads from background jobs to the whole product.

What a strong answer covers

  • The backlog lives in a durable, disk-backed log, so enqueue does not depend on consumers.
  • Redis holds only a small working set, filled at a controlled rate.
  • Workers are unchanged, so the change can be rolled out incrementally.
  • A gateway keeps web servers from needing Kafka clients and connections.Supporting

The reasoning

  1. Put the backlog in a durable, disk-backed log so enqueueing never depends on consumers keeping up.
  2. Keep the fast in-memory queue for handing out work, filled at a controlled rate.
  3. Designs that keep existing consumers unchanged can be rolled out incrementally and reversed.

The new design gives each store the job it is good at: Kafka holds the backlog (durable, cheap, days long), Redis hands out work (fast, small). The relay between them is where flow control lives, which is exactly what the old design lacked: nothing ever said "Redis is full, stop filling it".

Keeping the workers unchanged was not a compromise; it was the design. It turned a rewrite into a rollout. See Append-only logs and Backpressure and capacity.

Where another engineer could land differently

Starting from scratch, a managed queue (SQS, Cloud Tasks) or a database-backed job table can serve the same role with far less to operate.

Stage 4 of 9 · Decide

One slow job type

Dropbox's task framework gives each combination of task type and priority its own queue. Slack's relay can rate-limit each job type independently.

What you need to know first

Isolation means giving each kind of work its own capacity, so that a problem in one can't consume another's. In a job system that means separate queues, separate worker pools, and separate rate limits per job type or priority.

A priority is different: it changes the order work is taken in, but everything still shares the same workers.

Notifications have high priority and indexing low. All workers are already busy with slow indexing jobs when a notification arrives. When does it run?

Only when one of the slow indexing jobs finishes and frees a worker

Priority decides what a free worker picks next. It can't interrupt jobs already running, so if every worker is stuck on a slow job, the notification waits for one of them.

Isolation also protects downstream systems. If search is degraded, every indexing job makes it worse. A per-type rate limit caps how fast indexing jobs are released, so the degraded search cluster gets room to recover instead of a flood.

Search is degraded for an hour. Should indexing jobs be dropped, or delayed until it recovers?

Delayed. Dropping an indexing job means that message is never searchable. It's a permanent loss to avoid a temporary slowdown.

Shedding (dropping) is for work that can be lost or that loses its value with time. Indexing keeps its value, so it waits.

What the stage asks

How do you stop slow indexing from delaying notifications?

  1. Flawed

    Keep one shared queue and add workers

    New workers also fill up with slow indexing jobs, and every one of them adds load to the degraded search cluster. Notifications still wait in the same line.

  2. Sound

    Separate queues and worker pools per job type (or priority), with per-type rate limits in the relay

    Indexing can only occupy its own workers, and its rate limit stops it hammering the degraded search cluster. Notifications have their own queue and workers, so they keep flowing. A failure in one dependency stays in one lane.

  3. Defensible

    Add a priority to each job; workers always take the highest priority first

    Notifications jump the line, but slow indexing jobs already running still hold workers, and under sustained load low-priority work can starve completely. Priorities help; separate capacity isolates.

  4. Flawed

    Drop indexing jobs while search is degraded

    Messages would be missing from search permanently. Shedding is right for work that can be lost; indexing must be delayed, not discarded.

What a strong answer covers

  • In a shared pool, slow jobs accumulate in the workers and starve fast ones.
  • Separate queues and pools give each type its own capacity.
  • Per-type rate limits protect the degraded dependency from its own jobs.
  • Distinguishes work that can be delayed from work that can be dropped.Supporting

The reasoning

  1. Separate queues, worker pools and rate limits per job type stop one slow dependency from starving the rest.
  2. Priorities change order but share capacity; under sustained load they don't isolate.
  3. Delay work that keeps its value; only shed work that can be lost.

Isolation is how you stop a local failure becoming a global one. In a shared pool, the slowest job type always wins, because slow jobs are exactly the ones that stay in the workers. Give each type its own capacity and its own throttle, and a degraded dependency only slows the work that depends on it.

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.

What you need to know first

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 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.

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

The lease expires and another worker runs the job again.

The queue only knows the job wasn't acknowledged. It can't tell 'crashed before starting' from 'crashed after finishing', so it runs the job again.

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.

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?

About 127 seconds.

1 + 2 + 4 + 8 + 16 + 32 + 64 = 127 seconds, about two minutes. Doubling makes the total roughly twice the last wait.

Backoff gives a struggling dependency minutes, not milliseconds, to recover, while still retrying.

What the stage asks

Which statements hold?

  1. Holds

    With leases, a worker crashing mid-job does not lose the job.

    The job was never acknowledged, so its lease expires and another worker takes it. That is the point of leasing instead of popping.

  2. Holds

    Because a crashed or timed-out job runs again, handlers must be idempotent or deduplicate.

    The first worker may have finished the work and died before acknowledging. Use the job ID as an idempotency key. See Idempotency.

  3. Fails

    Retrying a failed job immediately, in a loop, gets work done fastest.

    If the failure is a struggling dependency, immediate retries multiply its load. Retry with exponential backoff and jitter. See Retries, backoff and jitter.

  4. Fails

    A job that crashes every worker that runs it is harmless if it is retried forever, since each attempt is cheap.

    A poison job crashes workers repeatedly and takes their in-flight work with them. Cap attempts and move the job to a dead-letter queue for a human to inspect.

The reasoning

  1. Leases make crashes safe by returning unacknowledged jobs, which also means jobs can run twice.
  2. At-least-once delivery requires idempotent handlers, usually keyed on the job ID.
  3. Retry with exponential backoff and jitter, cap attempts, and dead-letter jobs that keep failing.

At-least-once is the practical guarantee, and it has two obligations attached: handlers must tolerate repeats, and retries must be bounded and spaced out. Everything that goes wrong with job systems in production is one of those obligations being skipped. See Delivery guarantees.

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.

What you need to know first

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.

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?

About 83 minutes.

20,000,000 ÷ 4,000 = 5,000 seconds ≈ 83 minutes.

Slower than the workers could go, but the database stays up. A faster drain that knocks it over again takes longer in total.

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.

During the drain, a user sends a message. Its notification job is enqueued behind 40 minutes of old notifications. What should happen?

Fresh interactive work should be able to go first; old ones expire or wait.

A notification that arrives 40 minutes late is nearly useless. Letting new work through, and dropping expired work, keeps the product working during the drain.

What the stage asks

How do you drain it?

  1. Flawed

    Remove all rate limits and drain as fast as the workers can go

    The recovering database receives several times its normal write load at once and falls over again. Meanwhile new jobs wait behind 40 minutes of old ones, so even fresh notifications arrive late.

  2. Sound

    Drain each job type at a rate its downstream can take; let fresh interactive work go first; skip jobs whose usefulness has expired

    The dependency recovers instead of collapsing again. New notifications are not stuck behind old ones. And a job like 'push a typing indicator' or 'notify about a message the user has already read' is dropped at the start of processing because it is past its deadline, which frees capacity for work that still matters.

  3. Defensible

    Process the newest jobs first for interactive types, and the old ones afterwards

    Fresh work gets fresh latency, which is often what users notice. It breaks job types that rely on order, and old jobs still need throttling when their turn comes.

  4. Flawed

    Purge the backlog and start fresh

    Billing events, exports and search indexing would be lost forever. Purging is only acceptable for job types designed to be lossy.

What a strong answer covers

  • Drain at the downstream's capacity, not the workers' capacity.
  • New interactive work is not stuck behind old work.
  • Jobs past their usefulness are dropped cheaply instead of run.
  • Watches the age of the oldest job per type, not just queue length.Supporting

The reasoning

  1. Drain backlogs at the rate the downstream can take, not the rate workers can go.
  2. Let fresh interactive work go first, and drop jobs whose usefulness has expired.
  3. Track the age of the oldest job per type, not just queue length.

A backlog is work promised to the past. Some of it still matters (billing), some has expired (a typing indicator), and all of it competes with the present. Draining well means deciding, per job type, how fast, in what order, and whether at all. See Load shedding.

The durable buffer is what makes these choices possible: with the backlog safely on disk, you can afford to drain it slowly.

Stage 7 of 9 · Break it

Write the worker loop

Write the loop each worker runs for one job type. The queue supports leases, acknowledgements, delayed retries and a dead-letter queue.

What you need to know first

A worker loop handles one job at a time, and each job ends in exactly one of these ways:

OutcomeWhenQueue call
Donethe handler succeededack
Retry laterit failed, attempts remainretryLater with a delay
Dead letterit failed too many timesdeadLetter
Expiredit's too old to matterack without running

Where should the acknowledgement go?

After the handler succeeds

If the worker crashes before acknowledging, the lease expires and the job runs again. A repeat is safe with idempotent handlers; a lost job is not recoverable.

Jitter spreads retries out. If 1,000 jobs fail at the same moment and all wait exactly 4 seconds, they all retry at the same moment, too. "Equal jitter" waits half the backoff plus a random amount up to the other half:

base = min(1000 × 2^attempts, cap)
delay = base / 2 + random() × base / 2

With equal jitter, attempt 3 has base = 1000 × 2³ ms. What's the longest delay it can get, in seconds?

About 8 seconds.

base = 1000 × 8 = 8,000 ms. The delay is between 4,000 ms and 8,000 ms (8 s): half fixed, half random.

What should the loop do when lease() returns no job?

Sleep briefly (a few hundred milliseconds, ideally with a little randomness) before asking again. Without the sleep, every idle worker polls the queue in a tight loop and adds load to it for nothing.

What the stage asks

Implement runWorker.

Reference implementation

const MAX_ATTEMPTS = 8;
const LEASE_MS = 60_000;
const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms));

export async function runWorker(type: string): Promise<never> {
  const handler = handlers[type];
  if (!handler) throw new Error(`no handler for ${type}`);

  for (;;) {
    const job = await queue.lease(type, LEASE_MS);
    if (!job) {
      await sleep(200 + Math.random() * 200);
      continue;
    }

    if (job.expiresAfterMs !== undefined && Date.now() - job.enqueuedAt > job.expiresAfterMs) {
      await queue.ack(job); // nobody needs it any more
      continue;
    }

    try {
      await handler(job.payload, job.id); // handlers dedupe on job.id
      await queue.ack(job);
    } catch (err) {
      if (job.attempts + 1 >= MAX_ATTEMPTS) {
        await queue.deadLetter(job, String(err));
      } else {
        const base = Math.min(1_000 * 2 ** job.attempts, 10 * 60_000);
        await queue.retryLater(job, base / 2 + Math.random() * (base / 2));
      }
    }
  }
}
  • The acknowledgement comes after the handler. A crash between the two means the job runs again when the lease expires, which is why the handler receives the job ID to deduplicate.
  • Backoff doubles per attempt up to ten minutes, with "equal jitter" (half fixed, half random) so retries from a burst of failures spread out.
  • After eight attempts the job goes to the dead-letter queue with the error, rather than retrying forever.
  • Leases must outlast the slowest normal run of the handler; long jobs should extend their lease while they work.

What a strong answer covers

  • Leases a job and acknowledges only after the handler succeeds.
  • Retries failures with exponential backoff and jitter.
  • Caps attempts and dead-letters jobs that keep failing.
  • Drops (acknowledges without running) jobs past their expiry.
  • Passes the job ID to the handler for idempotency.
  • Sleeps briefly when the queue is empty instead of spinning.Supporting

The reasoning

  1. Acknowledge only after the handler succeeds; a crash then means a repeat, never a loss.
  2. Back off exponentially with jitter, cap attempts, and dead-letter what keeps failing.
  3. Check expiry before running, and pass the job ID so handlers can deduplicate.

A worker loop is a small state machine: leased → done, retry later, dead-lettered or expired. Each transition answers a failure the previous stages raised: crashes (leases), repeats (idempotency keys), struggling dependencies (backoff), poison jobs (dead letters), stale work (expiry).

Stage 8 of 9 · Change it

Switch over without an outage

Slack's rollout used a shadow mode in which the relay read jobs from Kafka and discarded them instead of pushing them to Redis.

What you need to know first

Changing a critical path safely has a standard shape:

  1. Build the new path and prove it is alive, with synthetic traffic.
  2. Run it on real traffic in parallel with the old path, with its output thrown away (shadow mode).
  3. Compare the two until they agree.
  4. Move a small, low-risk part over, with a way back.
  5. Move the rest in batches.

Why send heartbeat jobs through every Kafka partition?

So you notice when one partition stops flowing, even if it carries no real traffic

A silent partition looks the same as an idle one. A heartbeat that should arrive every few seconds turns 'nothing happened' into an alert.

In shadow mode, every job is enqueued both the old way and through the new gateway, and the relay discards what it reads. What does that let you check that synthetic tests can't?

That the new path handles the real load and the real mix of jobs (sizes, types, bursts) end to end, while the old path still does all the work. If the counts at each stage match the old path's, nothing is being lost or duplicated. If the new path falls over, no job is affected.

What the stage asks

Put the rollout steps in order.

In this order

  1. 1Deploy the gateway, Kafka and relay; send heartbeat jobs through every partition and alert if any stops arriving
  2. 2Enqueue every job both the old way and through the gateway; the relay runs in shadow mode, discarding jobs
  3. 3Compare job counts at every stage of both paths until they match
  4. 4Move a few low-risk job types to the new path: stop the old enqueue for them and turn off shadow mode
  5. 5Move the remaining job types in batches, keeping the old path available to switch back

The new path carries real traffic long before anything depends on it: heartbeats prove every partition is alive, and shadow mode proves the gateway and relay keep up with the full load while the old path still processes every job. Counts at each stage are the verification. Only then do job types move, a few at a time, each one reversible. It is the same pattern as a data migration, applied to a pipeline. See Online data migrations.

The reasoning

  1. Prove a new critical path on real traffic in shadow mode before anything depends on it.
  2. Heartbeats through every partition turn silent failures into alerts.
  3. Move a few low-risk parts first, keep every step reversible, and verify with counts at each stage.

You cannot test a critical path's new version with synthetic load alone. Running it in parallel with the old one, on real traffic, with its output thrown away, is how Slack (and many others) gained confidence before the switch.

Stage 9 of 9 · Defend it

Defend keeping Redis

Your interviewer: "You now run Kafka and Redis and two new services. Kafka alone could be the queue, or you could have used SQS. Why this?"

What you need to know first

"Why not use X instead?" is a standard interview follow-up. A strong answer has three parts:

  1. The constraint that decided it, which is often not technical elegance but risk, time or what already exists.
  2. What each part of your design is for, so nothing looks accidental.
  3. What the alternative does better, and when you would choose it.

Conceding real costs reads as judgement, not weakness.

Which reason best explains putting Kafka in front of Redis instead of replacing Redis?

Thousands of handlers and the worker fleet already spoke Redis, so this changed the failure mode without rewriting them.

The design is shaped by what existed. It turned a risky rewrite into an incremental rollout.

What the stage asks

Defend the hybrid, concede what the alternatives offer, and say when you would choose them.

Reference answer

The deciding constraint was not technical purity but risk: thousands of handlers and the whole worker fleet already spoke Redis. Putting a durable log in front changed what happens on a bad day without touching any of them, and it could be rolled out job type by job type.

Each part has one job. Kafka holds the backlog on disk for days. Redis hands individual jobs to workers fast, with per-job leases and retries the handlers already understood. The relay is flow control between them: rate limits per job type, and the place a backlog is drained deliberately.

Why not Kafka alone? Kafka consumers track an offset per partition. A slow or failing job blocks the jobs behind it in that partition unless you build retry topics and out-of-order acknowledgement on top. That is solvable, and some companies run job systems that way, but it is a rewrite of every consumer.

What I concede. Two more services and a Kafka cluster are real operational load. Starting from scratch, a managed queue (SQS, Cloud Tasks) with per-type queues, visibility timeouts and dead-letter queues would give most of this with far less to run, and a small system could use a jobs table in its main database.

What a strong answer covers

  • The existing workers and handlers stay unchanged, so the change could ship incrementally with no big-bang rewrite.
  • Explains each part's role: Kafka durable backlog, Redis fast per-job handout, relay flow control.
  • Explains what Kafka-only costs: per-partition head-of-line blocking and per-job retry machinery.
  • Concedes the operational cost and names when a managed queue or Kafka-only design would be better.

The reasoning

  1. Defend a design by naming the constraint that shaped it, including what already existed.
  2. Explain each component's single job: Kafka holds the backlog, Redis hands out work, the relay controls flow.
  3. Concede what alternatives do better, and say when you'd pick them.

A good defence is honest about why the design looks the way it does, including history. "We kept Redis because the workers depended on it, and it let us ship safely" is a better engineering argument than pretending the hybrid is what anyone would design from scratch.

How you did

Now try it as an interview question

  • “Design a distributed job queue or task scheduler.”
  • “Design a background job system like Sidekiq, Celery or SQS.”
  • “Your queue is backing up. What do you do?”
  • “How do you guarantee a background job runs exactly once?”

The interview mode mixes stages from this and other investigations with concept recall and questions about your own projects.

Back to the last stage