System Design LabSystem Design QuestionsDesign Distributed Job Scheduler

Design Distributed Job Scheduler

MediumDistributed Systemsschedulingqueuedistributedconcurrency

Question Overview

Design a shared scheduling service that runs one-off and cron jobs on a worker fleet. The key challenges are finding due jobs efficiently, coordinating schedulers without double-firing, surviving crashes with at-least-once execution, and absorbing top-of-hour bursts.…

Sign up to see the full question and AI interviewer

Requirements

  • Create, update, pause, and delete one-off and cron jobs, with cron evaluated in each job's own time zone
  • Per-job payload, timeout, and retry policy with exponential backoff and a maximum attempt count
  • Record every run's status, attempts, timings, and output; support manual triggers and history queries
  • Runs start within 2 s of their scheduled time p99 under normal load
  • At-least-once execution with a stable idempotency key; no run silently skipped after a node or worker failure
  • No single point of failure in scheduling or dispatch, and 90 days of durable run history

Back-of-the-envelope numbers

  • Executions: 100M/day ÷ 86,400 s ≈ 1,160 runs/s average, ~3.5K/s at 3× peak
  • Concurrency (Little's law): 1,160 runs/s × 20 s ≈ 23K runs in flight on average, ~70K at 3× peak
  • Worker capacity: 8,000 hosts × 20 slots = 160K slots, more than double the ~70K peak concurrency
  • Midnight burst: 100K runs due at 00:00 UTC; starting them within 10 s needs 10K dispatches/s (~8.6× average) and ~123K slots with the 23K baseline
  • Job definitions: 50M × ~1 KB ≈ 50 GB, modest enough for a lightly sharded database with the next_run_at index in memory
  • Run history: 100M/day × ~500 B ≈ 50 GB/day × 90 days ≈ 4.5 TB, partitioned by day so expiry is a partition drop
  • Heartbeats: 23K in-flight runs ÷ 10 s interval ≈ 2.3K/s average, ~7K/s at peak, kept in Redis rather than the history store

Key components

  • Job store: sharded database holding job definitions and next_run_at, indexed by (partition, next_run_at) for cheap due-job scans
  • Scheduler partitions: jobs hashed into ~1,024 partitions, each owned by exactly one scheduler node through leases in etcd or ZooKeeper
  • Fencing tokens: each lease carries an increasing epoch checked on every write, so a paused former owner cannot fire runs
  • Run creation: insert a run keyed by (job_id, scheduled_time) under a unique constraint and advance next_run_at in the same transaction
  • Dispatch queue: due runs go to a durable queue, pre-enqueued minutes ahead with delayed delivery to absorb top-of-hour bursts
  • Workers: claim runs with a lease, heartbeat every 10 s, and runs whose heartbeats stop are marked lost and retried per policy
  • Retries and dead letters: exponential backoff with jitter, a dead-letter queue after max attempts, and alerts to the job owner

Common mistakes

  • One scheduler polling SELECT * WHERE next_run_at <= now, which is both a single point of failure and a table-scan bottleneck
  • Running several uncoordinated scheduler replicas for availability, so every replica fires every job
  • Promising exactly-once execution; crashes and partitions allow only at-least-once delivery plus idempotent handlers
  • Leader election without fencing tokens, letting a GC-paused old leader wake up and dispatch alongside the new one
  • Computing cron times in server local time, so daylight-saving transitions skip or double-fire runs
  • Retrying failures immediately without backoff or jitter, hammering an already struggling downstream dependency

Likely follow-ups

  • How would you support dependencies, where job B runs only after job A succeeds?
  • If scheduling is down for 10 minutes, should missed runs fire all at once, fire once, or be skipped?
  • How would you stop one tenant with a million jobs from starving everyone else?
  • How would you rebalance partitions when adding scheduler nodes without double-firing runs?
  • How would you cancel a long-running job that is already executing on a worker?
  • How would you spread top-of-hour load when thousands of teams all schedule at minute zero?

No community solutions yet

Be the first to publish your solution

Practice ‘Design Distributed Job Scheduler’ with an AI Interviewer

Get scored feedback on your diagram, scalability approach, and trade-offs. Free while we grow — up to 3 full interviews a day.