Design a distributed job scheduler: due times, leases and safe retries
Explain when an occurrence becomes durable, how workers claim it, and why a lease cannot make an arbitrary external side effect exactly once. Includes cron changes, missed runs, cancellation and tenant isolation.
1. Turn schedule a job into a precise contract
I would first separate deciding when work becomes eligible from executing the work. A scheduler can control dispatch, but cannot guarantee a partner API finishes at an exact wall-clock instant. I will ask which timing and completion guarantees the caller actually needs.
| Candidate asks | Illustrative interviewer reply | Design consequence |
|---|---|---|
| What can be scheduled? | One-time jobs and recurring cron jobs. | Persist schedule definitions separately from individual occurrences. |
| What does on time mean? | p99 dispatch within one second of due time under admitted load. | Measure due-to-dispatch lag; execution duration is a different metric. |
| How large and frequent? | 10 million active schedules, 10,000 occurrences per second average. | Shard due-time indexes and plan for bursts at minute boundaries. |
| Can work run twice? | Retries are allowed but duplicate business effects are undesirable. | Stable occurrence keys and idempotent targets are required; leases alone are insufficient. |
| What happens after downtime? | Run only the latest missed occurrence by default. | Explicit coalescing policy prevents unbounded catch-up storms. |
| How should overlapping runs behave? | One active run per schedule; skip or defer another run by policy. | Admission needs per-schedule concurrency state and a documented skip reason. |
| Which features can wait? | DAG dependencies and arbitrary uploaded code are out of scope. | Registered targets with scoped credentials; this is not a general workflow engine. |
2. Define dispatch, execution and cancellation separately
I acknowledge schedule creation only after durable storage. I do not acknowledge job success just because a queue accepted a message. The execution ledger records the distinction so users can tell pending dispatch from a target timeout.
| Requirement | Agreed target or scope | Design implication |
|---|---|---|
| Schedule lifecycle | Create, edit, pause, cancel; recurrence uses IANA timezone plus an explicit DST policy. | Revision number fences stale due-index records. |
| Execution | At-least-once attempts within retry age and attempt limits. | Target receives stable occurrence ID plus attempt and fencing token where supported. |
| Timing | p99 due-to-dispatch under one second in normal admitted load; no early dispatch. | Small look-ahead window, clock monitoring and a final due-time check. |
| Durability | Accepted definitions survive one availability-zone loss. | Replicated authoritative store; dispatch intents cannot be lost between ledger and queue. |
| Cancellation | No new attempts after cancellation is observed at claim; in-flight effects may already have occurred. | Recheck schedule revision at claim and support cooperative cancellation where targets allow it. |
| Isolation | Per-tenant rate and concurrency caps. | Separate heavy tenants and reserve capacity for time-sensitive work. |
3. Account for run history and concurrency
I use the interview workload below consistently. The number of stored schedules does not determine occurrences per second: one recurring schedule can generate many runs. History and executor concurrency dominate this example, so both need explicit retention and admission policies.
| Quantity | Calculation | Interpretation |
|---|---|---|
| Definitions | 10M × 2 KB metadata = 20 GB logical. | Large job payloads live in object storage with an immutable reference. |
| Occurrences | 10,000/s × 86,400 = 864M/day. | A 3x dispatch peak is 30,000/s; minute-boundary skew can exceed that unless admission smooths it. |
| History | 864M/day × 1 KB × 7 days = 6.048 TB logical. | Three copies = 18.144 TB before indexes; archive summaries rather than retaining every attempt forever. |
| Concurrent execution | 10,000/s × 10-second mean duration = 100,000 in-flight runs. | 300,000 at a sustained 3x peak if duration stays constant; use this to challenge the admitted workload. |
| Dispatch traffic | 10,000/s × 500-byte envelope = 5 MB/s average. | 15 MB/s at 3x before protocol and replication; do not put large payloads in every retry. |
| Due-index shards | 128 logical shards gives about 78 occurrences/s average per shard. | This is load distribution, not a proven shard limit; test hot minutes and skewed tenants. |
4. Separate definitions, occurrences and attempts
The stable identity of a recurring run is schedule ID plus schedule revision plus scheduled UTC instant. An occurrence can have several attempts but one business idempotency key. Edits must not reinterpret an already-dispatched occurrence as a new job with the same identity.
| Record and key | Fields / invariant | Access pattern |
|---|---|---|
| Schedule(tenantId, scheduleId) | revision, status, timezone, expression, DST policy, nextDueAt, targetRef, concurrency and missed-run policy. | Owner-authorized management writes with optimistic revision checks. |
| DueIndex(bucket, shard, dueAt, scheduleId) | Definition revision and next due instant. | Bounded range scans; index entries are hints and must be validated against the definition. |
| Occurrence(scheduleId, revision, dueAt) | runId, state, attemptCount, retryDeadline, leaseUntil, fencingEpoch. | Unique insert prevents two scanners creating two logical occurrences. |
| Attempt(runId, attemptNo) | start/end, target result, traceId and classified error. | Append audit trail; final run state changes only under the current fence. |
| Outbox(runId, dispatchGeneration) | Queue envelope and publication state in occurrence transaction. | Repeated publication is safe because workers claim the same run. |
5. Make the management and worker contracts reviewable
All management calls are tenant scoped. I would return an execution-history URL with a schedule so clients do not have to infer success from a queue acknowledgment. The worker APIs below can be internal RPCs rather than public endpoints.
| Operation | Contract | Failure or retry behavior |
|---|---|---|
| POST /schedules | Idempotency-Key; targetRef, payloadRef, cron/timezone or dueAt. Return 201 with revision. | Validate mutually exclusive one-time/cron fields, target permission and future schedule bounds. |
| PATCH /schedules/{id} | If-Match revision; update cron, pause or cancellation state. | 409 on concurrent edit; caller rereads instead of overwriting a newer schedule. |
| GET /schedules/{id}/runs | Cursor ordered by scheduled instant and runId. | Show QUEUED, RUNNING, RETRY_WAIT, SUCCEEDED, FAILED, SKIPPED or CANCELLED distinctly. |
| POST /runs/{id}/claim | Atomically claim an eligible run and increment fencingEpoch. | Reject current live lease, cancelled definition or exhausted retry budget. |
| POST /runs/{id}/complete | Current fencingEpoch, outcome and target reference. | Reject stale workers; a business effect outside the ledger still needs target deduplication. |
6. Draw durable dispatch before adding a timing wheel
My baseline uses a due-time index, a transactional occurrence ledger, an outbox relay and a durable execution queue. Pollers own logical shards through renewable leases, but unique occurrence insertion remains the correctness guard if a poller pauses and another takes over.
7. Follow one due occurrence into the queue
Consider a daily report due at 09:00. The occurrence identity includes that date and instant, not the time a poller happens to wake up. Two pollers may race, but only one durable logical run should be created.
Persist the definition
Commit the schedule revision and due-index entry together, or publish an index-update intent transactionally. A repair scan handles lagging projections; a heap in one process is not the source of truth.
Scan a bounded window
Query current and recent unfinished time buckets for the owned shards. Fetch slightly ahead into memory but recheck authoritative time and definition state before making work eligible.
Create the occurrence
In one transaction, insert the unique occurrence, advance nextDueAt for that revision and record dispatch intent. A conflicting occurrence means another scanner already succeeded.
Publish with retry
Relay sends runId to the queue, then marks the intent delivered. A crash between these operations may publish twice; both envelopes still reference one ledger record.
Handle recurrence changes
An old revision’s index entry is ignored. Apply the chosen DST behavior: for example skip nonexistent local times and run once at the first occurrence of a repeated local time. Store resolved UTC instants for audit.
8. Trace execution, lease renewal and uncertain success
A worker receives a queue message and attempts a conditional claim. Receipt itself is not ownership. If it gets a lease, it starts a heartbeat, invokes the registered target with runId and only acknowledges the queue after a durable ledger transition.
Bound concurrency
Acquire tenant and schedule permits before claiming expensive work. Release permits on terminal completion or fenced lease expiry; do not let a slow tenant occupy the entire worker fleet.
Protect against stale workers
Increment a fencing token on each claim. Completion updates require the current token. A downstream store must also reject stale tokens if we rely on fencing to prevent stale business writes.
Handle success with lost acknowledgment
If the target committed but the worker crashed, retry uses the same business idempotency key. A cooperating target returns the previous result; an arbitrary HTTP target may duplicate the effect.
Retry deliberately
Classify transient and permanent failures; use exponential backoff with jitter, deadline and maximum attempts. A sweeper reclaims expired leases through conditional state transitions, never a blind reset of all old RUNNING rows.
Present honest history
Distinguish target-confirmed success from dispatch and from unknown outcome. Cancellation can stop future attempts but cannot roll back a completed email or payment.
9. Scale due-time discovery without creating a global bottleneck
At 10,000 occurrences per second I partition time buckets by a stable schedule hash and rebalance logical shards independently of machines. I measure oldest-due age rather than only queue length; a small queue can contain very late work.
| Decision | Chosen baseline | Alternative and trade-off |
|---|---|---|
| Near-term discovery | Persistent bucket index plus in-memory look-ahead heap. | A timing wheel lowers timer overhead but must be rebuildable after process loss. |
| Cron bursts | Per-tenant admission and optional flexible windows agreed by the caller. | Silently jittering an exact schedule changes its contract; reject excess load instead. |
| Recurring overlap | Default one active run and coalesce missed occurrences. | Parallel or full catch-up policies cost more capacity and need explicit limits. |
| Multi-region | One active home region per schedule with replicated recovery data. | Cross-region failover requires fencing old dispatchers and stating replication RPO; active-active ownership is not free. |
10. Test the cases that fool a happy-path scheduler
I monitor due-to-dispatch p99, oldest retry age, stale completion rejections, lease churn, tenant throttling and unknown target outcomes. Correlate scheduleId, runId, attempt and traceId without logging secrets from payloads.
| Failure | Detection | Recovery and remaining limitation |
|---|---|---|
| Poller dies after transaction | Outbox age rises; occurrence exists without queue progress. | Another relay republishes; uniqueness prevents a second logical run. |
| Worker pauses beyond lease | Heartbeat stops; later completion carries an old fence. | New attempt may run; reject stale ledger updates and rely on target dedup/fencing for business correctness. |
| Target succeeds but response is lost | Timeout with no authoritative target result. | Retry idempotently or reconcile via target status; label unknown where the target lacks support. |
| Database unavailable | Claim/commit errors and due-age growth. | Stop new unrecorded dispatches, preserve queue backlog, recover with admission-controlled catch-up. |
| Clock or DST mistake | Early-run checks and schedule-preview mismatch. | Use synchronized time, explicit timezone/DST rules and tests around transitions; alert on skew. |
| Tenant floods minute boundary | Per-tenant queued age and permit saturation. | Enforce quotas and isolate queues; one tenant cannot consume all dispatch capacity. |
| Cancel races with execution | Definition revision changed after worker claim. | Cooperative cancellation when possible; report that an already-started side effect may finish. |
11. State the guarantees you can actually defend
I separate durable scheduling, at-least-once dispatch and target execution. Occurrence identity prevents duplicate logical runs; leases allocate work; fences reject stale writes; idempotency handles replayed business operations. Each solves a different problem.
For a managed implementation I would compare a scheduler service with this custom due-index design against required precision, quotas and target support. I would not claim a managed product inherits our illustrative one-second target without checking its documented contract.
Technical references
Primary references explain underlying mechanisms. Workloads and architecture choices above remain proposed interview assumptions.