Ch.13: System Design: Distributed Job Scheduler
Outline
- 0:00 Introduction
- 0:29 The Architecture Map
- 1:20 Why Cron Plus a Queue Breaks
- 2:06 Functional Requirements
- 2:44 Non-Functional Requirements
- 3:29 Core Job Record
- 4:19 Scheduler Loop
- 5:13 Partitioning Due Work
- 5:51 Leases and Heartbeats
- 6:35 Crash Recovery and Retries
- 7:18 The Exactly-Once Trap
- 8:02 Recurring Jobs
- 9:08 Priority and Starvation
- 10:00 Worker Pools and Routing
- 10:57 Observability
- 11:46 Architecture Payoff
- 12:43 Why the Design Holds
- 13:14 Closing
Transcript
0:00 A distributed job scheduler sounds like cron with more machines. Take a billing job scheduled for midnight: it charges a customer, then the worker dies before recording success. Wait, so the retry itself becomes the bug, not just the fix attempt? Right. For example, the scheduler sees an unfinished job, retries it, and now the recovery path can double-charge the customer. Yes. That's the trap. The key is separation: the scheduler owns durable intent, and workers make side effects replayable or ignorable.
0:27 That sets up the whole scheduler problem. Now start with the finished map. A client writes a job into durable storage. Scheduler shards scan due-time buckets. A claim layer takes a lease. Workers run the payload. Heartbeats keep the lease alive, and timeouts return abandoned work. That looks like a calendar swallowed a queue. A calendar that has to remember to actually do the thing, not just show you the date. The shape is time on one side, execution on the other. The map also has recurring-job expansion, priority lanes, worker pools, retry policy, dead-letter storage, and observability wrapped around the whole thing.
1:04 It really does. And we come back to this map after each box has earned its place. The diagram is not a shopping list. The contract for the rest of the map is simple: no lost work, no duplicate work, no invisible backlog. Every box is there to answer one of those promises. From that map, the tempting shortcut is cron plus a queue. It works for 5 jobs and one worker, which is why it is so seductive. My instinct with reminders is to trust that shape, until cancellation turns out to be, More of a polite suggestion than a guarantee.
1:39 Once two scheduler processes can see the same due row, both can enqueue it unless the claim itself is atomic. Waking up every minute is the easy part. What actually bites you is ownership: if two schedulers both say yes, the work is already duplicated. That is where the nightmare starts: two machines, one timestamp, and no real owner. The first design rule is atomic claim before enqueue. That creates the clean line: claim first, enqueue second. Okay, let's pin down what this thing has to do before we touch architecture.
2:10 The system needs one-time delayed jobs, recurring jobs, cancellation, retries, priorities, audit history, and a way to route different jobs to different workers. Cancellation too? Is that really scheduler logic? A cancel has to race with claim and execution, not just mark a row canceled. Right. Take an email campaign: a user cancels it while a worker already has it. That is where the demo breaks. From support's view, audit history gives the answer a week later. The product promise is delayed intent that stays durable, visible, cancellable, and retryable.
2:41 The operator finally has something concrete to promise. Now the non-functional requirements are where the scheduler becomes a distributed system. It has to handle millions of scheduled jobs, survive node loss, avoid duplicate claims, keep schedule lag low, isolate tenants, and recover without manual cleanup. The precision question is sneaky: if every scheduler host has a slightly different clock, noon is not exactly the same moment everywhere. The stored due time is the source of truth, but each scanner still, you know, compares that value against its local clock.
3:14 Noon is a range, not a point, and the design needs clock-skew tolerance plus a lag metric that tells operators what actually happened. Two scanners can genuinely disagree about whether it's noon yet. The metric tells operators how late work is and where. That leads to a core data model that is small but, I mean, loaded. It starts with identity, tenant, payload reference, due time, priority, and status. Then recovery fields show up: attempt count and max attempts for retries, lease owner and lease expiry for stale ownership, plus a request key and audit timestamps.
3:50 The temptation is to store only the schedule timestamp and payload because that is all the first demo needs. And that demo is a trap, because the row has to explain recovery, retries, and outside side effects. Those fields are not decoration. They identify the current worker, stale ownership, and which outside request a retry belongs to. The job row carries the ledger of what the system believes about this unit of intent. That gives recovery a durable reference point. The scheduler loop is the part you have to get exactly right.
4:23 Scan for due jobs in your shard, try to claim a small batch atomically, publish claimed jobs to execution, and move on. Right. Walk me through one claim from the storage row, because that is where ownership becomes real. What actually changes? Let's say one due row is visible to two scheduler hosts. The row moves from scheduled to leased, receives an owner plus an expiry, and increments the attempt count. That update must be conditional: only claim it if the job is still schedulable right now. Oh, that's the whole game right there.
4:55 The claim is not a read followed by a write; it is one compare-and-set style move. It is the lock on the shared sink, not two roommates both calling dibs. Claiming is where duplicates are born or prevented. The update is the real guard; the scan is just how the job got found. Once claims are safe, one scheduler still cannot scan every due job. So the table is partitioned by time bucket and shard. Time bucket finds when; shard spreads who. So not one giant queue ordered by due time? Not at scale.
5:28 You might have a bucket for this minute, then a hash shard inside it. Scheduler process 12 owns bucket now plus shard 12, and another process owns a different shard. That dullness is intentional. You want predictable ownership ranges so failover can reassign them, not every scheduler politely scanning the same universe. Then teams can rebalance ownership before one scan turns into a global stampede. Once a worker receives the job, it owns a lease, not the job forever. The lease says, I am working on this until this expiry time.
6:00 And completion is a separate write. Lease is not completion. A long job heartbeats to extend the lease, like renewing a checked-out book before someone else can claim it. If the worker stops heartbeating, the scheduler can treat the lease as expired and make the job eligible again. Started is just the library receipt; completion is actually returning the book. That prevents the most common mental-model bug: thinking a popped job is safe just because someone started it. A start only proves someone touched it.
6:30 Crash recovery mostly asks what actually finished. That gives retry evidence. Now picture this: the worker claimed a job, ran for 30 seconds, and disappeared. No completion write. No heartbeat. Eventually the lease expires. I got burned by this exact shape with a payment job. The charge succeeded, the success write failed, and the retry looked completely reasonable from the scheduler's point of view. The scar is simple: the scheduler can only see its own state. It cannot magically know whether the external side effect happened.
7:01 I mean, that is why retries heal lost attempts, but they can also repeat a side effect unless the handler treats duplicate requests as the same request. Retry is the recovery tool. Duplicate safety is the damage control. Production needs both, which is exactly what finance is about to complain about. And then finance or customer support pushes back. They say, "Can you just make sure this job runs exactly once? We cannot charge people twice." They are right. That is the correct business requirement and the dangerous technical sentence. True. The scheduler can record durable intent, one active lease at a time, bounded retries, and an idempotency token. It cannot alone guarantee that an external system observed the side effect exactly once.
7:44 Right. So exactly once moves to the boundary: the payment provider, email service, or database write needs to accept that operation token and collapse duplicates. Correct. The scheduler makes duplicates detectable and safe to retry. The handler and external system make repeated attempts harmless. The closing rule is scheduler remembers; the handler owns the outside world. After the exactly-once boundary, recurring jobs look like strings, but at scale they are generators of future state. Every cron rule has to turn into concrete future runs.
8:14 Do we immediately create every single future run upfront for every tenant and every time zone, before the schedule window even arrives? Definitely not. That would turn a calendar rule into a storage bill. Most systems, I mean, expand a rolling window: generate the next few occurrences, run them through the same job table, and extend the window as time moves. And then time zones show up, because of course they do. Oh, man. Daylight saving time, skipped hours, duplicated hours, tenant-local calendars. For example, a recurring scheduler has to decide what "daily at 2:30" means when 2:30 disappears or when 1:30 happens twice in one night.
8:52 Oh, right. That calendar bug sleeps for 6 months. Yep. Calendar policy is the work hiding inside the cron string. Local-time weirdness needs one explicit answer per tenant, region, and daylight-saving policy. The calendar decision is explicit instead of accidental. Next, priority is where time stops being the only ordering dimension. Priority lanes decide which due job runs first when capacity is limited. Password reset email beats weekly digest. Payroll beats delete old cache files. That ordering is the user-facing guarantee.
9:24 Consider a low-priority cleanup job. How does it ever win if high-priority traffic keeps arriving? Usually the scheduler needs aging, fair-share quotas, or reserved worker capacity. Maybe base priority plus wait time multiplied by an aging factor. Urgent work goes first, but low-priority work still makes progress. The aging trick is the part I always forget exists. This is one of those designs where urgent is easy and fair is hard. Yeah, teams almost always skip that. The dashboard should show both: how fast urgent jobs clear and whether background work is stuck.
9:57 That turns background cleanup into visible work. Now routing decides which workers are even allowed to run a job. For example, some jobs are compute-heavy. Some are input-output-heavy. Some need a region, a tenant isolation boundary, or a machine with a specific dependency. Right. The scheduler is not just picking time. It is dispatching to the right pool. Sure. The job record carries route labels, and worker pools advertise what they can run. The scheduler matches the two. Hold on. What happens when a pool is saturated but another pool is empty?
10:31 Right. You only borrow capacity if the labels allow it. Otherwise you create a correctness bug while trying to fix a utilization graph. Right. The utilization graph is the seduction, and the labels are the brake. Yeah. That gives us the model: dispatch plus memory. The job remembers what it needs; the pool advertises what it has. That is the constraint I would want in the design review. Once jobs route correctly, observability is where the design stops being clever and starts being operable. A scheduler can fail silently, which is what makes it unpleasant.
11:06 No one pages because "the scheduler is down." They page because invoices are late, reminders did not send, or retries exploded. What do we measure before the customer notices? For example, start with schedule lag, claim latency, queue depth, and age of the oldest due job. Then track retry rate, dead-letter count, lease expiries, and worker saturation. The maximum age is the smoke alarm. Oh, wow. Average lag can look fine while one shard is 12 hours behind. The oldest item tells you which promise has already been broken.
11:38 Observability has to show that broken promise before the runbook starts. Yeah. You catch it before the customer's email does. Now every box on the map has a job. The durable job store holds intent, status, attempts, leases, priorities, routing labels, explicit recovery hints, and idempotency tokens. The map is doing real work now. Scheduler shards own time buckets and hash shards; they scan due work, claim atomically, and publish leased jobs to the right worker pool. Workers heartbeat while they run, complete explicitly when they finish, and let expired leases return to schedulable state when they disappear.
12:14 Recurring expansion generates future runs, retry policy decides when attempts come back, the exhausted-job area keeps failed work inspectable for humans, and observability watches lag, depth, retries, stuck leases, and queue age by lane. It's a lot of boxes, but each one protects a promise: keep work from being found too late, duplicated blindly, starved, or invisible. The promise check is simple: timing, ownership, fairness, and recovery each have a box. With the map assembled, the design has one job: keep the promises separate.
12:47 That separation is the whole trick. If you blur any two, the system lies to you. Blur lease and completion, and started becomes done. Blur retry and idempotency, and recovery creates duplicate side effects. Blur priority and fairness, and background work vanishes. It is kind of elegant, honestly. The scheduler is not trying to be magic. It is keeping enough state that every failure has a next move. The architecture keeps durable intent separate from messy outside effects. Finally, this is not cron scaled sideways.
13:16 It is durable intent plus safe ownership: claim one job, lease the work, retry without turning recovery into a second side effect, and measure the promises users feel. Next video, we move from delayed execution to content moderation: fast decisions, uncertain model confidence, human review queues, appeals, and the cost of being wrong in either direction. That is a very different kind of hard. Fast uncertainty instead of delayed certainty. That is it for this one. Durable intent, safe ownership, and measured promises are the model.
13:48 Thanks for listening to Learning Podcasts.