Skip to content
BytePatterns

Design a Distributed Job Scheduler: Leases, Retries, Dead Letters

9 min readBytePatterns

Design a distributed job scheduler: run-at times, leases with heartbeats, fencing tokens, idempotent handlers, backoff, dead letters and the top-of-hour spike.

A job scheduler accepts work now and runs it later: send this reminder at 9:00, rebuild that report every night, retry this webhook in five minutes. The hard part is not the clock. It is the worker that picks up a job and dies halfway through, and the promise that, despite it, no job quietly disappears.

The problem it solves

The lesson's illustrative numbers: five million jobs a day, about 58 a second on average, with a spike on every hour boundary because people schedule things "at 9:00". The requirements:

  • Durable acceptance. Once the caller has a job id, the job will run or be visibly parked; never lost.
  • On time. No job runs before its run-at time, and not long after it.
  • Crash-safe. A worker can die at any moment, mid-job, and the job still completes.
  • Contained failure. A job that can never succeed must not block or starve the rest.

The intuition

Store each job as a row with a run_at time, indexed on it. Then the design is about how workers take rows:

  1. Claim with a lease, not a delete. A worker sets lease_until = now + 30s and its own id on a batch of due rows. If it dies, the lease lapses and the row is simply due again. In PostgreSQL, SELECT ... FOR UPDATE SKIP LOCKED lets many workers claim different rows without waiting on each other (from memory, as of October 2026).
  2. Heartbeat while running. Long jobs extend the lease every ten seconds or so, well inside the lease.
  3. Assume every job can run twice. From outside, a slow worker and a dead worker look identical, so a lapsed lease can hand a job to a second worker while the first is still going. Handlers must be idempotent, and the same at-least-once reasoning as in message queues applies.
  4. Fencing tokens. Each claim increments a token. Completing or failing a job requires the current token, so a worker that wakes up after its lease was taken cannot overwrite the new owner's result.
  5. Back off, then park. A failure reschedules the job with an exponential delay, sparing whatever dependency failed, and after a fixed number of attempts the job moves to a dead-letter table for a human.
  6. Spread the spike. Add a little random jitter to "on the hour" jobs that do not need the exact second.

Recurring jobs fit the same table: store the schedule, and when one run is claimed, insert the next.

Watch it run

The animation follows two jobs. With five million jobs a day, the one thing you cannot promise is that the worker survives the job. A job arrives with a time attached; it is not run now, it is written down, and two rows sit due with nothing handed to anybody yet. A free worker claims j1 by taking a thirty-second lease over the row, not by deleting it. While the handler runs, a heartbeat keeps pushing that lease forward; it finishes, the row is deleted, and that job can never run again. Then j2: same claim, same lease, and this worker dies mid-run. Nothing heartbeats, the lease lapses, and the row is simply due again. But a slow worker looks the same from outside, so the handler must be safe to run twice. A handler that raises instead backs off for two seconds, then four, then eight. After the last attempt the job is parked in the dead-letter queue, so it stops blocking everything queued behind it.

Design a Job Scheduler

Step 1 of 11

Five million jobs a day, and the one thing you cannot promise is that the worker survives the job.

The same interactive animation as the lesson — step through it with the controls.

The code

A toy model: the jobs table is a dictionary and time is a number passed in, so every lease and backoff is deterministic. claim scans for due, unleased rows the way an indexed query would select them:

import random
from collections import Counter

class Scheduler:
    """Jobs table with leases. Time is a number of seconds passed in by the caller."""
    LEASE, MAX_ATTEMPTS = 30, 5

    def __init__(self):
        self.jobs, self.done, self.dead, self.next_id = {}, [], [], 0

    def submit(self, run_at):
        jid = "j%d" % self.next_id
        self.next_id += 1
        self.jobs[jid] = {"run_at": run_at, "lease_until": 0, "attempts": 0, "token": 0}
        return jid

    def claim(self, now, limit=10):
        """Due and not leased: take a lease and a fresh fencing token."""
        due = sorted((j["run_at"], jid) for jid, j in self.jobs.items()
                     if j["run_at"] <= now and j["lease_until"] <= now)
        out = []
        for _, jid in due[:limit]:
            j = self.jobs[jid]
            j["lease_until"], j["token"] = now + self.LEASE, j["token"] + 1
            out.append((jid, j["token"]))
        return out

    def _mine(self, jid, token):
        j = self.jobs.get(jid)
        return j if j is not None and j["token"] == token else None

    def heartbeat(self, now, jid, token):
        j = self._mine(jid, token)
        if j is None or j["lease_until"] <= now:
            return False                      # the lease is gone: stop working
        j["lease_until"] = now + self.LEASE
        return True

    def complete(self, jid, token):
        if self._mine(jid, token) is None:
            return "stale"                    # someone else holds this job now
        del self.jobs[jid]
        self.done.append(jid)
        return "done"

    def fail(self, now, jid, token):
        j = self._mine(jid, token)
        if j is None:
            return "stale"
        j["attempts"] += 1
        if j["attempts"] >= self.MAX_ATTEMPTS:
            del self.jobs[jid]
            self.dead.append(jid)
            return "parked"
        wait = 2 ** j["attempts"]
        j["run_at"], j["lease_until"] = now + wait, 0
        return "retry in %ds" % wait

s = Scheduler()
s.submit(run_at=100)
print(s.claim(now=99), s.claim(now=100))          # [] [('j0', 1)]
print(s.heartbeat(110, "j0", 1), s.claim(now=135))   # True []
print(s.claim(now=141))                           # [('j0', 2)]
print(s.complete("j0", 1), s.complete("j0", 2))   # stale done

At 100 worker one takes the job with token 1; its heartbeat at 110 pushes the lease to 140, so nobody can claim it at 135. Then worker one stalls, and at 141 a second worker takes it with token 2. When the first worker finally reports success, its token is stale and the result is refused. Next, a handler that always raises, and the top-of-hour spike:

t, log = 0, []
s.submit(run_at=0)
while s.jobs:
    for jid, token in s.claim(t):
        log.append(s.fail(t, jid, token))       # a handler that always raises
    t += 1
print(log, t - 1)   # ['retry in 2s', 'retry in 4s', 'retry in 8s', 'retry in 16s', 'parked'] 30

rng = random.Random(36)
hourly = Counter(3600 for _ in range(10_000))
jittered = Counter(3600 + rng.randrange(60) for _ in range(10_000))
print(max(hourly.values()), max(jittered.values()))                 # 10000 193

Five attempts over 30 seconds, then parked. Ten thousand jobs due at exactly 3600 hit the workers in one second; a minute of jitter caps the worst second at 193. Finally, 200 seeded simulations with crashing, failing and stalling workers, checked against what must hold: no job runs early, every job ends exactly once in done or dead, and nothing is left behind:

ok, twice = True, 0
for seed in range(200):
    rng = random.Random(seed)
    s, due = Scheduler(), {}
    for _ in range(rng.randint(1, 25)):
        t0 = rng.randrange(120)
        due[s.submit(t0)] = t0
    effects, running = Counter(), []             # successful handler runs per job
    for now in range(3000):
        for jid, token in s.claim(now, limit=rng.randint(0, 3)):
            ok &= now >= due[jid]                # never early
            outcome = rng.choice(["crash", "fail", "stall", "ok", "ok", "ok", "ok"])
            running.append((jid, token, now + rng.choice([5, 20, 45]), outcome))
        for r in running[:]:
            jid, token, finish_at, outcome = r
            if outcome in ("fail", "ok") and now % 10 == 0:
                s.heartbeat(now, jid, token)     # healthy workers beat every 10 s
            if now >= finish_at:
                running.remove(r)
                if outcome in ("ok", "stall"):   # a stalled worker still finishes
                    effects[jid] += 1
                    s.complete(jid, token)
                elif outcome == "fail":
                    s.fail(now, jid, token)
        if not s.jobs and not running:
            break
    ok &= sorted(s.done + s.dead) == sorted(due) and not s.jobs   # none lost, none twice
    ok &= all(effects[jid] >= 1 for jid in s.done)
    twice += sum(c > 1 for c in effects.values())
print(ok, twice)                                 # True 136

Every invariant holds, and 136 jobs had their handler succeed twice: a stalled worker finished after its lease was taken. The scheduler refused the stale result, but the side effect had already happened, which is why the handler must make a second run harmless, for example by keying its write on the job id.

The complexity

  • Claim: an index range scan on run_at for due rows, O(log n + batch) per poll.
  • Heartbeat, complete, fail: one row update each, guarded by the token.
  • Throughput: limited by how fast the store can update rows; partition the table by job id hash when one database is not enough.

Where it goes wrong

  • Delete on claim. A crash after the delete loses the job.
  • A lease shorter than the heartbeat gap. Healthy jobs get stolen and run twice.
  • Counting attempts only on failure. A job that crashes its worker every time retries forever; count claims too.
  • Retrying in lockstep. Without jitter, a thousand failed jobs retry in the same second and knock the dependency over again.
  • Trusting worker clocks. Compare leases against the database's time, not each worker's.

When it shows up in interviews

As "design a job scheduler", "design a cron service" or "design delayed tasks", and inside larger designs such as notification sends and webhook retries. Follow-ups: a dying worker, exactly-once (honestly: at-least-once plus idempotency), recurring jobs, priorities and the spike. How many jobs one worker runs at once is thread pool sizing.

How to say it in an interview

"Jobs are durable rows with a run_at index. Workers claim due rows by taking a short lease, using SKIP LOCKED or an equivalent so they do not block each other, and heartbeat while running. A lapsed lease makes the job due again, so delivery is at-least-once and handlers must be idempotent; a fencing token stops a stale worker from overwriting a newer result. Failures back off exponentially with jitter and park in a dead-letter table after a fixed number of attempts. Top-of-hour spikes get jitter where exact timing is not required."