Delayed jobs, cron jobs, and the machinery that runs them while you sleep.
"Design a distributed job scheduler. I want one-shot delayed jobs and cron jobs, retries when things fail, and a way to deal with jobs that never succeed. How do you make sure every job runs on time, even when machines die?"
A distributed scheduler is three boring pieces bolted together: a durable table of (run_at, job) rows, workers that lease due jobs with heartbeats instead of just grabbing them, and retries with exponential backoff ending in a dead-letter queue. "Distributed cron" is just "cron with a database and a lease."
The sentence that carries the interview: "At-least-once delivery with idempotent jobs, leases so a dead worker's job gets picked up, and backoff plus a DLQ so one poison job can't wedge the whole system." Everything else — cron parsing, dashboards, timezones — is UI on top of that core.
"Cron on one server is fine — I'll scale later." A single box running cron has three silent killers: no retries (a 3am deploy kills the 3:05am job and nobody notices until Monday), no failover (the box dies, the jobs die with it), and the thundering herd (fifty 0 * * * * jobs all fire at the top of the hour and OOM the box together). The fix isn't "a bigger cron box" — it's making scheduled time a row in a database that any worker can claim. That one move buys you retries, failover, and observability in a single stroke.
Analogy first: the scheduler is a restaurant kitchen's ticket rail. The math tells you how long the rail needs to be and how many cooks you need before tickets start falling on the floor.
| Quantity | Assumption | Result |
|---|---|---|
| Job rows | 1M delayed jobs alive at once, 1 KB payload each | ~1 GB — fits in Postgres comfortably; the index on run_at is the hot path |
| Due-job poll | Poller grabs jobs with run_at <= now every 500ms, batch of 500 | Sustains 1,000 jobs/sec dispatch; poll query must be an index-only scan or it melts |
| Workers needed | 100 jobs/sec arriving, avg job takes 2s | Little's law: 100 × 2 = 200 concurrent slots — size the fleet from this, not vibes |
| Lease heartbeats | 200 workers × heartbeat every 10s | 20 writes/sec — trivial; heartbeats are cheap, lost leases are expensive |
| Retry storm | 1% of 10k jobs fail, naive immediate retry | 100 extra jobs/sec hammering an already-sick downstream — this is why backoff exists |
Takeaway after the math: the system is dominated by two numbers — the due-job index scan and Little's-law worker sizing. Get those right and the rest is plumbing.
Picture an airport departures board: the schedule is written down centrally, planes (workers) claim flights, and if a plane breaks down the flight goes back on the board for another plane.
flowchart LR
A["Client"] --> B["Scheduler API
validate + persist"]
B --> C[("Job store
run_at indexed")]
D["Cron leader
elected"] --> C
D -- "expands cron into
one-shot rows" --> C
E["Due-job poller"] --> C
E -- "claim with lease
SKIP LOCKED" --> F["Pending queue"]
F --> G["Worker 1"]
F --> H["Worker 2"]
F --> I["Worker N"]
G -- "heartbeat
every 10s" --> C
G -- "fail: backoff" --> J["Retry with
exponential backoff"]
J --> F
G -- "fail: attempts exhausted" --> K["Dead-letter queue
human review"]
G -- "success" --> L["Done + metrics"]
Multiple pollers (or poller threads) racing to claim due jobs need a way to split work without a distributed lock. SELECT ... FOR UPDATE SKIP LOCKED lets each poller grab a batch of unlocked rows and skip rows another poller already locked — no coordination, no waiting, near-perfect sharding of the due set. It's the closest thing databases have to a built-in work queue, and it's why "just Postgres" is a legitimate answer for surprisingly large schedulers.
A worker doesn't take a job — it leases it, like renting a car with the keys due back. It stamps lease_owner + lease_expires on the row and heartbeats to extend. If the worker dies, the lease expires and the poller re-queues the job. No lease, no safety.
sequenceDiagram
participant W as Worker
participant S as Job store
W->>S: claim due jobs (SKIP LOCKED)
S->>W: job #42, lease 30s
W->>W: execute payload...
W->>S: heartbeat (lease +30s)
W->>W: execute payload...
W->>S: complete, status=done
Note over W,S: if heartbeats stop, lease expires
and the poller re-queues job #42
A failed job doesn't retry instantly — that just DDoSes whatever is already sick. Each attempt waits 2^attempt seconds (plus random jitter so a thousand failures don't all retry at the same instant), and after the cap it's parked in the DLQ with its error history attached.
sequenceDiagram
participant W as Worker
participant R as Retry planner
participant D as Dead-letter queue
W->>R: job #7 failed (attempt 1)
R->>R: wait 2s + jitter, requeue
W->>R: job #7 failed (attempt 2)
R->>R: wait 4s + jitter, requeue
W->>R: job #7 failed (attempt 5, max)
R->>D: park with error history, alert human
POST /jobs
{ name, run_at | cron, timezone,
payload, max_attempts, backoff_base }
-> { job_id, next_run_at }
GET /jobs/{id} # status + attempt history
DELETE /jobs/{id} # cancel a scheduled job
POST /jobs/{id}/trigger # run now, out of schedule
GET /dlq # poison jobs awaiting humans
POST /dlq/{id}/requeue # forgive and retry
iduuid pkrun_atindexed — the hot pathstatusscheduled · pending · running · done · deadattempt / max_attemptsretry countinglease_owner / lease_expiresnullablecron / timezonenull for one-shotspayloadjsonbCron jobs are just a template row; the leader materializes the next run_at after each firing. DLQ can be a status, not a separate table — but alerting queries it like one.
| Decision | Option A | Option B | Verdict |
|---|---|---|---|
| Due-job discovery | Poll the DB every 500ms — simple, but polling burns CPU at idle | Push via LISTEN/NOTIFY or a message queue — instant, but another moving part | Poll with SKIP LOCKED until 10k+/sec proves otherwise |
| Cron leadership | DB advisory lock — free if you have Postgres | ZooKeeper/etcd election — proper, but new infra | Advisory lock; promote when you outgrow one DB |
| Delivery semantics | At-least-once + idempotent jobs | Exactly-once (transactions + dedupe) | At-least-once. Exactly-once schedulers are where interviewers go to watch you suffer |
| Where workers live | Long-lived fleet — predictable, but idle cost | Serverless per job — infinite scale, cold starts | Fleet for steady load; burst to serverless |
(cron_id, scheduled_for) so the double-insert fails loudly.run_at by seconds, rate-limit the poller batch.SELECT ... FOR UPDATE SKIP LOCKED for claiming, advisory lock for the single cron leader. One database, zero new infra.min(2^attempt, 300)s + jitter(0–30%), cap 8 attempts, then DLQ + PagerDuty.(cron_id, scheduled_for) with a unique constraint — double-leader inserts die loudly instead of double-firing.lease_owner, lease_expires columns on the whiteboard. Interviewers remember the drawing.A working miniature of the architecture above: a due-job poller, 3 workers with leases, exponential backoff, and a dead-letter queue — all ticking live. Add jobs and watch. Try the poison job to see backoff (2s → 4s → 8s → 16s) end in the DLQ.
empty
empty
empty
empty — poison jobs land here after 5 attempts
Without jitter, a thousand jobs that all failed at 12:00:00 retry at 12:00:02, then 12:00:06 — synchronized stampedes that re-overwhelm the recovering downstream. Adding ±30% random jitter smears retries across time. It's the difference between a retry storm and a gentle rain. The simulator uses pure powers of two so you can see the doubling clearly; production adds the jitter.
Built as a single self-contained file · diagrams render locally, no CDN · back to the takeaway ↑