← All writing
articleJul 06, 202419 min read

Distributed Job Scheduler: A Cron Table That Ten Thousand Services Write To

Leader-elected scheduling, a timing wheel, DAG dependencies, and the exactly-once claim that is actually at-least-once plus idempotency.

SchedulingDistributed SystemsDAGArchitecture
Distributed Job Scheduler: A Cron Table That Ten Thousand Services Write To cover illustration

A cron table is a sorted list of things that need to happen and a process that wakes up to do them. Distributed, it is the same thing with three additional problems: several processes are awake, they must not duplicate each other’s work, and the jobs themselves are not reliably repeatable.

The jobs are the easy part. They are the part everyone has already built. The scheduler is the part where the failure modes live, and the honest version of this system is a design about at-least-once delivery and idempotent handlers.

The scale

A platform where ten thousand services each register a few recurring jobs, plus a much larger number of one-off and workflow jobs.

10,000 services x 3 recurring jobs      = 30,000 schedule entries
                                          fired every minute
                                          = ~500 jobs/sec average

largest tenants: 500 jobs per second
peak: 5,000 jobs/sec in bursts

one job = 1s to 30min of work
  most are 1-5s
  a few are hours (reports, exports, reindexing)

workflow jobs: 1,000 to 100,000 nodes
  each node is a job with dependencies
  a failed node may retry 3 times
  = up to 100,000 job executions per workflow

Three numbers shape the design.

The next-due lookup happens constantly. Every scheduler tick asks “what is due in the next minute”, over 30,000 entries. Doing that with a sorted array and a binary search is fine; doing it by scanning is not, and the scan is what most people write first.

Bursty, not steady. 5,000 jobs in a second followed by nothing. The scheduler has to absorb a backlog without falling behind, which means the due-time index has to be a structure that supports “give me everything due before T” efficiently rather than a per-tenant polling loop.

The long tail is the hard tail. A job that runs for three hours needs a deadline, a heartbeat, a cancellation, and a place to record partial progress. A job that runs for one second needs none of it, and a system designed only for the fast case will strand the slow ones.

Structure one: the timing wheel

The core operation is pop_all(due_before = now). The requirements are: find everything due, and find it in time order. There are three standard answers and they suit different scales.

Sorted set, for the general case. A balanced tree keyed by next-fire-time, with a range query for everything up to now.

  score  next_fire_time
  -----  ---------------
   1200  2024-12-10T14:00:00Z   nightly-report
   1201  2024-12-10T14:00:05Z   billing-sweep
   1202  2024-12-10T14:30:00Z   cache-warm
  ...

A range query for the first 500 items is O(log n + 500), which is exactly right. This is the default choice and it is correct for tens of thousands of entries.

Timing wheel, for high volume and coarse precision. Buckets by time quantum, with a wheel of fixed granularity.

  1-second ticks, 60 buckets (one minute)

  bucket[0]   ->  fire now
  bucket[7]   ->  fire in 7 seconds
  bucket[42]  ->  fire in 42 seconds

  advance the wheel: cursor++, pop bucket[cursor]
  on schedule: place in the correct bucket

The advantage is O(1) per insert and remove, and no comparisons at all. The cost is precision: you can only fire at bucket boundaries, and a job with a 250ms schedule rounds up to a second. That is fine for periodic jobs and wrong for anything latency-sensitive. Timing wheels are also more complex to get right, particularly on removal — a cancelled job still has an entry in a bucket, and it needs a tombstone or a version check so a cancelled job does not fire.

Heap, for the single-leader case. A binary heap keyed by next-fire-time, with a pop until the top is in the future.

  pop while heap.top <= now:
     job = pop()
     execute or dispatch(job)
     reschedule(job)

This is the simplest correct implementation for a single process and it is O(log n) per operation with no range query needed at all. If your entire schedule fits in one process’s memory — 30,000 entries is about 6 MB — a heap is genuinely the right answer, and reaching for a timing wheel at that size is complexity without benefit.

The honest ordering is: heap for a single leader, sorted set for a distributed or persistent design, timing wheel when you have enough volume that log n matters and you can accept the precision loss.

Structure two: the DAG

Workflows — where job B must run after job A — are a directed acyclic graph, and the scheduler’s job is to run nodes whose dependencies are satisfied.

      extract          transform        load
     /       \            |             |
  source    validate -----+---------- audit
     \       /                             |
      cleanup ----------------------------+

  extract -> transform -> load -> audit
  validate must finish before transform
  cleanup waits for everything

Three operations, and the choice between them is the design:

Walk the DAG at dispatch time. The workflow controller holds the graph in memory, evaluates which nodes are ready, and dispatches them.

  • Pro: the graph is in memory, evaluation is fast, and dependencies are evaluated against live state.
  • Con: a workflow with 100,000 nodes is a large in-memory object, and a controller restart means rebuilding it from persistence.

Store dependencies in the database and query them. Each node is a row with a depends_on reference; finding ready nodes is a query.

SELECT n.id
FROM   nodes n
WHERE  n.state = 'pending'
  AND  n.workflow_id = $1
  AND  NOT EXISTS (
         SELECT 1 FROM nodes d
         WHERE  d.node_id = n.depends_on
           AND  d.state <> 'succeeded'
       );
  • Pro: no in-memory graph, survives restarts, and the query is the source of truth.
  • Con: this query runs per node dispatch and needs a good index on (workflow_id, state), or a 100,000-node workflow becomes a sequential scan per node.

Precompute levels. If the graph is a DAG and you compute the topological levels, every node in level N runs in parallel once all of level N-1 has succeeded. Execution becomes “for each level, wait, then run everything in it”.

  • Pro: extremely simple, embarrassingly parallel within a level, and trivially explainable.
  • Con: a single slow node in level 2 blocks the entire level 3, even if nothing in level 3 depends on it. This is the critical path problem, and precomputed levels throw away all the available parallelism.

The choice that matters is whether you optimise for simplicity or for utilisation. Levels are simpler and waste capacity; graph-walking is more code and uses the machine better. For workflows that fan out widely with variable durations, the waste from levels is real and the graph walk pays for itself.

A fourth structure worth knowing is a lightweight state machine for workflows with a small number of states and long waits — approval workflows, onboarding, anything with a human in the loop. A state machine stored per instance is far simpler than a DAG and does not need a graph at all. Reach for it when the workflow has a shape rather than a graph.

The exactly-once question, answered honestly

There is no general exactly-once guarantee for an arbitrary side effect in another system when the job store and that system do not share an atomic commit protocol. A message can be redelivered, and a worker can crash after doing the work but before acknowledging it. Both produce the same observation — the side effect may have happened, and the scheduler cannot tell from its own state. A transactional destination, an idempotency key, or a destination-side compare-and-set can narrow that boundary; a label on the queue cannot.

The impossibility follows from the ordering requirement. If the worker must do the work and record that it did the work in one atomic step, and the two are in different systems, then a crash between them is unrecoverable without the worker being able to tell which side it got to.

So the design is two things together, and every correct implementation has both:

  at-least-once delivery
    the scheduler retries until it observes success
    a job may run 1, 2, or 5 times
    never 0 times

  idempotent handler
    running it twice has the same effect as once

Idempotency in practice looks like one of these:

  • Idempotency key. The job gets a stable key, usually workflow_id:node_id:attempt_group, and the handler’s first action is an insert into a table with a unique constraint. A duplicate insert fails, the handler returns success, and the effect happens once.
INSERT INTO job_executions (idempotency_key, started_at)
VALUES ($1, now())
ON CONFLICT (idempotency_key) DO NOTHING;
-- if rowcount = 0, someone already did this; exit successfully
  • Natural idempotency. SET x = 5 is idempotent. INSERT is not. x = x + 1 is not. DELETE FROM t WHERE id = $1 is. Most “we need idempotency” problems are solved by choosing an operation that is naturally idempotent, and this is much cheaper than a dedup table.

  • Optimistic concurrency. UPDATE t SET balance = balance - 100 WHERE id = $1 AND balance >= 100. Run it twice and the second one affects zero rows, because the balance is no longer sufficient. The check is in the statement, and no coordination is needed.

  • Set-based operation. A job that recomputes an aggregate from source data is naturally idempotent if it overwrites rather than increments.

The guarantee you can make locally is one durable row per scheduled occurrence. Two schedulers may both decide that job 4711 is due; they are allowed to race as long as the occurrence key makes the write idempotent:

  leader holds the scheduler lease
      |
      v
  for each schedule due before now:
      INSERT INTO job_due (schedule_id, fire_time, claimed_by)
      VALUES ($schedule, $time, $this_leader)
      -- unique on (schedule_id, fire_time)
      -- a second attempt hits the constraint and is skipped

The unique constraint on (schedule_id, fire_time) makes occurrence materialisation idempotent: one row survives no matter how many schedulers attempt it. Be explicit about time zones, daylight-saving overlap, and schedule revisions; sometimes the durable key should include a schedule-version or generated occurrence ID rather than a wall-clock timestamp. Execution is then normally at-least-once and handled by the worker’s side-effect contract.

Leader election, and why the lease is not optional

Only one scheduler should be creating jobs. Multiple schedulers are fine — they are stateless workers — but the tick that inserts into job_due must be single-writer, or you get duplicates.

  +---------------------------+
  |  scheduler (n replicas)   |
  |                           |
  |  A: leader                |  holds a 10s lease
  |     - claims due jobs     |  renews every 3s
  |     - inserts job_due     |
  |                           |
  |  B, C: standby            |  watch the lease key
  |     - on lease loss,      |  one of them takes over
  |       contend to acquire  |
  |     - become leader       |
  +---------------------------+

  guarantees:
    - at most one leader at a time (quorum)
    - a new leader within the lease TTL
    - an old leader that lost the lease must stop
      (fence yourself, exit on renewal failure)

That last point is the same self-fencing discipline as a heartbeat. A scheduler that loses its lease must stop claiming jobs immediately, because the new leader is already claiming them, and a stale leader with a valid-looking lease is a duplicate-job bug. The failure to guard against is the partition where the old leader still believes it holds the lease — which is why lease renewal failure has to be fatal rather than a warning.

The window during which a job can be late is the lease TTL plus a tick interval. With a 10-second lease and a 1-second tick, worst case is roughly 11 seconds of delay. If your schedule has finer granularity than that, the lease is too long, and the trade is availability of the scheduler against schedule precision.

Failure stories worth testing

Kill the leader mid-tick

A new leader takes over within the lease TTL. Confirm no jobs are lost and no jobs are claimed twice, and measure the actual takeover time.

Partition the leader from the quorum

The old leader must stop claiming jobs when it cannot renew. If it keeps going, you get two schedulers and duplicate jobs, and the unique constraint on (schedule_id, fire_time) is what catches it.

Kill a worker after the side effect, before the acknowledgement

The job is redelivered and runs twice. This is the at-least-once test and the idempotency handler is what makes it correct. If this test shows duplicate effects, the design is not done.

Set the clock back five minutes on the scheduler

Jobs scheduled in that window fire immediately and en masse. Then set it forward five minutes: jobs fire late, in a batch. The clock is the thing the design distrusts, and a monotonic clock for intervals plus wall-clock for scheduling is the mitigation.

Burst 5,000 jobs due at the same instant

Confirm the queue absorbs the backlog and the next-due lookup is still fast. This is where a scan-based scheduler dies.

Cancel a job that is already in a timing wheel bucket

It must not fire. This is a tombstone or a version check, and it is a bug that appears only under cancellation, which is exactly when you do not want a bug.

Fail one node of a 100,000-node workflow

Confirm the workflow retries the node, that retrying does not re-run the successful side effects of the other 99,999, and that the DAG evaluation is still fast.

Put a 30-minute job in the same queue as a 1-second job

Confirm the 1-second jobs are not stuck behind it. Per-job queues and deadlines, or the scheduler is head-of-line blocked.

Complete a workflow’s final node and lose the completion write

Confirm the workflow is not re-run from the start. This is a distributed transaction between the job result and the workflow state, and it needs to be resumable.

Run two scheduler instances against the same database with no leader election

Confirm the unique constraint catches every duplicate. This should fail; if it passes, the protection is not what you think it is.

A production-ready architecture

  +---------------------------+
  |  API                      |  register, pause, resume,
  |  register schedule        |  trigger now, query history
  |  Cron / natural language  |
  +------------+--------------+
               |
               v
  +---------------------------+
  |  scheduler leader         |  lease-based, 10s TTL
  |  every 1s:                |  renew every 3s
  |    pop due from heap      |  exit on renewal failure
  |    INSERT job_due         |  unique (schedule_id, fire_time)
  |    (conditional)          |  -> exactly-once claim
  +------------+--------------+
               |
               v
  +---------------------------+
  |  job_due table            |  the durable hand-off
  |  (schedule_id, fire_time) |  scheduler and workers are
  |  state, attempts, lease   |  decoupled through this
  +------------+--------------+
               |  claim: UPDATE ... SET worker = me
               |         WHERE id = $1 AND (lease IS NULL
               |                OR lease < now())
               v
  +---------------------------+
  |  workers (n replicas)     |  stateless, any worker
  |  - per-job-class queue    |  takes any due job
  |  - per-job concurrency    |  limits, not a global pool
  |  - deadline + heartbeat   |  long jobs report progress
  |  - idempotency key        |  INSERT first, exit on conflict
  |  - retry with backoff     |  attempts tracked in job_due
  +------------+--------------+
               |
               v
  +---------------------------+
  |  workflow controller      |  if the job is a DAG node
  |  - DAG in DB or memory    |  evaluate ready nodes
  |  - dependency state       |  dispatch what is unblocked
  |  - resume from any state  |  crash-safe
  +---------------------------+

  delivery guarantee:
    at-least-once execution
    exactly-once claim
    idempotent handlers
    -> effectively-once effect

A sensible delivery checklist:

  1. Make occurrence materialisation idempotent with a unique occurrence key. Do not rely on leader election alone; it reduces contention, while the constraint is the correctness boundary.
  2. Assume duplicate execution and require an idempotency key on every handler. The only acceptable proof is a test that kills the worker after the side effect.
  3. On lease renewal failure, exit. Do not log a warning and continue; a scheduler that cannot prove it is the leader must not claim jobs.
  4. Size the lease against your schedule precision and know the resulting worst-case lateness.
  5. Use a heap if one leader holds the whole schedule, a sorted set for persistence or distribution, a timing wheel only at a volume where log n is measurable.
  6. Isolate jobs: per-class queues, per-job concurrency limits, and deadlines. A long job must never block a short one.
  7. Require long jobs to heartbeat progress and make that visible, so a stuck three-hour job is distinguishable from a slow one.
  8. Make workflow execution resumable from any partial state, and test resume after every node type.
  9. Use graph walking over precomputed levels for wide, variable-duration workflows; use levels for simple ones and accept the waste.
  10. Alert on schedule lag, not just on job failure. A scheduler that is running fine and falling behind is the failure you want to catch.
  11. Cap retries per job and give up explicitly, sending to a dead-letter path that a human can see. Infinite retries hide bugs.
  12. Treat the clock as untrusted. Monotonic time for intervals, wall clock only for the schedule itself, and a test that moves the clock in both directions.

Common mistakes

Mistake What actually happens Better decision
Claiming exactly-once execution Duplicates happen and are discovered in production At-least-once plus idempotent handlers
Leader election as the only duplicate protection A partition or a race produces duplicate jobs Unique constraint on (schedule_id, fire_time)
Non-idempotent handlers Retries double-charge, double-send, double-insert Idempotency key, checked first
Leader continues after lease loss Two schedulers claim the same work Renewal failure is fatal, exit immediately
Scanning the schedule every tick O(n) per tick, falls over at scale Heap, sorted set, or timing wheel
Timing wheel with no tombstone Cancelled jobs still fire Version check or tombstone on remove
One global worker pool A 30-minute job blocks a 1-second job Per-class queues, per-job concurrency
No per-job deadline A hung job runs forever Deadline plus progress heartbeat
Infinite retries Bugs are hidden, queues never drain Attempt cap, dead-letter, visible to humans
DAG in memory only A controller restart loses the workflow Persist it, or rebuild deterministically
Precomputed levels for a wide workflow One slow node blocks an entire level Graph walk for variable durations
No cycle detection A bad workflow definition hangs forever Validate acyclicity at registration
Monitoring job failures, not schedule lag A falling-behind scheduler looks healthy Alert on lag and backlog depth
Trusting the wall clock for intervals NTP adjustment fires or delays the whole schedule Monotonic clock for intervals
Relying on the queue for exactly-once A redelivered message runs a non-idempotent handler Idempotency at the handler, always
No per-tenant limits One tenant’s jobs starve everyone else Per-tenant quotas and concurrency
No manual trigger path Re-running a failed job requires a code change First-class trigger-now with the same path
Comparing only start and end times A job that ran for 30 minutes looks identical to one that took 30s of work Record attempts, duration, and exit reason
Storing the schedule only in memory A restart loses every registration Persist; memory is a cache of the DB

The complete story in one minute

A distributed scheduler is a cron table with two data structures in it. The first answers “what is due before T” and is a heap for a single leader, a sorted set when the schedule is distributed or persisted, or a timing wheel at a volume where log n is measurable and you can accept second-level precision. The second answers “what is unblocked” and is a DAG, walked at dispatch time if you want utilisation and precomputed into levels if you want simplicity — with levels wasting real parallelism whenever one node in a level is slow.

Exactly-once execution is not a general property across two systems that do not share a transaction. A worker can crash after the side effect and before the acknowledgement, leaving the scheduler unable to distinguish “committed” from “not attempted.” The honest default is at-least-once delivery plus a destination-enforced idempotency key, a naturally idempotent operation, or an optimistic-concurrency predicate. If a destination can join the same atomic commit, say so specifically instead of exporting that claim to every side effect.

For scheduled creation, insert with a unique occurrence key so concurrent schedulers converge on one durable row. Leader election makes the race cheaper, but the constraint is the correctness argument — a partitioned leader that keeps trying is caught by the same key, not by the lease.

Isolation is per-job, never global. Per-class queues, per-job concurrency limits, deadlines, and progress heartbeats, so a three-hour report cannot sit in front of a one-second cache warm. And alert on schedule lag rather than job failure, because a scheduler that is working perfectly and falling behind is the one you want to catch before it is a missed contract.

heap: pop all due, claim with unique insert
workers: any worker takes any due job, idempotency key first
delivery: at-least-once execution, exactly-once claim
isolation: per-job queues, limits, deadlines

The hard part was never parsing a cron expression. It was making sure that a worker dying between doing the work and saying it did the work could not produce a second effect.

Technical references

Keep reading
Browse everything