System designby Learnastra

System-design interview · Extended interviews

Design a distributed job scheduler

By Anup Rai

Design a scheduler that records each due run, limits dispatch to available capacity, replaces failed workers safely and recovers report publication or delivery after a lost response.

You will learn to

  • Distinguish a recurring schedule, one due occurrence, and its execution attempts.
  • Recover dispatch and worker failure without assuming exactly-once side effects.
  • Calculate start lag under bursts and define timezone, catch-up, and cancellation behavior.

Practice in this chapter

8 interview questions with model answers and follow-ups.

Go to interview practice

Useful 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.

01Problem and scope

A distributed scheduler records each run due under a recurring rule and assigns workers to attempt it. A schedule defines when work is due; an occurrence identifies one intended run; attempts are executions or retries of that run. Remembering due work is different from preventing duplicate external effects. A weekday sales report at 09:00 America/New_York provides the example: restart must not lose the resolved occurrence, and result or email publication needs a stable identity.

A schedule is the recurring rule. An occurrence identifies one intended run, such as (sales-report,Tuesday-09:00-resolved-UTC). An attempt is one worker's try at that occurrence. Preserve the occurrence identity across retries. A scheduler can guarantee durable tracking and retry policy; an arbitrary job's external side effects still require their own idempotency or reconciliation.

I clarify the product before choosing a queue. The interviewer says, “Run reports every weekday.” I ask whether a report means invoking a registered function, launching arbitrary user code, or completing an external email. We choose approved report jobs with versioned parameters. The scheduler guarantees durable tracking of each intended run and safe publication of its result; an email adapter separately handles repeatable requests. Retrying a run may execute it again; the promise concerns its saved result and protected effects.

The scheduling client can therefore see Tuesday's report as one run with two attempts, rather than two unrelated reports. We exclude a general dependency graph and interactive workflows from the first design. They would require additional dependency state, not merely a more elaborate cron expression, the compact calendar notation used to specify matching minutes, hours and days.

02Functional requirements

Schedule changes affect future work; occurrence operations affect a specific intended run. We materialize an occurrence when we persist the record for one due instant. After that point, retries keep its recorded identity and payload.

  1. Create or revise a schedule. Return its ID, revision, and next resolved instant. Revision affects future unmaterialized occurrences.
  2. Run once. Return a durable occurrence ID. Repeating a request key returns the same occurrence.
  3. Execute when due. Show queued, running, then terminal status. A retry remains the same occurrence.
  4. Inspect history. Show scheduled time, attempt history, and result. An attempt log is not proof of external delivery.
  5. Cancel. Report whether cancellation was only requested or actually accepted. Completed side effects cannot be recalled.
  6. Pause and resume. Preserve the rule and apply its declared misfire policy—the rule for scheduled times missed while paused or unavailable. Do not silently replay every missed instant.

Scope and acceptance boundaries

For the scheduling client, we choose one occurrence per local weekday at 09:00, no logical overlap between occurrences of this schedule, and “latest only” catch-up after a prolonged outage. A skipped older occurrence is represented as skipped with a reason. It is not quietly deleted from history.

Schedule edits carry an expected revision. If the scheduling client and an administrator both change revision 3, one obtains revision 4 and the other receives a conflict with the current rule. An already materialized Tuesday occurrence retains its revision-3 payload. Otherwise a retry could execute different business work under the same identity. The UI previews the next five resolved instants before saving a time-zone rule.

03Non-functional requirements

  1. API latency and availability. Assume schedule creation/status p95 below 200 ms and 99.9% monthly availability.
  2. Start deadline. For admitted ordinary jobs, 99% should start within 10 seconds of their due instant. Rejected work is excluded; accepted jobs remain in the denominator during infrastructure overload.
  3. Durability. Preserve accepted occurrence state through one availability-zone failure using synchronously replicated storage.
  4. Retention. Keep run metadata for 30 days and report objects for 7 days.
  5. Ownership and result consistency. During authority loss, delay claims rather than issue two current ownership tokens. A completed run has one canonical result pointer: the stored reference to the output the scheduler accepted for that run.
  6. External-effect limits. Execution may happen twice. An email provider without suitable support requires an explicit uncertain-outcome policy; an expired lease does not prove the old worker stopped.

Scheduling semantics

“Every hour” can mean the top of each clock hour, every sixty minutes from a fixed anchor, or sixty minutes after the previous run finishes. Choose the recurrence policy before calculating due times.

Schedule type How the next due time is derived
Calendar/cron Named wall-clock time in a time zone
Fixed rate Anchor plus a multiple of the interval
Fixed delay Interval after the previous run completes

Start time is not completion time

The targets are exercise objectives, not properties of cron syntax. A three-hour report cannot have a ten-second completion promise. Measure queue delay and execution duration separately; safe result publication takes priority over punctual execution.

04Capacity estimates

Assume ten million occurrences/day, five-second jobs in the worked burst, and the stated retention policies.

Estimate Arithmetic Consequence
Average dispatch Ten million/day ≈ 116 starts/s Does not describe a synchronized burst
09:00 burst 50,000 due jobs / 10,000 slots = five waves at 0, 5, 10, 15, 20 seconds Last start lags 20 seconds; last completion is at second 25
Ten-second start bound At least 16,667 slots give three ideal waves at 0, 5, 10 Real overhead and uneven durations require margin
Run metadata Ten million × 2 KB = 20 GB/day; 30 days = 600 GB Before indexes and replicas
Burst claim rate 50,000 claims / ten seconds ≈ 5,000/s Before completion and heartbeat writes
Heartbeats 10,000 live attempts / five seconds = 2,000 updates/s Keep claim transactions short
Worker memory 10,000 active attempts × 256 MB = 2.56 TB Across the whole fleet
Report output Ten million × 1 MB = 10 TB/day; seven days ≈ 70 TB Before redundancy; dwarfs metadata
Safe utilization 10,000 slots × assumed 80% = 8,000 planned active slots Nominal slots are not all schedulable headroom

Start deadline versus completion time

The 10,000-slot fleet cannot start the entire burst within ten seconds. Options are more capacity, user-approved jitter, capacity reserved for strict jobs, or rejecting an impossible admission promise. A queue preserves the burst but does not create execution capacity.

Separate control and execution resources

Job duration and resource mix determine execution capacity; scheduler QPS does not. Report execution must not hold a database row lock for five seconds. Keep large results in object storage and control metadata/result retention independently.

Provision before the known 09:00 wave, negotiate jitter, or admit fewer deadline-bound jobs. The resource calculation determines which deadline promises are feasible.

05APIs and contracts

The API separates schedule management from worker ownership. A worker holds a lease, permission valid until a deadline; a heartbeat asks the authority to extend it. Each new attempt receives a fencing token, an increasing ownership number that the protected store checks before accepting its result.

API Example Meaning
Create schedule POST /schedules with {"idempotencyKey":"schedule-request-7","cron":"0 9 * * 1-5","timeZone":"America/New_York","overlap":"forbid"} Return s7 with explicit missed-run/DST policy
Inspect runs GET /schedules/s7/runs?after=... Paginated occurrence history
Cancel POST /runs/r7/cancel Request cancellation; report current boundary
Worker heartbeat POST /attempts/a2/heartbeat with lease token Extend current ownership if valid

In the example, 0 9 * * 1-5 means 09:00 on Monday through Friday in the named time zone. The stored records connect this recurring rule to each concrete run and its execution attempts.

Record Key and fields Query
Schedule scheduleId, rule, zone, nextRunAt, version Indexed due-time scan
Occurrence unique (scheduleId,scheduledInstant), state, payload version One logical intended run
Attempt run ID, attempt number, lease expiry, fencing token Claim, timeout, and recovery
Outbox / result dispatch ID; result URI and checksum Recover handoff and inspect output

Store timestamps as instants for execution and retain the calendar rule/timezone for future calculation. Updating a schedule changes future occurrences under an explicit version; it should not silently rewrite completed history.

Creation returns 201 with {scheduleId:"s7", revision:1, nextRunAt:"2026-09-23T13:00:00Z"} for an appropriate Eastern daylight-time weekday. Invalid cron or time-zone identifiers receive 400; unauthorized job types receive 403; an exhausted tenant quota receives 429 before acceptance. The same idempotency key with different parameters receives 409, rather than silently modifying the old schedule.

A run response includes {runId:"r7", state:"running", attempt:2, scheduledAt:..., startedAt:..., cancellationRequested:false}. A timeout creating a run means the client must retry its key or query it; it does not mean creation failed. History uses an opaque cursor over (scheduledAt,runId) and returns a stable upper-bound time so new runs do not make the client skip older ones.

Worker claim, heartbeat, and completion endpoints require worker credentials, run ID, attempt ID, and fencing token. Clients never choose their own token. A stale completion returns a distinct conflict. The worker can stop retrying because its result can no longer become the run’s accepted output.

06Data model and access patterns

The no-overlap requirement applies across different occurrences of the same schedule. Its active-run guard records which occurrence currently holds permission to run, so the database must check that guard and claim the occurrence in one transaction. The outbox records committed dispatch intent for later delivery to workers.

The schedule ID determines the authority shard. The following records live in one transactional database partition for that schedule:

Record and key Material fields Query or invariant
Schedule (tenant,scheduleId) revision, rule, zone, nextRunAt, activeRunId Conditional rule edit; one logical active occurrence
Occurrence (scheduleId,scheduledInstant) runId, payloadRevision, state, currentFence Unique materialization of a due instant
Attempt (runId,attemptNo) worker, fence, leaseUntil, outcome Inspect retries; reject obsolete completion
Outbox (eventId) runId, eventType, publishedAt Recover committed dispatch intent
Request key (tenant,createKey) or (tenant,scheduleId,runKey) requestHash, runId or scheduleId Repeat a timed-out creation safely

A due bucket is a stable group of schedules scanned together. Bucketing lets scanners divide the search for due work as the schedule population grows.

A due index on (bucket,nextRunAt,scheduleId) finds a bounded batch without scanning every schedule. A run-history index serves one schedule in time order. Tenant-wide history is built from run records and may lag. The individual run endpoint reads current state from its owning database.

Report bytes belong in object storage under immutable names such as r7/a2/result. The occurrence contains the chosen object's checksum and URI after successful completion. A stale attempt can create an orphan object, but cannot overwrite the canonical result. The run authority also stores an object registry row linking each staged upload to its attempt and fence. Publication atomically changes that row from staged to referenced with the canonical result pointer. Cleanup locks the same row and run state and may mark it deleting only when no retained reference exists and its attempt no longer has publication authority. Publication rejects deleting rows. Waiting before cleanup avoids repeatedly creating and deleting recent uploads, but safety comes from the transaction ordering: publication wins the transaction and protects the object, or deletion wins and publication fails.

The queue contains dispatch hints, not the only record of a due run. Losing or duplicating one hint is recoverable from the outbox and occurrence state.

Create-schedule requests route deterministically by authenticated tenant and creation key to a stable creation bucket; the resulting schedule stays in that bucket or follows its versioned owner mapping. This lets schedule creation and its retry result share one transaction. Manual runs of an existing schedule instead scope their request key to that schedule. If request keys lived in an unrelated database, schedule creation and its retry result could not commit in the same local transaction.

A result reference includes the verified immutable object version, not only a reusable object name. Enforce conditional creation or pin a storage VersionId and read that exact version; an attempt-scoped upload credential alone does not prevent overwriting its own path. A result download records a bounded pin that prevents cleanup of that object version until the promised download URL expires. The staged-object registry exists before upload begins and is checked atomically during publication.

07Basic working design

The smallest useful service consists of a schedule API, one database, a scanner that polls its due index once per second, and a bounded worker pool. At Tuesday 09:00, the scanner locks schedule 7, verifies its revision and next due instant, inserts occurrence r7, inserts dispatch intent dispatch-r7, and advances nextRunAt in one transaction. Only then is the recorded run eligible for dispatch.

The baseline uses an outbox table and poller without a broker. Restarting after commit preserves dispatch intent; restarting before commit leaves the instant due. Execution capacity and scanning are its limits.

We use database time for lease comparisons within this authority. Scheduling wall-clock instants and measuring elapsed lease duration are different concerns; worker clocks do not decide whether another worker owns a run.

architecture · baselineOne scanner and durable run state

Materialization and dispatch intent share a transaction; report execution does not hold its locks.

One scanner and durable run stateMaterialization and dispatch intent share a transaction; report execution does not hold its locks. client to api: Create schedule 7; api to db: Commit rule and request key; scan to db: Materialize r7 + dispatch intent; scan to worker: Dispatch r7; worker to db: Claim / conditional completion; worker to objects: Upload r7/a1 resultCreate schedule 7Commit rule and request keyMaterialize r7 + dispatch intentDispatch r7Claim / conditional completionUpload r7/a1 resultACTORSchedule clientsSERVICESchedule APISTORESchedule and rundatabaseWORKERDue scanner /dispatcherWORKERBounded worker poolSTOREImmutable resultobjectssyncasync
Read each connection in order
  1. syncCreate schedule 7Schedule clients → Schedule API
  2. syncCommit rule and request keySchedule API → Schedule and run database
  3. syncMaterialize r7 + dispatch intentDue scanner / dispatcher → Schedule and run database
  4. asyncDispatch r7Due scanner / dispatcher → Bounded worker pool
  5. syncClaim / conditional completionBounded worker pool → Schedule and run database
  6. syncUpload r7/a1 resultBounded worker pool → Immutable result objects

08Find the baseline flaws

Consider the 50,000-run morning burst. The baseline's 10,000 slots start waves at seconds 0, 5, 10, 15, and 20 under the idealized five-second duration assumption. Forty percent start after the ten-second objective even before scanner and claim overhead. Polling faster cannot fix occupied slots. Longer reports delay the last jobs' start times further, so a single average duration is insufficient for admission.

Now suppose worker A claims r7 with fence 41, pauses during a runtime stall, and misses its lease. A replacement worker B claims fence 42 and finishes. A later resumes. Merely checking that A's lease was valid when it started does not stop A from overwriting B's report or sending another email. The database must reject A's completion against the current token, and external effects need their own protection.

Another error appears with a second scanner: both read Tuesday as due and independently send messages before updating nextRunAt. Two workers then execute what should be one occurrence. The unique occurrence key and materialization transaction prevent that race even when scanners overlap. Leader election alone does not.

Finally, a scheduler outage lasting two days can release several million overdue jobs at once. Treating catch-up as an unbounded loop turns recovery into overload. The missed-run policy must therefore be chosen when defining the schedule, before an outage occurs.

09Improve the design, step by step

First, separate durable dispatch from execution. The trigger is scanner delay while workers are busy. An outbox relay sends small run IDs into ready queues; dispatchers claim from the authoritative run store before allocating a worker. Scanning now remains responsive during a report burst. The cost is extra delivery latency, queue storage, and duplicate messages. We retain occurrence checks because a relay can publish twice. At modest load, the rejected alternative—a database-backed work queue—remains simpler and adequate.

Second, partition due scanning and metadata. The trigger is a saturated due index or claim write path. A stable hash of schedule ID assigns due buckets and their schedules to database shards. A bucket directory assigns scanners, and each scanner owns a short renewable lease to reduce redundant work. Unique occurrence creation remains authoritative even if two scanners overlap during reassignment. This increases aggregate scan and write throughput, but adds routing, rebalancing, and uneven-bucket risk. Keep one larger database while it meets the targets without the added work of managing shards.

Third, isolate resource classes and tenants. The trigger is short reports waiting behind hour-long exports. Dispatch applies per-tenant active limits and separate pools for short, long, and memory-heavy work. A weighted fair policy reserves capacity for smaller tenants while allowing bounded borrowing. Short jobs wait less behind long jobs, and admission becomes more predictable. Some reserved slots may sit idle, reducing total utilization. A single FIFO queue is appropriate when jobs are homogeneous and strict arrival fairness is the product requirement.

Fourth, prepare near-term work and planned bursts. The trigger is a large population of far-future schedules making tight polling expensive. Scanners load only a short future horizon into an in-memory timer structure and periodically refresh it from durable state. Workers are started and made ready before known daily peaks; after a restart, timers are rebuilt from the due index. This lowers polling work and cold-start delay, but timers can be stale after edits, so materialization still verifies schedule revision. We reject making an in-memory timer wheel, which groups timers into time slots, the sole authority: it would forget work on restart. A dedicated durable workflow engine becomes attractive when dependencies, signals, and long-running state exceed these schedule semantics.

Each step preserves the same occurrence identity and result-publication guard. The service evolves its execution machinery without redefining what Tuesday's report means to the scheduling client.

10Detailed architecture

Scheduling authority

The API authenticates each request, finds the schedule’s shard and checks the tenant’s capacity limits. In the scheduling authority, replicated databases own rules, occurrences, attempts, active-run guards, and outbox rows. Scanner leases distribute work across due buckets; they do not replace transactional uniqueness.

Execution pools and attempts

Execution uses the relay, ready queues, dispatchers and worker pools for each resource class. Queue messages may repeat. Dispatchers turn a hint into an authorized attempt by calling the run authority, which atomically allocates the next fence. Workers heartbeat through that same authority and upload immutable output directly to the object store using credentials restricted to their run and attempt prefix.

Canonical external-effect intent

For report email, successful result publication creates an outbox intent carrying stable action identity send-r7 and the canonical object version. The effect adapter accepts only this committed intent, so a stale worker cannot email a different private output under that key. Other external job effects still require their own authorization and identity protocol. The effect adapter records requests and provider outcomes. If a provider supports durable idempotency within the required retry period, it reuses that identity. Otherwise the product explicitly handles uncertain delivery rather than promising exactly one email.

Timing boundaries and implementation

Creation, claims, heartbeats, completion and current-status reads wait for the owning database’s decision. Dispatch, report execution, history indexing, and object cleanup are asynchronous. Database replicas span failure domains within a region; a failover mechanism must preserve current ownership and committed state. We avoid a multi-region active-active schedule writer in this version because two independent schedule writers would need to agree on occurrence creation and the active-run guard.

A PostgreSQL implementation can scan indexed due rows and use short FOR UPDATE SKIP LOCKED transactions to distribute independent claim work; its skipped-row view is appropriate for this work queue, not a general consistent report. A broker is optional until dispatch load justifies it. Kubernetes CronJob is useful for simpler periodic container jobs but documents approximate scheduling and the need for idempotent jobs; it is not a substitute for this custom result/effect protocol. A durable workflow engine such as Temporal becomes attractive when persisted dependencies, timers and signals dominate the product.

architecture · finalScheduling authority and execution pools

The queue may repeat a hint; only the run authority can allocate a current attempt and publish its result.

Scheduling authority and execution poolsThe queue may repeat a hint; only the run authority can allocate a current attempt and publish its result. client to api: 1. Schedule / status / cancel; api to db: 2. Route and transact; db to replica: Committed state; scan to db: 3. Materialize due occurrence; relay to db: Read durable dispatch intent; relay to queue: 4. Publish run ID; queue to dispatch: Ready hint; dispatch to db: 5. Claim with guard and fence; dispatch to worker: 6. Start authorized attempt; worker to db: Heartbeat / publish result; worker to objects: Upload immutable bytes; relay to effect: Committed send-r7 + canonical version; effect to provider: Stable external action; api to objects: Authorize result download1. Schedule / status / cancel2. Route andtransactCommitted state3. Materialize due occurrenceRead durable dispatch intent4. Publish run IDReady hint5. Claim with guard and fence6. Start authorized attemptHeartbeat / publish resultUpload immutable bytesCommitted send-r7 + canonicalversionStable external actionAuthorize result downloadACTORSchedule clientsG1SERVICEAuthenticated API +shard routerG1STORESchedule / runauthorityG2STOREAuthority replicasG2WORKERDue-bucket scannersG2WORKEROutbox relayG2QUEUEReady queues byresource classG3SERVICEFair dispatcherG3WORKERSandboxed workerpoolsG3STOREImmutable resultstoreG3SERVICEIdempotent effectadapterG3EXTERNALEmail providerG4syncreplicationasyncG1 Entry and accessG2 Schedule shard ownershipG3 Execution and publicationG4 External delivery
Read each connection in order
  1. sync1. Schedule / status / cancelSchedule clients → Authenticated API + shard router
  2. sync2. Route and transactAuthenticated API + shard router → Schedule / run authority
  3. replicationCommitted stateSchedule / run authority → Authority replicas
  4. sync3. Materialize due occurrenceDue-bucket scanners → Schedule / run authority
  5. syncRead durable dispatch intentOutbox relay → Schedule / run authority
  6. async4. Publish run IDOutbox relay → Ready queues by resource class
  7. asyncReady hintReady queues by resource class → Fair dispatcher
  8. sync5. Claim with guard and fenceFair dispatcher → Schedule / run authority
  9. async6. Start authorized attemptFair dispatcher → Sandboxed worker pools
  10. syncHeartbeat / publish resultSandboxed worker pools → Schedule / run authority
  11. syncUpload immutable bytesSandboxed worker pools → Immutable result store
  12. asyncCommitted send-r7 + canonical versionOutbox relay → Idempotent effect adapter
  13. syncStable external actionIdempotent effect adapter → Email provider
  14. syncAuthorize result downloadAuthenticated API + shard router → Immutable result store

11Write path and acknowledgement

Create each due occurrence once in durable scheduler state, then dispatch recoverably. Attempts may repeat, while accepted state transitions require current ownership.

  1. The scheduler reads due schedule s7. In a transaction it inserts occurrence r7 for the resolved Tuesday instant, records outbox dispatch-r7, and advances nextRunAt. A competing scheduler hits the same unique occurrence key and cannot create another logical Tuesday run.

  2. The dispatcher publishes r7 to the ready queue. Publishing twice is possible after a lost acknowledgment; workers therefore do not treat each queue message as a new occurrence.

  3. Worker A atomically claims the schedule’s active-run guard for r7, changes r7 from ready to running, creates attempt a1, and obtains lease token 41 with an expiration. If another occurrence of s7 still holds that guard, this occurrence stays pending under the chosen overlap policy. A lease is time-limited ownership that must be renewed; it is not proof that a process has stopped when time expires.

  4. A builds the report and writes a versioned result object. Completion is accepted only if A still owns token 41. No attempt sends its private output directly. Only successful canonical publication creates the report-delivery intent.

  5. The service records succeeded, immutable result version and completion time, creates unique outbox intent send-r7 for that canonical result, and releases the schedule guard in the same transaction; queue acknowledgment follows. The scheduling client can inspect which intended time ran, how late it started, and which attempts occurred.

  6. If the worker's completion response is lost, it queries r7 before doing more work. A repeated completion from the current attempt with the same checksum returns the saved terminal result; a conflicting checksum receives a conflict. The canonical result never changes because a reply disappeared.

  7. The effect adapter reads the committed send-r7 intent and canonical result version, submits that stable logical action, and stores the resulting external identifier, and reconciles timeouts against that identifier or provider idempotency key. Report completion and email delivery can be separate visible states. The scheduling client sees “report ready, email pending” instead of a misleading all-or-nothing success.

  8. Queue acknowledgment follows the durable claim or recognized terminal state according to the dispatch contract. Run recovery does not depend on a broker continuing to retain an unacknowledged message: the authority sweeps expired running attempts and writes fresh dispatch intents.

If report computation fails transiently, the authority schedules the next attempt with bounded exponential delay and jitter. An invalid report query is a permanent failure and does not consume an unlimited retry budget. The original payload revision remains attached to every retry.

12Read and delivery path

The history list is derived from authoritative run changes and may lag them. Its watermark reports how far that processing has progressed, while opening an individual run reads its owner directly. This explains why a just-created run can exist before it appears in history.

Run history reports one occurrence with its attempts and authoritative outcome. Result reads use the committed publication, not a worker-local completion claim.

  1. The scheduling client requests /schedules/7/runs?after=.... The API verifies tenant ownership and reads a cursor-bounded history projection. It includes the projection watermark, so a recent creation that is still missing from the list can be distinguished from an absent run.
  2. Opening r7 routes to its schedule shard. The authoritative response identifies the one occurrence, current attempt, due and actual start times, and any requested cancellation. A stale read replica is not used when deciding whether a cancel or manual retry is still legal.
  3. For completed r7, the API checks access again and issues a short-lived download URL for the canonical immutable result object and checksum. It does not guess the latest object by lexicographic filename; a stale attempt may have uploaded a newer-looking orphan.
  4. Attempt history explains a1 as expired and a2 as successful. It does not expose worker credentials, raw secret parameters, or another tenant's object paths.
  5. If the scheduling client cancels while r7 is running, the API records cancellationRequested and returns that intermediate state. Workers observe it at checkpoints. Once the authority accepts a canceled terminal result, no new retry is dispatched, although an already-started external operation may still complete.

List caching is allowed for a few seconds under the declared freshness target. The UI can refresh the one active run from its owner while keeping older immutable history cached. This avoids turning every dashboard poll into a full scan of the attempt table.

13Correctness deep dive

The database checks result ownership in the same transaction that changes the run and schedule. Uploading bytes is not publication. Assume r7 is running, schedule 7's activeRunId is r7, and its current token is 42.

complete(run=r7, attempt=a2, token=42, object=O2, hash=H2):
  begin transaction; lock schedule 7, then occurrence r7
  if r7 is terminal with this accepted attempt and same hash:
      return its recorded result
  require state == running and currentFence == 42
  require leaseUntil > freshAuthorityTimeAfterLockWaits()
  require verified object version has a live staging grant, not deleting
  require activeRunId == r7 and not cancellationAccepted
  transfer staging grant to canonical result reference
  set r7 = succeeded, result = (O2,versionId,H2), acceptedAttempt = a2
  insert unique send-r7 intent referencing that canonical version
  clear activeRunId only if it still equals r7
  commit

Claims, heartbeats, completion and cancellation use a consistent schedule-then-occurrence lock order. Lease expiry is checked with fresh authority time after lock waits, not a timestamp captured when the transaction began.

The object must already exist and match the claimed checksum; the worker cannot publish a pointer to an unfinished multipart upload. This check does not require object creation and the database to share a transaction: unreferenced uploads are allowed, missing canonical objects are not.

Time Actor Durable state after the operation
09:00:00 A claims a1 r7 running; token 41; active guard r7
09:00:15 Authority expires A and B claims a2 r7 running; token 42; same active guard
09:00:18 B uploads O2 and completes r7 succeeded; canonical O2; guard released
09:00:19 A uploads O1 and completes with 41 Completion rejected; O1 is an orphan

If A's completion races lease expiry and B's replacement claim, the serialized authority transaction selects one ordering. A completes first and B cannot claim a terminal run, or B advances the token first and A cannot complete. There is no ordering in which both results become canonical.

sequence · stale-attemptA stale worker cannot publish

Worker A may upload an unreferenced object, but the database updates the run's result pointer only for the current attempt's valid token.

A stale worker cannot publishWorker A may upload an unreferenced object, but the database updates the run's result pointer only for the current attempt's valid token. a to db: Claim r7; receive fence 41; a to a: Pause past lease; b to db: Replace expired attempt; fence 42; b to obj: Upload immutable O2; b to db: Complete with current fence 42; db to b: Commit canonical O2; a to obj: Upload immutable O1 after resuming; a to db: Complete with stale fence 41; db to a: Reject; canonical O2 remainsPARTICIPANTWorker APARTICIPANTRun authorityPARTICIPANTWorker BPARTICIPANTResult objects1. Claim r7; receive fence 412. Pause past lease3. Replace expired attempt;fence 424. Upload immutable O25. Complete with currentfence 426. Commit canonical O27. Upload immutable O1 after resuming8. Complete with stale fence419. Reject; canonical O2remainssyncreturnblocked
Read each connection in order
  1. syncClaim r7; receive fence 41Worker A → Run authority
  2. syncPause past leaseWorker A → Worker A
  3. syncReplace expired attempt; fence 42Worker B → Run authority
  4. syncUpload immutable O2Worker B → Result objects
  5. syncComplete with current fence 42Worker B → Run authority
  6. returnCommit canonical O2Run authority → Worker B
  7. syncUpload immutable O1 after resumingWorker A → Result objects
  8. syncComplete with stale fence 41Worker A → Run authority
  9. blockedReject; canonical O2 remainsRun authority → Worker A

14Failure and recovery

Failure or race Required response and boundary
Expired worker resumes Worker A pauses long enough for its lease under token 41 to expire. Worker B claims the same occurrence with token 42 and completes. Then A resumes. A local “my lease was valid earlier” check cannot prevent its stale side effect. A fencing token is an increasing ownership number that the protected destination checks; once it accepts token 42, it rejects token 41. Completion updates in our database must compare the current token.
External destination cannot fence For an external email or payment API that does not understand fencing, use a stable logical operation key such as send-r7 with the destination's idempotency contract, or reconcile ambiguous outcomes. If the destination supports neither, exactly-once effects are not guaranteed. A unique run-table row does not prevent an external provider from performing the same action twice.
Crash before or after an effect If the scheduler crashes after committing r7 but before publishing, the outbox dispatcher recovers it. If a worker dies before any effect, retry after lease expiry and backoff. If it dies after an effect but before recording completion, replay the same logical operation identity. Classify permanent failures separately from transient ones and cap attempts; a poison job must not consume the fleet forever.
Different occurrences overlap A unique occurrence prevents duplicate Tuesday records; it does not stop Wednesday from starting while Tuesday is still active. The per-schedule guard enforces the declared no-overlap rule across distinct occurrences and remains assigned to the same occurrence during retries. If an old process may still run after lease expiry, fencing or destination idempotency protects accepted effects; claiming replacement work does not prove that the old process physically stopped.
Worker or database partition A partition between a worker and its authority makes heartbeats fail. The worker stops initiating new protected actions and attempts cooperative cancellation, but process termination is not the safety proof. The authority may later grant a new token, and the canonical-result guard handles a paused process that ignores cancellation. If the database loses its write quorum, new claims pause; existing computation can produce staged objects but cannot publish them authoritatively.
Overload and controlled catch-up During overload, already accepted runs remain durable and display delayed status. Admission rejects new deadline commitments before creating them. A recovery controller applies the declared misfire policy in bounded pages and records which instants were skipped. It does not hide backlog by resetting due timestamps to now. An operator can separately choose a controlled backfill whose new action identity makes its extra business execution explicit.

15Operations, security, and cost

The primary alert is due-to-start lateness by tenant and resource class. Queue depth alone is ambiguous: 100 one-second jobs and 100 hour-long jobs imply very different delay. Track estimated work seconds, oldest admitted due time, expired lease rate, rejected stale completions, skipped misfires, and output publication failures. A rising fence-rejection rate can reveal worker pauses even while aggregate throughput looks healthy.

Roll out a new cron parser in shadow mode against saved rules, comparing the next month of resolved instants across daylight-saving transitions before activation. A migration copies a bucket, drains or redirects its owner through a versioned routing change, and preserves unique occurrence keys and fence counters. Test scanner death after materialization, worker death after upload, and completion-response loss with a recovery drill. Restore testing must include outbox and request-key records, not only schedule rows.

Resource cost is dominated by execution and output retention. At 10 million five-second runs daily, the workload consumes 50 million slot-seconds, about 13,889 slot-hours per day before idle capacity. Keeping 10,000 slots warm continuously supplies 240,000 slot-hours per day. That large gap motivates burst-aware provisioning, while the ten-second start target limits how aggressively we can scale to zero.

16Decision ledger and limitations

The scheduler remembers each intended run and accepts a result only from its current attempt. The choices below support that promise while allowing attempts to repeat and requiring separate protection for external effects.

Decision Benefit Cost
Durable unique occurrence Recovery without inventing a second run Transaction and history storage
Lease plus checked fencing Replace failed workers safely Destination cooperation is required
At-least-once attempts Recover from uncertain failure Jobs need duplicate-safe effects
Explicit missed-run policy Predictable outage recovery Some occurrences are intentionally skipped or coalesced

The main remaining bottleneck is a single very busy schedule with overlap forbidden: adding workers cannot make its serial business work concurrent. We can split it into independent schedules only if the report semantics allow separate partitions and a later merge. A workflow graph is the next design when jobs depend on one another, not a hidden feature of the ready queue.

We favor same-region authoritative ownership for correctness and bounded failover. Disaster recovery to an asynchronous remote replica would have a declared recovery-point loss unless accepted state also survives there. We would stop, reconcile, and explicitly account for affected occurrences before resuming side effects. “Highly available” alone does not answer which Tuesday reports might be repeated or missing.

17Interview closing

“I have separated the recurrence rule, a scheduled occurrence, and an execution attempt. Each scheduled instant has one durable occurrence identity even when execution requires multiple attempts. I start with a transactional schedule store: materializing a due instant, advancing the rule, and creating dispatch intent commit together. A queue and fair worker pools then absorb bursts, while schedule-based partitioning distributes scanning and state writes.

“The hard guarantee is one canonical result for an occurrence. A current fencing token is checked atomically when the run publishes its immutable output pointer, so a resumed old worker cannot overwrite the replacement's result. That does not make arbitrary external effects execute once; the email adapter needs its own stable action identity and reconciliation policy.

“The 50,000-run morning burst fails a ten-second start target with 10,000 five-second slots, so I would measure runtime tails and warm capacity before promising that deadline. Misfire, overlap, time-zone, and cancellation behavior are explicit product policies. My next test is to pause a worker past its lease, complete a replacement, and prove that the old worker's result remains unreferenced.”

If the interviewer adds month-long workflows with human approvals, I would keep the run and effect identities but introduce durable workflow state and event history. If the requirement instead becomes arbitrary untrusted code, sandbox isolation, network policy, resource metering, and secret access become central parts of execution rather than small additions to the scheduler.

Practise the interview questions

Say your answer aloud before opening the model answer. Then answer the follow-up and compare the reasoning.

Foundation · Question 1

What is the difference between a schedule and a job run?

Reveal a model answer

The schedule is a rule such as weekdays at 09:00 in New York. A run is one resolved intended occurrence. That run can have several attempts after failures without becoming several logical Tuesday reports.

What the answer must demonstrate: Do not deduplicate all future recurrences together.

Applied · Question 2

Can 10,000 slots start 50,000 five-second jobs within ten seconds?

Reveal a model answer

No under the given equal-duration model. Start waves are at 0, 5, 10, 15, and 20 seconds. I need more reserved slots, permitted jitter, or a weaker admission promise.

What the answer must demonstrate: Check the arithmetic before promising an SLA.

Applied · Question 3

The old worker resumes after its lease expired. What prevents a second effect?

Reveal a model answer

Our state updates compare the current fencing token, and cooperative downstream storage rejects older tokens. For external APIs, a stable logical action ID can provide deduplication if supported.

What the answer must demonstrate: Show the pause between check and effect.

Foundation · Question 4

How do you avoid losing a run between database insert and queue publish?

Reveal a model answer

Create the occurrence and an outbox record in one transaction. A dispatcher retries publishing that record, and workers deduplicate or atomically claim the stable occurrence ID.

What the answer must demonstrate: Close the handoff gap without claiming perfect queues.

Follow-up · Question 5

What does every day at 02:30 mean across daylight saving?

Reveal a model answer

It is ambiguous unless the product specifies a timezone and a skip/shift policy for nonexistent times plus a once/twice policy for repeated times. I store the rule and resolved execution instant.

What the answer must demonstrate: Calendar time is a product contract.

Follow-up · Question 6

Can cancellation guarantee the report email is never sent?

Reveal a model answer

Only before the external-send boundary. A running job can cooperate with cancellation at checkpoints, but a completed external send may be irreversible. The status should report that distinction.

What the answer must demonstrate: Cancellation and rollback are not synonyms.

Applied · Question 7

Wednesday becomes due while Tuesday is retrying. What does no overlap mean?

Reveal a model answer

I keep a schedule-level active-run guard owned by Tuesday r7 across its attempts. Wednesday can be materialized and remain pending, but its claim cannot acquire the guard until Tuesday reaches a terminal state. A lease expiry replaces an attempt of Tuesday; it does not make Wednesday independent.

What the answer must demonstrate: A run-level lock alone does not serialize different occurrences.

Follow-up · Question 8

The scheduling client edits a schedule while Tuesday is already queued. Which report should execute?

Reveal a model answer

I attach an immutable payload revision to the occurrence when it is materialized. Editing the schedule changes future unmaterialized instants, while r7 keeps its original parameters. Otherwise a retry with the same identity could perform different work.

What the answer must demonstrate: Separate schedule revision from occurrence identity and attempt identity.

Blank-page exercise · 45 minutes

Build the answer yourself

Design the scheduling client’s recurring report scheduler, then pause worker A after it starts and let worker B take over during a 09:00 burst.

  • Distinguish schedule, occurrence, attempt, and effect identity.
  • Calculate burst start lag and required capacity.
  • Trace an outbox handoff and a lease takeover.
  • Explain checked fencing and external idempotency limits.
  • Define DST, missed runs, overlap, and cancellation.

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 schedulerWhat stays the same across retries?Recall first, then reveal

The occurrence ID and logical external-action IDs; only the attempt identity changes.

One occurrence, many attempts.

Return to lesson
Design a distributed job schedulerA lease expires. Has the old worker stopped?Recall first, then reveal

Not necessarily. It may be paused; a destination must reject its stale token or deduplicate its effect.

Expiry is not a kill switch.

Return to lesson
Design a distributed job schedulerWhen does the last wave start: 50k jobs, 10k slots, 5s each?Recall first, then reveal

At 20 seconds; it completes at 25 seconds under the simplified model.

Start lag and completion time differ.

Return to lesson

Final revision

Summary and interview notes

A distributed scheduler durably materializes intended occurrences and treats retries as attempts of the same run. The database checks the attempt’s lease and fencing token before accepting its result. External effects need their own stable IDs and recovery rules.

Remember these points

  • Materialize the occurrence, advance the schedule and write dispatch intent in one authority transaction.
  • Calendar rules follow a clock time; fixed-rate runs follow an anchor; fixed-delay runs wait after completion.
  • A schedule-level active-run guard prevents logical overlap across different occurrences; a run lease alone does not.
  • Publish only a verified immutable object version under a current unexpired attempt, and emit delivery intent from that commit.
  • Queues preserve accepted work but cannot create worker capacity to satisfy a start deadline.

Interview tips

  • Calculate burst start waves separately from final completion time.
  • Pause a worker beyond expiry, complete a replacement and trace both object publication and email delivery.
  • State timezone, DST, misfire, overlap and cancellation policies before selecting a cron engine.

Important qualifications

  • A paused worker can resume while its replacement runs. The design protects accepted results and cooperating destinations, rather than promising that only one process is executing.
  • Result retention must honor active download pins, and an upload prefix alone does not make bytes immutable.

Technical references

Practice marks stay in this browser.