System-design interview · Extended interviews
Design a distributed job scheduler
Turn recurring rules into durable run records, retry failed attempts and accept results only from the current worker, with explicit timing, overlap and side-effect limits.
You will learn to
- Distinguish a recurring schedule, one intended occurrence and its execution attempts.
- Trace atomic due-run creation, worker claim and fenced completion.
- Explain calendar rules, missed-run policy and capacity limits before promising start deadlines.
Practice in this chapter
8 interview questions with model answers and follow-ups.
Go to interview practiceUseful foundations: Message queues, event logs, delivery guarantees, and backpressure · Databases, data models, and ACID transactions · Quorums, consensus, leases, and fencing
Workload and timing examples are interview assumptions.
Dotted concept links open the relevant explanation in a new tab.
01Define the run the customer expects
Design a scheduler for approved report jobs with versioned parameters. A customer asks for a sales summary every weekday at 09:00 in America/New_York. A schedule is that recurring rule. An occurrence is one intended Tuesday run at a resolved instant. An attempt is a worker's try at that occurrence. A retry creates another attempt, not another Tuesday report.
Support create/edit/pause schedules, run once, inspect history, cancel and retry according to policy. Exclude arbitrary untrusted code and multi-step workflow graphs initially. Choose jobs returning a small structured result that can be stored transactionally with run status; large report objects are a follow-up with their own publication rules.
Promise durable tracking and one accepted result per occurrence, not that a process executes physically once. A paused worker may resume after replacement. Any external email or payment effect needs a destination-supported identity or reconciliation beyond the scheduler's local result guarantee.
Target 99% of ordinary admitted runs starting within ten seconds of their due time. Require accepted run state to survive one zone failure under synchronous durable database replication and safe failover. Starting is different from finishing: a three-hour report cannot inherit a ten-second completion target. During unsafe authority loss, delay claims and completion rather than accept competing owners.
Clarify whether jobs are approved handlers or arbitrary code, whether the deadline means start or finish, and how overlapping or missed calendar runs should behave. Choose approved bounded-result reports, a start-time objective, no logical overlap and latest-only catch-up for the running example.
02Functional requirements
Manage schedules. Create, edit and pause versioned calendar or interval schedules with named time zones and an explicit missed-run policy.
Run approved jobs. Save a record for each intended occurrence, dispatch it to a compatible worker and accept its bounded structured result.
Inspect and control execution. Provide run history and status, manual runs, supported cancellation and retries of the same intended occurrence.
Recover unfinished work. Resume dispatch after crashes, replace failed attempts and deliver only the accepted result through the committed notification path.
03Non-functional requirements
These are illustrative interview assumptions, not product facts or measured benchmarks. Confirm them before choosing components, then validate the completed design under the stated workload. Here p95 means the 95th-percentile latency: 95% of measured requests take no longer than that value. Report errors and rejected work alongside latency; a fast failure is not a successful outcome.
Workload and resources. Plan for ten million daily occurrences and a 50,000-job 09:00 burst. The five-second job and 10,000-slot example is deliberately insufficient for the target; provision measured execution capacity or renegotiate admission/jitter.
Start-time objective. Target starting 99% of ordinary admitted occurrences within ten seconds of their due time under the agreed runtime/resource mix. Misses remain visible; this does not promise job completion within ten seconds or every physical retry meeting the original deadline.
Execution correctness. Create one logical occurrence per scheduled instant and accept at most one result. Enforce the selected no-overlap policy across occurrences; retries may physically execute more than once.
Durability and safe failure. Preserve accepted schedule/run state across one zone failure using synchronous durable database replication and safe failover. Without safe authority, delay claims and completion instead of accepting competing owners.
Security and effect boundary. Authorize schedule ownership and approved worker types, constrain resources and secrets, and require destination-supported identity or reconciliation for external effects. A fenced local result does not prove an arbitrary external call happened once.
04Build a complete scanner, database and worker flow
Start with a schedule API, one relational database, a scanner polling an indexed due-time column and a bounded worker pool. At Tuesday 09:00, the scanner locks schedule S7, verifies its current revision and due instant, inserts occurrence R7, creates durable dispatch work and advances the next due time in one transaction. A crash before commit leaves it due; a crash after commit leaves R7 discoverable.
A worker claims R7 through a short transaction. The database marks it running, records attempt A1, grants a deadline and an increasing ownership token. The worker releases locks before computing the report. Heartbeats may extend the deadline only while A1 is still the current owner.
On completion, another transaction checks the attempt and token against current authority, stores the bounded result, marks R7 succeeded and saves any notification work. A repeated completion with the same accepted identity returns the saved outcome. A stale worker cannot replace that result.
The first dispatcher can poll durable work directly; no broker is necessary to make the system recoverable. A queue becomes useful later for worker isolation and throughput, while the run database remains the authority about which attempt may commit.
The database owns occurrences and completion; the worker performs computation outside locks.
Read each connection in order
- syncCreate / edit / inspectSchedule owner → Schedule API
- syncVersioned schedule and request resultSchedule API → Schedules, runs, attempts, outbox
- syncAtomic occurrence + next due timeDue scanner → Schedules, runs, attempts, outbox
- syncClaim token and heartbeatBounded report workers → Schedules, runs, attempts, outbox
- syncGuarded result completionBounded report workers → Schedules, runs, attempts, outbox
05Choose calendar, interval and outage behavior explicitly
“Every hour” can mean the top of each clock hour, sixty minutes from a fixed anchor, or sixty minutes after the prior run completes. These are different scheduling contracts.
| Rule | How the next occurrence is determined |
|---|---|
| Calendar | Match a wall-clock rule in a named time zone, such as weekdays at 09:00. |
| Fixed rate | Add an interval to an anchor, independently of job completion. |
| Fixed delay | Wait the interval after the previous run completes. |
Store the named time zone, rule revision and resolved execution instant. Daylight-saving changes can skip or repeat a local time; define whether to skip, shift or run one/both occurrences. Preview several future instants so users can see what the rule means.
For this sales summary, choose no logical overlap and latest-only catch-up after a long outage. Older missed occurrences are explicitly recorded as skipped under that policy, rather than launching every missed day simultaneously. Other products may need catch-all or skip-all, but those choices affect capacity and business meaning.
A schedule edit affects future occurrences not yet created. An already materialized R7 keeps its original parameter revision through retries. Otherwise the same logical run could produce different business work just because someone edited tomorrow's schedule while today's worker was recovering.
06Prove whether the morning start target is feasible
Ten million daily occurrences average about 116 starts/s, but calendar schedules synchronize work. Suppose 50,000 jobs become due at 09:00, each takes five seconds and the fleet has 10,000 execution slots. Under ideal conditions, starts occur in five waves at seconds 0, 5, 10, 15 and 20. The last jobs miss a ten-second start target even before scanner delay or worker startup.
At least 16,667 slots would allow three ideal waves at 0, 5 and 10 seconds. Real duration variance, setup overhead and failure reserve require more margin. Alternatively negotiate start-time jitter, reserve capacity for strict schedules or reject an impossible deadline commitment. A queue stores the burst; it does not create compute capacity.
At 2 KB per run, ten million runs produce 20 GB/day and 600 GB over thirty days before indexes and replicas. Ten thousand active attempts heartbeating every five seconds produce 2,000 state updates/s in addition to claims and completions. Keep those transactions small.
Worker resources often dominate scheduler metadata. Ten thousand 256 MB jobs need roughly 2.56 TB of fleet memory. Measure resource classes and runtime tails, not only job counts. A hundred hour-long jobs are a different queue from a hundred one-second summaries.
07Keep recurrence, occurrence and attempt records distinct
Interfaces
| Request or message | Contract |
|---|---|
POST /schedules with request key |
Creates one schedule and returns revision and next resolved instant. |
PUT /schedules/S7 with expectedRevision and rule changes |
Changes future rule interpretation without overwriting a concurrent edit. |
POST /schedules/S7/runs with key |
Creates or recovers one manual occurrence. |
GET /runs/R7 |
Returns authoritative status, due/start times, attempt and accepted result. |
Inspect the worked occurrence
GET /runs/R7
Schedule, occurrence and attempt in the example
schedule: S7 — sales summary, weekdays at 09:00, America/New_York
occurrence: R7 — one resolved Tuesday run with frozen parameters
attempt: A1 — the worker’s current try at R7
Schedule update request fields
expectedRevision: the current schedule revision
rule changes: the intended changes for future, uncreated occurrences
Stored records
| Record | Fields or identity | Purpose |
|---|---|---|
| Schedule | Schedule ID | Rule, zone, revision, next due time and active occurrence for no-overlap policy. |
| Occurrence unique by schedule/resolved instant | Schedule and resolved instant (unique) | One logical run with frozen parameters and current attempt token. |
| Attempt | Occurrence and attempt identity | Worker, number, lease deadline and outcome. |
| Outbox | Dispatch or notification work identity | Durable dispatch or post-completion notification work. |
Keep a schedule and its occurrences where they can commit together in one database transaction. A due index finds bounded batches without scanning every rule. History uses scheduled time plus run ID with a stable page boundary. Authenticate schedule ownership and worker credentials; callers cannot choose their own ownership tokens.
Create and run-once keys include a payload fingerprint. Matching retries recover the original resource, while changed payload reuse conflicts. Scope request identities so their saved result and resource creation share one transaction. Do not put a global deduplication table on unrelated shards and then assume it atomically protects every schedule write.
08Explain leases and fencing through a paused worker
A lease is time-limited permission to work. It lets the scheduler eventually replace a worker that stops heartbeating, but its expiry does not physically stop that process. Fencing is the enforcement step: the protected store rejects an outdated ownership token when accepting a result.
The stale-worker example proceeds as follows:
A claims R7 with token 41, then pauses long enough for its deadline to expire.
B claims a replacement attempt with token 42 and completes.
When A resumes, its completion still carries 41. The database sees that 42 is current and refuses A’s update.
Checking “my lease was valid when I started” would not prevent A overwriting B’s result.
Claims, heartbeats and completion lock the run record and check the authoritative database clock after waiting for locks. Completion verifies current attempt, token, unexpired ownership and allowed state in the same transaction that stores the result. If A completes before a replacement is authorized, the run becomes terminal and B cannot claim it. If replacement wins first, A cannot complete. Those are the two permitted outcomes.
For no-overlap, also keep an active-occurrence guard on S7. It remains attached to R7 during retries, preventing Wednesday's distinct occurrence from starting logically while Tuesday is unfinished. This does not prove that an expired Tuesday process physically stopped; fenced results and protected external effects still matter.
Token enforcement occurs in the same database transaction as result acceptance.
Read each connection in order
- syncClaim R7; receive token 41Worker A → Run authority
- syncLease 41 expiresRun authority → Run authority
- syncClaim replacement token 42Worker B → Run authority
- syncComplete with 42; accept resultWorker B → Run authority
- syncResume and complete with 41Worker A → Run authority
- returnReject stale attemptRun authority → Worker A
09Keep external actions behind the accepted result
The chosen report job computes a bounded result and returns it to the scheduler. It does not email its private attempt output directly. The successful completion transaction stores the canonical result and a unique delivery intent such as send-R7. A notification adapter then processes that committed intent with the accepted content.
This arrangement prevents stale A from publishing a different report through the normal notification path after B's result wins. The provider may still accept email and lose its reply; the adapter must recover that uncertain result. Reuse the logical action identity under the provider's supported idempotency contract or reconcile it. Without those capabilities, disclose the category's duplicate-versus-missing-delivery policy.
Arbitrary jobs that write external databases need destination cooperation. A destination can enforce the scheduler's ownership token with its protected write, or it can recognize an appropriate stable logical operation identity. A worker can pause after checking its token, then resume the remote call after replacement. The local check therefore cannot guarantee one external effect.
Cancellation also has a boundary. Record cancellation requested and let workers stop cooperatively. Once the authority accepts a canceled terminal state, reject later completion and dispatch no retries. Already completed external effects cannot necessarily be recalled, and killing a process does not undo them.
10Distribute scanning and execution for measured bottlenecks
Separate scanners from workers when execution load delays finding due work. An outbox relay publishes run IDs into ready queues. Repeated queue messages remain hints for the same R7; workers must obtain the current run claim before working. Recovery also scans expired attempts, so correctness does not depend entirely on one unacknowledged broker message surviving forever.
Partition schedules into stable buckets and distribute due-index scans. Scanner leases can reduce redundant effort, but unique occurrence creation and the materialization transaction remain the backstop when two scanners overlap during reassignment. Do not dispatch first and update nextRunAt later: that can create duplicate intended runs after a crash.
Use separate resource pools and tenant active-job limits for short, long and memory-heavy work. Weighted fairness prevents one large customer from occupying all slots; the cost is scheduling complexity and sometimes idle reserved capacity. A single first-in-first-out queue is adequate when jobs are homogeneous and that fairness policy is desired.
Prewarm workers for known peaks and load only a short future horizon into timer structures if database polling becomes expensive. In-memory timers speed up dispatch and can be rebuilt from the saved schedules. Before materializing work, recheck the current rule revision so an obsolete timer cannot execute a canceled or edited occurrence.
Scanners create each due occurrence and its dispatch work in one database transaction. The relay publishes run IDs to resource-class queues; workers claim current authority before computing. Results and terminal status commit under the same ownership check, while API reads expose the durable run state.
Read each connection in order
- syncRules and status requestsSchedule caller → Schedule and status API
- syncPersist rules / read runsSchedule and status API → Schedule / run / result DB
- syncMaterialize due occurrencesDue-bucket scanners → Schedule / run / result DB
- syncRead dispatch outboxDispatch outbox relay → Schedule / run / result DB
- asyncPublish stable run IDsDispatch outbox relay → Resource-class queues
- asyncWake a run attemptResource-class queues → Bounded worker pools
- syncClaim / fenced completionBounded worker pools → Schedule / run / result DB
11Recover work without silently changing its meaning
| Failure | Required recovery |
|---|---|
| Scanner dies before materialization commit | The due instant remains eligible for a later scan. |
| Scanner dies after commit but before dispatch | Durable outgoing work republishes the same occurrence. |
| Worker fails before completion | Expire its attempt, back off and retry the same frozen occurrence within limits. |
| Completion response is lost | Query or repeat the same completion identity; return the saved result. |
| Database loses safe authority | Delay new claims and accepted results until authority recovers. |
Classify permanent errors, such as an invalid report query, separately from transient failures. Limit attempts and total retry age, retain the reason and expose exhausted runs for review. A poisoned job must not consume the fleet indefinitely.
Apply outage catch-up policy in bounded pages. A latest-only schedule records skipped older work and creates the appropriate current occurrence; it does not rewrite all due timestamps to now and hide the missed history. A deliberate backfill uses a new action ID because it requests an additional run.
Preserve schedule guards, request keys, attempt tokens and outgoing work in backup/restore. Restoring only recurring rules can recreate already executed occurrences. Reconcile uncertain external actions before blindly resuming them after a disaster outside the promised durability boundary.
12Measure deadline behavior and state the next limits
Monitor due-to-start delay by tenant and resource class, runtime tails, active slots, oldest accepted due time, lease expiry, stale-completion rejection, skipped occurrences and permanent failure rates. Queue length alone cannot estimate delay without knowing execution durations and resource needs.
Test two scanners reading the same due instant, a worker pausing past its lease, completion racing cancellation and no-overlap across two different occurrences. Simulate a two-day outage and verify the chosen catch-up policy. Shadow-test recurrence parser changes across time zones and daylight-saving transitions before enabling them for existing rules.
Approved job types still need resource limits, restricted secrets and tenant isolation. Allowing arbitrary user code would make sandboxing, network policy and metering central requirements, not a small follow-up flag. Large report objects similarly need verified immutable publication and retention; the bounded database result keeps the first design's guarantee concrete.
The remaining single-schedule limit is intentional: a no-overlap job cannot gain logical parallelism merely by adding workers. Split its business work only if independent partitions and a merge step are valid. Dependencies, human approvals and month-long execution are signals to introduce a durable workflow model rather than stretching a cron rule into one.
13Check the design against its requirements
Use the numbered requirements to check the final design. FR refers to the functional list; NFR refers to the non-functional list. Performance rows specify tests still required, not achieved benchmark results.
| Requirement | Design mechanism | Verification and remaining limit |
|---|---|---|
| FR1; NFR3: intended calendar work | Versioned rules, resolved instants, unique occurrences and a schedule-level active-run guard. | Test daylight-saving boundaries, concurrent scanners, edits and two missed days. Verify the selected no-overlap/latest-only policy rather than silently inventing extra runs. |
| FR2–3; NFR1–2: timely execution | Due indexes, separate resource pools, admission and prewarmed workers. | Simulate the 09:00 burst and measure the fraction started within ten seconds. Ten thousand slots miss the target; a queue cannot supply missing compute. |
| FR3–4; NFR3–4: recover one accepted result | Atomic materialization and current-token completion; durable dispatch and safe database failover. | Pause A past its lease, let B complete and resume A. Reject A’s result; fail one zone and recover the same occurrence and accepted outcome. |
| FR4; NFR5: protected follow-up effects | Only committed result/outbox identities reach the notification adapter; external recovery follows its provider contract. | Lose the send response and race cancellation with completion. Preserve uncertainty where the destination cannot deduplicate; do not claim physical exactly-once execution. |
14Rapid revision
Remember: An expired lease does not stop a paused worker from resuming. The result store must reject its old token.
| Concern | Complete mechanism |
|---|---|
| Schedule | Save a versioned recurrence rule and time zone; decide how to handle overlapping and missed runs. |
| Occurrence | Record each intended run once, with its scheduled time and fixed inputs. |
| Materialization | In one transaction, save the run and dispatch task and advance the next due time. |
| Attempt | Give each worker claim a deadline and newer ownership token; retries keep the same intended run. |
| Completion | Check the current token and save the result and run state in one transaction. |
| No overlap | Allow only one active run per schedule; this is separate from rejecting duplicate attempts of that run. |
| External actions | Send only actions whose IDs are committed; recover uncertain results using the destination’s retry or lookup rules. |
| Capacity | Calculate how queued jobs fit into worker slots before their start deadlines; a queue adds no execution slots. |
| Recovery | Recover saved due times and runs; limit catch-up according to the chosen missed-run policy. |
Close with Tuesday R7 running twice but accepting only B's result after A pauses. Name the remaining external-effect boundary and the 09:00 capacity calculation. Those demonstrate a defensible scheduler without promising that every arbitrary task executes physically once.
Practise the interview questions
Say your answer aloud before opening the model answer. Then answer the follow-up and compare the reasoning.
How do schedule, occurrence and attempt differ?
Reveal a model answer
The schedule is the rule, an occurrence is one intended resolved run, and attempts retry that same run. This keeps history and business identity stable.
Interviewer follow-up
Can a retry use current edited parameters?
Reveal the follow-up answer
No. An already created occurrence retains its frozen revision.
What the answer must demonstrate: Preserve one intended occurrence and its frozen parameters across attempts.
Two scanners see Tuesday due. What prevents two runs?
Reveal a model answer
Both use the unique schedule/instant identity and transactionally create the occurrence, dispatch work and advance next due time.
Interviewer follow-up
Why not dispatch before committing?
Reveal the follow-up answer
A crash can leave the instant due after work has already been sent, recreating it.
What the answer must demonstrate: Atomically materialize the due instant, dispatch intent and next schedule time.
Why can an expired worker still be dangerous?
Reveal a model answer
Expiry changes permission but does not stop the process. A paused worker may resume and try to publish after replacement.
Interviewer follow-up
Where must fencing be checked?
Reveal the follow-up answer
At the protected result write, atomically with acceptance of the result.
What the answer must demonstrate: Enforce the current token at completion because expiry does not stop a process.
Does unique Tuesday occurrence identity prevent Wednesday overlapping?
Reveal a model answer
No. Those are distinct valid occurrences. A schedule-level active-run guard enforces the chosen no-overlap policy across them.
Interviewer follow-up
Does that prove no old process is running?
Reveal the follow-up answer
No. Stale physical execution still needs fenced results and protected effects.
What the answer must demonstrate: Use a separate schedule guard for distinct occurrences and retain physical-execution limits.
Can a unique run row guarantee one email?
Reveal a model answer
No. The email provider performs an external effect. Use a committed stable action identity and the provider’s supported idempotency or reconciliation.
Interviewer follow-up
Why prevent workers emailing private outputs?
Reveal the follow-up answer
Only the canonical accepted result should create the delivery intent, so stale attempts cannot publish alternate content.
What the answer must demonstrate: Protect external actions independently from local occurrence and result uniqueness.
Why do 50,000 five-second jobs exceed a ten-second start target with 10,000 slots?
Reveal a model answer
They start in five ideal waves at 0, 5, 10, 15 and 20 seconds. The last two waves start late before overhead is counted.
Interviewer follow-up
What are defensible options?
Reveal the follow-up answer
More prewarmed capacity, agreed jitter, priority reservation or rejecting an impossible promise.
What the answer must demonstrate: Compute execution waves and distinguish start-time targets from completion.
What must a calendar schedule say about daylight-saving changes?
Reveal a model answer
How nonexistent and repeated local times are handled, using a named zone and resolved execution instant.
Interviewer follow-up
How is fixed delay different?
Reveal the follow-up answer
Its next time depends on prior completion, unlike an anchored fixed-rate rule.
What the answer must demonstrate: State time-zone, daylight-saving and recurrence-mode behavior explicitly.
What does cancellation mean for a running job?
Reveal a model answer
It requests cooperative stopping and prevents later accepted completion/retries once cancellation becomes terminal. External effects already started may still complete.
Interviewer follow-up
What if completion already committed?
Reveal the follow-up answer
Return the committed outcome rather than pretend cancellation erased it.
What the answer must demonstrate: Serialize terminal cancellation with completion and disclose already-started external effects.
Blank-page exercise · 45 minutes
Build the answer yourself
Design a weekday report scheduler, then pause worker A beyond its lease, finish worker B and let A resume.
- Agree numbered functional and non-functional requirements, including start versus finish, overlap/catch-up policy and accepted-result guarantees. Then distinguish a recurring schedule, one intended occurrence and its execution attempts.
- Trace atomic due-run creation, worker claim and fenced completion.
- Explain calendar rules, missed-run policy and capacity limits before promising start deadlines.
- Trace a timeout and a concurrent request using the actual durable records.
- Review the final architecture against every numbered requirement, including measured bottlenecks, targets still needing validation and remaining failure limits.
Check that each component and design decision follows from your requirements and workload.
Recall the key ideas
Answer from memory before opening each card. Explain why the choice works and what it costs. Revisit missed cards tomorrow.
Design a distributed job schedulerA pauses with token 41; B completes R7 with token 42. What happens when A resumes?Recall first, then reveal
R7 and its frozen parameters stay the same, but the database rejects A’s old token. B’s saved result remains the accepted result.
One intended run may execute twice; only the current attempt can save its result.
Return to lessonDesign a distributed job schedulerHow is a late result from a replaced worker rejected?Recall first, then reveal
The result store checks the ownership token atomically when saving completion and rejects an old token.
The result store checks the token when accepting completion.
Return to lessonDesign a distributed job schedulerDoes a ready queue let more jobs execute at once?Recall first, then reveal
No. It buffers and schedules jobs; workers still need enough execution capacity.
Queue is not compute.
Return to lessonFinal revision
Summary and interview notes
Save each intended run, retry failed attempts, and accept results only from the current worker. Define when jobs start, whether runs may overlap, and how external actions recover after uncertain outcomes.
Remember these points
- Give each intended run a durable identity.
- Create due work and advance recurrence atomically.
- Check the current token when saving results; an expired worker may still be running.
- State calendar, overlap, catch-up and capacity policies explicitly.
Interview tips
- Separate schedule, intended run and execution attempt.
- Pause A with token 41, let B finish with 42, then resume A and inspect R7.
Important qualifications
- Traffic and latency figures are interview assumptions, not claims about a named company's deployment.
Continue after the core interview
Explore the advanced version
The advanced lesson keeps the full detailed design. Use these sections when you want to examine the stronger requirements and failure cases.
- Large-object result publication
Immutable object verification, references and cleanup add a separate lifetime protocol.
- Full external-effect and worker failure analysis
Arbitrary destinations require their own enforcement and uncertain-outcome policies.
- Advanced partitioning and timing accelerators
Due buckets and timer structures become useful after measured scan bottlenecks.
- Workflow and untrusted-code extensions
Dependencies and arbitrary code change the state and security models substantially.
Technical references
- Kubernetes CronJob documentationOfficial discussion of time zones, concurrency policies, missed schedules, and job idempotency.
- Temporal workflow execution overviewConcrete durable-execution model that separates persisted workflow progress from worker processes.
- Transactional outbox patternSupports recoverable database-to-queue dispatch.
Practice marks stay in this browser.