Skip to main content

Distributed Job Scheduler (Distributed Cron)

Sharpened prompt. Design a scheduler that runs 100M scheduled jobs/day — cron-style recurring jobs, one-shot delayed tasks, and multi-hour workflows — firing within 1 second of the target time, never running a payroll job twice, and recovering cleanly when a worker dies mid-execution.

Two words carry this entire design: "exactly" and "on time." Neither is achievable in the strict sense, and knowing precisely what you can promise instead is the interview.

1. Problem framing​

Functional requirements​

  • Register recurring (cron expression) and one-shot (run-at) jobs.
  • At-least-once execution with a strong idempotency story.
  • Retries with backoff; a dead-letter path after exhaustion.
  • Cancel or reschedule a pending job.
  • Visibility: history, current state, duration, failure reason.

Non-functional requirements​

PropertyTargetConsequence
Timing accuracyP99 within 1s of scheduled timeCannot poll a database table every minute
Duplicate executionEffectively never for critical jobsConsensus + fencing + idempotency keys, all three
Scale100M/day ≈ 1,200/sec avg, 50k/sec at the top of the hourThe peak is 40× the average — that is the real requirement
Job duration100ms to 6 hoursHeartbeat leases, not fixed timeouts
DurabilityA registered job survives total cluster lossSchedule state is on disk before ack

Back-of-the-envelope​

Average: 100M/day = 1,157 jobs/sec
Peak: ~60% of cron jobs are on the hour -> ~50,000/sec in the first second of :00
Pending: 500M scheduled-but-not-yet-run rows × 500 B = 250 GB (disk, not RAM)
Near-term: jobs due in the next 60s ≈ 100k -> ~50 MB, fits in a timing wheel in RAM

The 40× peak at the top of the hour is the number that reshapes the design (see 5.6). Most candidates never compute it.

2. High-level architecture​

The split that matters: scheduling (deciding when, needs consensus, low volume) is separate from execution (doing the work, needs throughput, high volume). Conflating them is why homegrown schedulers fall over.

3. Component inventory​

ComponentConcrete choiceWhy this one
Consensus / leasesetcd (Raft)Linearizable compare-and-swap, native leases with TTL, and watch primitives — everything leader election needs
Schedule storeCassandra partitioned by (time_bucket, shard)Writes are append-heavy, reads are always a range scan over the next bucket. Perfect LSM workload
Near-term indexHierarchical timing wheel, in memoryO(1) insert and O(1) tick vs O(log n) for a heap; this is how Kafka and Netty do timers
DispatchKafka (ordered, replayable) or SQS (simple, has native delay)Decouples fire time from execution capacity; absorbs the :00 spike
Long workflowsTemporal / CadenceDurable execution: the workflow's own state survives worker death, so you don't rebuild it
IdempotencyRedis (fast path) + DynamoDB conditional write (durable)Two tiers: cheap check, authoritative claim

4. The toughest parts​

4.1 Split-brain: two leaders, one payroll run​

Why it's hard. Leader election with a lock is not enough. Node A holds the lock, then GC-pauses for 12 seconds. Its lease expires; node B legitimately becomes leader and starts firing jobs. Node A wakes up, still believing it is leader, and fires the same jobs. Both were correct according to their own view. Payroll runs twice. No amount of lock quality prevents this — the problem is that A's in-flight action outlives its authority.

Solution — fencing tokens. Every leadership term carries a monotonically increasing number. Every dispatch is tagged with it. The downstream resource rejects any token lower than the highest it has already seen.

// etcd gives you the fence for free: the key's ModRevision is globally monotonic.
sess, _ := concurrency.NewSession(cli, concurrency.WithTTL(10))
el := concurrency.NewElection(sess, "/sched/partition/7")
if err := el.Campaign(ctx, nodeID); err != nil { return err }

fence := el.Rev() // strictly increasing across all terms, forever

// The worker side enforces it:
// UPDATE job_runs SET status='dispatched', fence=:f
// WHERE job_id=:id AND run_at=:t AND (fence IS NULL OR fence < :f)
// Zero rows updated -> a newer leader already owns this. Drop it silently.

The point to make out loud: "Locks give mutual exclusion only while your process is healthy. Fencing tokens give it even when your process is wrong about being healthy — and a paused process is always wrong about that."

Partition the schedule space (by hash(job_id) % 64) and elect a leader per partition, so 64 leaders share the load and a single leader failure affects 1/64 of jobs.

4.2 Firing on time without scanning a table​

Why it's hard. The obvious design polls: SELECT * FROM jobs WHERE run_at <= now() AND status='pending'. At 500M pending rows this is a punishing index scan; poll every 100ms for 1-second accuracy and you have hammered the database into the ground. A priority heap in memory is O(log n) per operation and cannot hold 500M entries anyway.

Solution — a two-tier time index: durable buckets on disk, a hierarchical timing wheel in RAM.

The wheel is an array of slots and a cursor. Insert is "compute the slot, append to its list" — O(1). Each tick advances the cursor and fires one slot's list — O(1) amortised. Long delays cascade down from coarser wheels, exactly like the hands of a clock.

Wheel 0: 512 slots × 10ms = 5.12s span (millisecond precision)
Wheel 1: 512 slots × 5.12s = 43.7min span (overflow from wheel 0)
Wheel 2: 512 slots × 43.7m = 15.5 days (overflow from wheel 1)
type Wheel struct {
slots [][]*Job
tickMs int64
cursor int
overflow *Wheel // coarser wheel above this one
}

func (w *Wheel) Add(j *Job, delayMs int64) {
if delayMs < w.tickMs*int64(len(w.slots)) {
idx := (w.cursor + int(delayMs/w.tickMs)) % len(w.slots)
w.slots[idx] = append(w.slots[idx], j) // O(1)
return
}
w.overflow.Add(j, delayMs) // too far out; cascade upward
}

func (w *Wheel) Tick(now int64) []*Job {
w.cursor = (w.cursor + 1) % len(w.slots)
due := w.slots[w.cursor]
w.slots[w.cursor] = nil
if w.cursor == 0 && w.overflow != nil {
// A full revolution: pull the next batch down from the coarser wheel.
for _, j := range w.overflow.Drain(now) { w.Add(j, j.RunAt-now) }
}
return due
}

The durable tier holds only what the wheel cannot: a leader loads the next 60 seconds of jobs from Cassandra (WHERE time_bucket = :minute) into the wheel and re-loads every 30 seconds with overlap. Crash recovery is simply reloading the current bucket — nothing is lost because the wheel is a cache of the store, never the source of truth.

4.3 A worker dies six hours into a job​

Why it's hard. A fixed visibility timeout forces an impossible choice. Set it to 10 minutes and long jobs get re-dispatched while still running (duplicate execution). Set it to 7 hours and a crashed worker's job sits dead for 7 hours before anyone notices.

Solution — heartbeat leases with progress checkpoints.

async def run_with_lease(job, worker_id):
lease = await etcd.lease(ttl=30)
ok = await etcd.put(f"/running/{job.id}", worker_id,
lease=lease, prev_kv_absent=True) # atomic claim
if not ok:
return # someone else owns it

async def beat():
while True:
await etcd.keepalive(lease) # extends the 30s TTL
await checkpoint_store.put(job.id, job.progress)
await asyncio.sleep(10)

hb = asyncio.create_task(beat())
try:
await job.execute() # may take 6 hours
finally:
hb.cancel()
await etcd.revoke(lease) # release immediately on exit

The orchestrator watches for lease expiry (etcd delivers a delete event) and re-dispatches from the last checkpoint. Liveness is now detected in 30 seconds regardless of whether the job runs for 100ms or 6 hours.

For genuinely long, multi-step work, do not build this yourself — Temporal persists every step's result in an event history, so a worker crash resumes at the exact step boundary with all local variables restored. Saying "I'd use durable execution rather than reimplementing checkpointing" is a strong, practical signal.

4.4 At-least-once delivery, effectively-once effect​

Why it's hard. Every layer here retries: the queue redelivers on timeout, the dispatcher retries on network error, the operator retries after an incident. Exactly-once delivery does not exist in a distributed system — the FLP result and the two-generals problem see to that. What you can build is at-least-once delivery plus an idempotent consumer, which produces an exactly-once effect.

Solution — a deterministic idempotency key claimed by conditional write.

# The key must be deterministic: the same logical execution must produce the
# same key on every retry, from any node. Never use a random UUID here.
idem_key = f"{job_id}:{scheduled_run_at.isoformat()}:{attempt_group}"

try:
ddb.put_item(
Item={"pk": idem_key, "status": "running", "worker": wid,
"ttl": int(time.time()) + 7 * 86400},
ConditionExpression="attribute_not_exists(pk)",
)
except ConditionalCheckFailedException:
prior = ddb.get_item(Key={"pk": idem_key})["Item"]
if prior["status"] == "succeeded":
return prior["result"] # replay the answer, do not re-run
if prior["status"] == "running" and not lease_expired(prior):
return # a live worker owns it
# else: previous holder died; the lease check above lets us take over

The three-way distinction — succeeded (replay result), running with a live lease (back off), running with a dead lease (take over) — is what makes this correct rather than merely plausible.

For side effects you do not control (sending an email, charging a card), push the idempotency key into the downstream call — Stripe, SendGrid, and Twilio all accept one. The guarantee then holds end to end rather than stopping at your process boundary.

4.5 Retries that do not amplify an outage​

Why it's hard. A downstream service degrades. 50,000 jobs fail. Your scheduler retries all 50,000 after 1 minute — simultaneously, because they all failed simultaneously. The retry wave is larger than the original load and arrives exactly while the dependency is trying to recover. You have built a self-sustaining outage.

Solution — three mechanisms, all required.

  1. Exponential backoff with full jitter: sleep = random(0, min(cap, base × 2^attempt)). Full jitter, not "exponential plus a little noise" — the randomness is what decorrelates the wave, and it measurably beats equal-jitter variants.
  2. A circuit breaker per downstream. After N consecutive failures, stop dispatching to that target entirely and park the jobs. Probe with a single request periodically. This converts 50,000 hopeful retries into one.
  3. A retry budget. Cap retries at ~10% of normal traffic to that dependency (the Envoy/gRPC retry-budget model). Beyond the budget, fail fast. This is the backstop that holds when someone misconfigures the backoff.
retry_policy:
max_attempts: 6
backoff: {base_ms: 1000, cap_ms: 900000, jitter: full}
budget: {ratio: 0.1, min_concurrent: 5}
retry_on: [timeout, 5xx, connect_failure] # never retry 4xx — it will fail again
circuit_breaker: {failure_threshold: 20, open_duration_s: 30, half_open_probes: 1}

4.6 The thundering herd at the top of every hour​

Why it's hard. Humans write 0 * * * *. A large majority of cron jobs are scheduled on the hour, and a huge share of those on the hour and at midnight. Your average is 1,200/sec but 11:00:00.000 brings 50,000. Provisioning for the peak means 40× idle capacity; provisioning for the average means an hourly latency spike and a queue that never fully drains before the next hour.

Solution — deterministic jitter, plus a queue that is allowed to absorb.

# Spread jobs across the minute deterministically: the same job always lands on the
# same offset (so runs stay evenly spaced), but the population is uniformly smeared.
def jittered_run_at(job_id: str, scheduled: datetime, window_s: int = 60) -> datetime:
offset = int(hashlib.sha256(job_id.encode()).hexdigest(), 16) % (window_s * 1000)
return scheduled + timedelta(milliseconds=offset)

Expose this as a job property (jitter_window) and default it to 60 seconds for anything the user did not explicitly mark as time-critical. A report that runs at 11:00:37 instead of 11:00:00 is identical in value; a scheduler that stays up is not.

Then let the dispatch queue do its job: Kafka happily absorbs a 50k/sec burst and workers drain it at their own rate. Autoscale on queue depth and consumer lag, not CPU — CPU is a lagging indicator and by the time it rises the backlog has already formed.

4.7 Catching up after downtime without a stampede​

Why it's hard. The scheduler is down for two hours. On recovery, 8.6M jobs are past due. Firing them all is a self-inflicted DDoS. Skipping them all silently loses work someone depended on. And some jobs are meaningless late — a "send the 9am digest" job firing at 11am is worse than not firing.

Solution — a per-job catch-up policy, declared at registration.

PolicyBehaviourFits
skipDrop all missed runs, resume at the next scheduled timeMetrics snapshots, digests, cache warmers
run_onceCollapse N missed runs into a single executionIdempotent syncs, reconciliation
run_allExecute every missed occurrence, rate-limitedBilling, ledgers, anything where each period matters
deadlineRun only if within max_delay of the targetTime-sensitive notifications

Feed catch-up work through a separate, lower-priority queue with a hard rate limit so recovery traffic can never starve live traffic. And expose a "missed runs" count in the API — silent data loss is the worst outcome of the three.

5. What breaks first​

EventFirst failureMitigation
Leader GC pauseDuplicate dispatchFencing tokens rejected downstream (4.1)
Midnight cron alignmentQueue lag, missed 1s SLODeterministic jitter; scale on lag, not CPU
Downstream dependency downRetry amplificationCircuit breaker + full jitter + retry budget
Worker OOM mid-jobJob stalls until lease expiry30s heartbeat lease; checkpoint + resume
etcd quorum lossNo new leader elections; existing leaders drain their loaded windowWheel already holds 60s of jobs; degrade to read-only, alarm loudly
Clock skew across nodesEarly or late firingchrony with a tight step threshold; a job is authoritative on its run_at, not on any node's local clock

6. Cheat sheet​

  • Split scheduling from execution. Consensus on the small problem, throughput on the big one.
  • Election: etcd lease per partition (64 partitions) + fencing token = ModRevision, enforced by a conditional update downstream.
  • Timing: hierarchical timing wheel for the next 60s (O(1)), Cassandra time buckets for everything beyond.
  • Long jobs: heartbeat leases + checkpoints, or Temporal for durable execution.
  • Correctness: at-least-once delivery + deterministic idempotency key + conditional write = exactly-once effect.
  • Retries: full jitter, circuit breaker, retry budget. All three, not one.
  • Peak: deterministic per-job jitter across a 60s window kills the :00 herd.
  • The one-liner: "You cannot deliver exactly once, so I make the dispatcher at-least-once, the executor idempotent, and the leader fenced — the three together are indistinguishable from exactly-once from the outside."