Design a Distributed Job Scheduler and Task Queue
This page is one interview loop in three rounds. All three rounds design the same system. Each round opens with the interviewer raising the scope, and the design from the round before has to evolve to meet it.
| Round 1: Mid-level | Round 2: Senior | Round 3: Architect | |
|---|---|---|---|
| Story | An app sends reminder emails at chosen times and runs a few nightly jobs | Every team in the company schedules delayed, recurring and immediate tasks on one platform | A multi-tenant service that runs multi-step workflows across regions |
| Level (Amazon) | SDE II (L5) | Senior SDE (L6) | Principal (L7) |
| Volume | ~100K scheduled jobs; ~175K runs/day; 50 runs/s in the busiest minute | 100M active jobs; 500M runs/day (5,787/s average); ~36.8K runs/s for 5 minutes after midnight UTC | 5,000 tenants; 1B job runs + 500M workflow steps a day; 1.7B pending timers |
| Storage | ~3 GB | ~100 GB of active jobs, ~1.3 TB of recent runs; 15 TB/month of raw history to S3 | ~18.5 TB per home copy, ~48 TB with replicas |
| Footprint | 1 region, 3 AZs | 1 region, 3 AZs; survives losing an AZ | 3 regions; survives losing a region |
| Targets | Run within ~1 minute of the scheduled time; never lose a job; 99.9% | Dispatch P99 < 1 s after the target time; 99.99% | Per-tenant SLOs; region failover with duplicates bounded by replication lag |
| Reading time | ~35 min | ~40 min | ~45 min |
You can start at any round. Rounds 2 and 3 open with a "Where we left off" summary that catches you up.
Loop Opener: What Is a Job Scheduler?
You Already Know One: an Alarm Clock With a To-Do List
Your phone's alarm app does two things. It remembers when ("7:00 on weekdays"), and when the time comes it does something (plays a sound). A job scheduler is the same, for software:
| We say | The scheduler stores | At that time it |
|---|---|---|
| "At 9:00, send Alice her reminder" | a one-off job: one time, one action | calls the email service once |
| "Every night at 2:00, build the report" | a recurring job, written as a cron expression (a compact "every day at 2:00" rule) | calls the report service, then works out the next 2:00 |
| "Charge this card as soon as you can" | an immediate task: a job due right now | hands it to a worker at once |
The thing a job calls is its target: an HTTP endpoint, a queue, or a worker process. The job itself is small: when, what to call, with what payload, and what to do on failure.
What Makes It Hard
On one machine with one crontab, this is easy. In production, four things go wrong:
| On one machine | In a distributed scheduler |
|---|---|
| One clock, one list, one process | Many workers must agree who runs each job, without running it twice |
| If the machine dies, nothing runs (and we notice) | A worker can die halfway through a job. Did the email go out or not? |
| Few jobs | Millions of jobs, and people love round numbers: millions are due at exactly midnight |
| The job calls a target and waits | Targets get slow and fail. Retrying blindly makes them worse |
Every problem in the three rounds is one of three concerns:
- When: finding due jobs on time (timers, buckets, cron).
- Who: making sure one worker owns each run (claims, leases, fencing).
- What if it fails: retries, dead-letter queues, reapers.
The Question the Whole Loop Answers
A worker can crash after it called the target but before it wrote down "done". The next worker can't tell whether the call happened. So it has two choices: run the job again (maybe twice), or don't (maybe never). No design escapes this for side effects outside our own system. That's why "exactly once" is not something we can promise for an email or a payment.
How do we run every job at the right time, and make it safe when machines fail and a job runs twice?
The answer gets sharper every round:
- Round 1: at-least-once execution, plus an idempotency key the target uses to ignore repeats.
- Round 2: the same promise at 100,000 times the load, with fencing so a stale worker can't corrupt our own records.
- Round 3: a precise contract: exactly-once for the engine's own state transitions, at-least-once for everything it calls, across regions.
Round 1 · Mid-level · "Scheduled Emails and Nightly Jobs for One App"
~35 min · SDE II (L5) · 1 region, 3 AZs · ~100K scheduled jobs · 50 runs/s in the busiest minute · within ~1 min · 99.9%
R1.1 Establish Design Scope
The interviewer says: "Our app lets users set reminders, and we have a few nightly jobs. Design the scheduler." Before we draw anything, we ask questions and say what each answer changes.
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| What kinds of jobs? | One-off reminders at a chosen time, and recurring jobs: a daily digest per user, plus about 20 nightly system jobs. | Two schedule types: a timestamp, and a cron expression. After each recurring run we must compute the next one. |
| How precise must the time be? | Within a minute is fine. Never early. | We can poll every few seconds instead of keeping exact timers. "Never early" rules out running a job before its time to smooth load. |
| What does a job do? | It calls an internal HTTP endpoint, such as the email service, with a small JSON payload. | The target is an HTTP call we make and wait for. Its success or failure decides what we record. |
| How long do jobs run? | Seconds for reminders; a few minutes for nightly jobs. | A fixed claim time (a lease) of a few minutes is enough. No heartbeats yet. |
| What if a job fails? | Retry a few times. After that, someone should look at it. | Retries with backoff, then a dead-letter queue and an alarm. |
| Time zones? | The user's. "8:00 every day" means 8:00 where the user lives. | Each job stores an IANA time zone (like America/New_York), and daylight saving changes become our problem (step 1.5). |
Out of scope for this round:
- Priorities. Every job is equally important today.
- Very long jobs (hours). They'd need heartbeats and lease renewal.
- Workflows, where one job's output feeds the next. Each job stands alone.
The interviewer will widen this scope later. Write your out-of-scope list where you can see it: in a multi-round loop, some of it comes back.
R1.2 Functional Requirements, Derived Step by Step
| Phrase from the problem | Operation |
|---|---|
| "Users set reminders" | create(job) with a one-off time |
| "A daily digest", "nightly jobs" | create(job) with a cron expression and time zone |
| "At the chosen time, call the endpoint" | the scheduler runs the job at time T |
| "Retry a few times" | a per-job retry policy |
| "The user deleted the reminder" | cancel(job) |
| "Did my job run?" | get(job) returns its status and recent runs |
Not yet: priorities, long jobs with heartbeats, cancelling a run that's already running, and workflows.
R1.3 Non-Functional Requirements: the Questions
- Durability: never lose a job. Once we return "created", the job must run even if every machine restarts.
- Never early. A reminder at 9:00 must not fire at 8:59.
- A little late is fine. Within about a minute of the scheduled time.
- Avoid duplicates, but we can't rule them out. If a worker dies after calling the target and before recording success, we can't know whether the call worked. We'd rather run twice than never, so the target must be able to ignore a repeat. We give it an idempotency key: a value that is the same every time we retry the same run, so the target can recognize and skip a repeat.
- Availability: 99.9%. The API and the scheduler can be down about 8.8 hours a year (
0.1% × 8,760 h). Jobs due during an outage run late, not never.
R1.4 The API
Create a job
httpPOST /v1/jobs HTTP/1.1 Content-Type: application/json Idempotency-Key: 5d0c1f3a-8a7e-4b52-9a61-0e2f7c9b1d44 { "schedule": { "type": "CRON", "cron": "0 8 * * *", "time_zone": "America/New_York" }, "target": { "url": "https://email.internal/v1/digests", "method": "POST", "timeout_seconds": 30 }, "payload": { "user_id": 1001, "kind": "daily_digest" }, "retry_policy": { "max_attempts": 5, "initial_backoff_seconds": 10, "max_backoff_seconds": 600 }, "misfire_policy": "RUN_ONCE" }
json{ "job_id": "job_01J9ZK8V4T6Q", "status": "SCHEDULED", "next_run_at": "2026-09-28T12:00:00Z", "next_run_local": "2026-09-28T08:00:00-04:00" }
A one-off job uses "schedule": { "type": "ONCE", "run_at": "2026-09-28T09:00:00", "time_zone": "Europe/Berlin" }.
Read a job
httpGET /v1/jobs/job_01J9ZK8V4T6Q HTTP/1.1
json{ "job_id": "job_01J9ZK8V4T6Q", "status": "SCHEDULED", "next_run_at": "2026-09-28T12:00:00Z", "recent_runs": [ { "scheduled_for": "2026-09-27T12:00:00Z", "attempt": 1, "outcome": "SUCCEEDED", "finished_at": "2026-09-27T12:00:04Z" } ] }
Cancel a job
httpDELETE /v1/jobs/job_01J9ZK8V4T6Q HTTP/1.1
json{ "job_id": "job_01J9ZK8V4T6Q", "status": "CANCELLED" }
Why an Idempotency-Key on create? POST creates something new. If the network times out after we stored the job, the client retries, and without the key we'd store the reminder twice and send two emails. With the key, the server returns the job it already created.
What the target receives. Every call to the target carries Idempotency-Key: job_01J9ZK8V4T6Q:2026-09-28T12:00:00Z: the job ID plus the scheduled time of this run. It stays the same across retries of one run and changes for the next day's run.
Status codes
| Code | Meaning | What the client does |
|---|---|---|
201 Created | Stored | Nothing |
200 OK | Same idempotency key as before; here is the existing job | Nothing |
400 Bad Request | Bad cron expression, unknown time zone, payload too big | Fix the request |
404 Not Found | No such job | Treat as absent |
409 Conflict | The job is running right now and can't be edited | Retry the edit after the run |
429 Too Many Requests | Over the rate limit | Back off, retry |
Recap
- Two schedule types: a one-off time and a cron expression, both with a time zone.
- About 100K scheduled jobs; runs must start within a minute and never early.
- Jobs call internal HTTP endpoints for seconds to minutes.
- Retries, then a dead-letter queue. Duplicates are possible, so targets get an idempotency key.
Let's build it, starting with what most teams already have.
R1.5 Design Evolution: From One Crontab to a Worker Fleet
Every step follows the same pattern: a problem, your turn to think, the answer, and what the answer costs us.
Step 1.0: The Baseline
One server runs cron. Its crontab has the nightly jobs, and a small script reads reminders from the app database every minute and sends the due ones.
Synthesizing vector architecture diagram...
Everything lives on one machine: the schedule, the loop that finds due work, and the calls.
What's good: it's simple, and it works for a while. Everything below exists because this server will die, and nobody will notice until the reminders stop.
Step 1.1: The Cron Server Died and the Nightly Job Didn't Run
The problem: the cron server's disk failed at 1:00. The 2:00 report never ran, and 400 reminders were never sent. Nobody knew until users complained. What would you do? How do we make sure a single machine's death doesn't lose or stall jobs?
We use Amazon DynamoDB for the table. In DynamoDB an index must have a partition key, so we can't just "index next_run_at". We give each waiting job two extra attributes, timer_shard (one of D#0 … D#3, from a hash of the job ID) and timer_at (its next_run_at), and build a global secondary index (GSI, a second copy of selected attributes, keyed differently) called timer-index on them. A worker asks each of the 4 shards for items with timer_at <= now. The index is sparse: only items that have timer_shard appear in it, so it holds waiting jobs and nothing else.
Why 4 shards and not 1? One partition key value lives on one partition, and DynamoDB serves each partition at up to 1,000 write units and 3,000 read units per second. We're nowhere near that today; 4 shards is cheap insurance, and Round 2 shows why the count must grow with the load.
Step 1.2: Two Workers Picked the Same Job
The problem: at 8:00, three workers poll at almost the same moment. All three see the digest for user 1001 as due. All three send it. The user gets three emails. What would you do? How does exactly one worker take each due job?
The claim, as the JSON request the worker sends to DynamoDB's UpdateItem:
json{ "TableName": "jobs", "Key": { "job_id": { "S": "job_01J9ZK8V4T6Q" } }, "ConditionExpression": "#st = :scheduled AND next_run_at = :seen_due", "UpdateExpression": "SET #st = :claimed, lease_owner = :me, lease_expires_at = :lease_end, timer_shard = :lease_shard, timer_at = :lease_end, attempt = attempt + :one", "ExpressionAttributeNames": { "#st": "status" }, "ExpressionAttributeValues": { ":scheduled": { "S": "SCHEDULED" }, ":claimed": { "S": "CLAIMED" }, ":seen_due": { "N": "1790596800000" }, ":me": { "S": "worker-az-b-7f3c" }, ":lease_end": { "N": "1790597130000" }, ":lease_shard": { "S": "L#2" }, ":one": { "N": "1" } }, "ReturnValues": "ALL_NEW" }
Three details:
next_run_at = :seen_duemakes the claim fail if the job was rescheduled or cancelled after we read it from the index. The index is a copy that's updated a moment after the table, so what a worker sees there can be stale. The conditional claim is what makes stale reads harmless.- The claim moves the job in the index.
timer_shardchanges fromD#2("waiting to run") toL#2("claimed, lease ends attimer_at"). The same index now answers a second question: "which claims have expired?" (step 1.3). One index, one rule: every job has exactly one next timer. :lease_endcomes from the worker's clock:now + timeout (30 s) + 5 minutes. DynamoDB's condition language has no "current server time". We come back to that in R1.9.
ReturnValues: ALL_NEW hands the worker the whole job (target, payload, retry policy) in the same call, so it doesn't need a second read.
Synthesizing vector architecture diagram...
Both workers see the same due job. The database applies the conditional updates one at a time, so the second one finds status CLAIMED and fails. To waste fewer claims, each worker shuffles its list of due jobs before claiming.
Primitive: Distributed Locks & Leases
Step 1.3: A Worker Died Holding a Job
The problem: a worker claimed the 2:00 report, then its Fargate task was stopped for a deploy. The job is CLAIMED by a worker that no longer exists.
What would you do? How does the job get run, and what new risk does that create?
Synthesizing vector architecture diagram...
A job is always in one state, and every arrow is a conditional update. A cron job whose attempts are exhausted records a failed run, goes to the dead-letter queue, and is re-armed for its next occurrence.
How a target makes itself idempotent. The email service keeps a table of idempotency keys it has handled, with a unique constraint. On a request, it inserts the key first; if the insert fails because the key exists, it returns the stored result and sends nothing. The key must be stored in the same transaction as the side effect, or a crash between them brings the problem back.
Step 1.4: The Target Was Down
The problem: the email service is down for 10 minutes during a deploy. Every reminder due in that window fails. What would you do? When and how do we try again?
Only some failures are worth retrying:
| Target response | Retry? | Why |
|---|---|---|
Timeout, connection refused, 502, 503, 504 | Yes | The target may recover |
429 Too Many Requests | Yes, not before Retry-After | The target asked us to slow down |
400, 404, 422 | No: straight to the DLQ | The request itself is wrong; sending it again won't fix it |
409 with "already processed" | No: count it as success | The target recognized the idempotency key |
Primitive: Circuit Breaker, Bulkhead & Fault Tolerance
Step 1.5: A Daily Job Was Skipped in March and Ran Twice in November
The problem: one user's digest is set to 30 2 * * * in America/New_York. On the night the clocks jump forward in March, 2:30 doesn't exist, and the digest never ran. Another user's is 30 1 * * *. On the night the clocks fall back in November, the hour from 1:00 to 1:59 happens twice, and that digest ran twice. Support got complaints both nights.
What would you do? How do we store and compute recurring times?
Cron details that trip people up:
| Detail | What to know |
|---|---|
| Five fields vs six | Classic Unix cron has 5 fields: minute, hour, day of month, month, day of week. Quartz puts seconds first (6 or 7 fields). EventBridge Scheduler uses 6 fields with year last, and requires ? in either day-of-month or day-of-week. 0 8 * * * in one dialect is invalid or means something else in another, so our API accepts exactly one dialect (5-field Unix) and rejects the rest. |
| Day of month AND day of week | In classic Unix cron, if both are restricted (0 8 1 * 1), the job runs when either matches: the 1st of the month and every Monday. Most people expect "both". We warn at creation. |
| Every 15 minutes in a local zone | */15 * * * * in New York, evaluated in local time, runs four extra times in the repeated hour or loses runs in the skipped one, depending on the library. For sub-daily schedules we recommend UTC. |
| Time-zone rules change | Governments change DST rules, sometimes at short notice. When we deploy new IANA data, we recompute next_run_at for jobs in the affected zones. |
Trace: re-arming a daily job on the fall-back night (the night of Sunday, November 1, 2026 in New York)
text1. Job: cron "30 1 * * *", zone America/New_York, dialect Unix 5-field. 2. Run scheduled for 2026-10-31 01:30 EDT = 2026-10-31T05:30Z succeeds. 3. cron_next(after 2026-10-31T05:30Z) finds 01:30 on 2026-11-01. 4. 01:30 occurs twice that night: 01:30 EDT (05:30Z) and 01:30 EST (06:30Z). 5. Policy "run once, first occurrence" picks 05:30Z. next_run_at = 2026-11-01T05:30Z. 6. After that run, cron_next(after 2026-11-01T05:30Z) must return 2026-11-02 01:30 EST = 06:30Z, not the second 01:30 of the same night (06:30Z on Nov 1). The library's "after" logic, plus the policy, decide this. We test both DST nights for every zone we support.
Step 1.6: The Scheduler Was Down for an Hour: Run the Missed Jobs?
The problem: a bad deploy stopped every worker from 03:00 to 04:00. When they come back, 12,000 jobs are overdue, including an hourly cleanup job that missed its 03:00 run and is due at 04:00 too. What would you do? Which overdue runs happen, and how many times?
Round 1 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 1.0 | (baseline) | One cron server | Everything below |
| 1.1 | The cron server died | Jobs table with next_run_at; stateless workers polling a sparse index | Polling load; lateness up to one poll |
| 1.2 | Two workers took one job | Conditional claim with a lease | A lease to manage |
| 1.3 | A worker died holding a job | Lease expiry + reaper; idempotency keys | Possible duplicates |
| 1.4 | The target was down | Backoff with full jitter, max attempts, SQS DLQ + alarm | Late completion |
| 1.5 | DST ran a job twice or never | Cron + IANA zone, next run computed in zone, stored in UTC; gap and overlap policies | Time-zone data to maintain |
| 1.6 | Missed runs after an outage | Misfire policy per job | Semantics to define per job |
R1.6 Architecture v1
Synthesizing vector architecture diagram...
The API only writes jobs. Workers do everything else: find due jobs, claim them, call targets, record runs, re-arm cron jobs and reap expired claims. The table is the only shared state, so any worker can die at any moment.
Tables
jobs (partition key job_id)
| Attribute | Example | Notes |
|---|---|---|
job_id | job_01J9ZK8V4T6Q | Derived from the idempotency key on create, so a retried create hits the same item and attribute_not_exists(job_id) rejects the duplicate |
status | SCHEDULED | SCHEDULED, CLAIMED, COMPLETED, FAILED, CANCELLED |
schedule | { "type": "CRON", "cron": "0 8 * * *", "time_zone": "America/New_York" } | |
next_run_at | 1790596800000 | Number, UTC epoch milliseconds |
timer_shard, timer_at | D#2, 1790596800000 | Keys of the sparse timer-index GSI; L#2 and the lease end while claimed; removed when the job is finished |
lease_owner, lease_expires_at | worker-az-b-7f3c, 1790597130000 | Set by the claim |
attempt | 1 | Attempts of the current run |
target, payload, retry_policy, misfire_policy | Kept under ~1 KB so one write costs one write unit |
timer_at must be a Number. As a String, "999" sorts after "1000", and the index would return jobs in the wrong order.
job_runs (partition key job_id, sort key scheduled_for#attempt): one item per attempt with the outcome, the HTTP status and timings. A 30-day TTL cleans it up. TTL deletion is a background process that typically removes an item within a few days of its expiry time, so the table holds a little more than 30 days, and reads filter out expired items.
Trace 1: a one-off reminder
- At 08:40 the app creates a job for 09:00 Berlin time. The API converts it to
07:00:00Z, puts the item withtimer_shard = D#1, and returns201. - At 07:00:03Z, worker AZ-a's poll of
D#1returns the job. It claims it (step 1.2). - It calls the email service with
Idempotency-Key: job_…:2026-09-28T07:00:00Z, gets200, writes ajob_runsitem, and sets the jobCOMPLETEDwithtimer_shardremoved, conditional onlease_ownerbeing itself and status stillCLAIMED.
Trace 2: a cron re-arm
- The digest for 08:00 New York (12:00Z) succeeds.
- The worker computes the next occurrence from
scheduled_for = 12:00Z: 08:00 tomorrow New York = 12:00Z tomorrow. - One conditional update sets
status = SCHEDULED,next_run_atandtimer_atto tomorrow,timer_shardback toD#…, and clears the lease.
Trace 3: a worker crash
- Worker AZ-c claims the 02:00 report; its lease ends at 02:05:30.
- The task is killed at 02:01.
- At 02:05:47 (the next 30-second reaper pass after the lease end plus a 5-second margin), worker AZ-a finds it in
L#…, returns it toSCHEDULEDdue now, and claims it on its next poll. - The report runs again from the start with the same idempotency key. If the report service had already finished the first attempt's work under that key, it returns the stored result.
R1.7 Numbers
Targets
| Quality | Target | Why this number |
|---|---|---|
| Timeliness | Start within 60 s of the scheduled time; never early | From R1.1 |
| Durability | A created job always runs or ends in the DLQ | From R1.3 |
| Availability | 99.9% | About 8.8 hours a year; late jobs still run |
Jobs and runs (assumptions: the app has 1M users; these are our estimates, not given numbers)
| Item | Math | Result |
|---|---|---|
| One-off reminders | 150K created a day; each waits 12 h on average. Jobs waiting = arrival rate × wait (Little's law) = 150K × 0.5 day | 75K waiting |
| Daily digests | 25K users opted in; one cron job each | 25K jobs, 25K runs/day |
| Nightly system jobs | Report, cleanup, exports | ~20 |
| Active jobs | 75K + 25K + 20 | ≈ 100K |
| Runs per day | 150K + 25K + 20 | ≈ 175K |
| Average rate | 175K ÷ 86,400 s | ≈ 2 runs/s |
The busiest minute. Averages hide clustering. Users keep the default digest time (08:00) about 20% of the time (assumption), and half our users live in one time zone: 25K × 20% × 50% = 2,500 digests due at 08:00 in that zone. Add about 500 reminders people set for a round time like 09:00 in the same minute (assumption): 3,000 jobs due in one minute.
- If all 3,000 had to start in the first second: 3,000 runs/s.
- Our target is "within 60 s", so we can spread them over the minute:
3,000 ÷ 60 s = 50 runs/s. That's the number we size for.
Workers
| Item | Math | Result |
|---|---|---|
| Calls in flight at peak | 50 runs/s × 2 s average call (assumption) | 100 |
| Per worker | 50 concurrent calls (0.5 vCPU, 1 GB; I/O-bound) | 3 workers = 150 slots, one per AZ |
| Losing one worker | 2 × 50 = 100 slots | Still covers 100 in flight |
| Lateness | ≤ 5 s poll + claim + queueing behind 100 calls | well under 60 s |
Polling load
| Item | Math | Result |
|---|---|---|
| Due polls | 3 workers × 4 D# shards ÷ 5 s | 2.4 queries/s |
| Reaper polls | 3 workers × 4 L# shards ÷ 30 s | 0.4 queries/s |
| Per month | 2.8/s × 2.59M s | ≈ 7.3M queries |
Writes per run (DynamoDB bills one write unit per 1 KB written; changing a GSI key costs a delete plus a put in the index)
| Write | Units |
|---|---|
Claim: table + index move D#→L# | 1 + 2 |
Finish: table + index move back to D# (cron) or delete (one-off) | 1 + 2 (or 1 + 1) |
job_runs item | 1 |
| Create (one-offs only): table + index | 1 + 1 |
That's 7 per cron run (claim 3, finish 3, run item 1) and 8 per one-off job over its life (create 2, claim 3, finish 2, run item 1): 150K × 8 + 25K × 7 = 1.375M writes a day, about 41M a month.
Storage: jobs about 100K × 1 KB = 100 MB. job_runs: 175K a day × 30 days × ~0.5 KB = 2.6 GB, plus the few days TTL may lag, so plan for ~3 GB.
Rough monthly cost (us-east-1 list prices; check the AWS Pricing Calculator before quoting)
| Line | Math | ≈ Monthly |
|---|---|---|
| DynamoDB writes (on-demand) | 41.25M × $0.625 per million | $26 |
| DynamoDB reads | 7.3M polls × 0.5 read unit + ~30M status reads × 0.5 (1M GETs/day, assumed) ≈ 19M units × $0.125 per million | $2 |
| DynamoDB storage | ~3 GB × $0.25 | $1 |
| Workers | 3 Fargate Graviton tasks (0.5 vCPU, 1 GB) × $0.0198/h × 730 h | $43 |
| API | 2 tasks of the same size | $29 |
| Application Load Balancer | $16 base + about 1 capacity unit | $22 |
| CloudWatch, SQS, KMS | logs, a few alarms, a near-empty DLQ | $22 |
| Total | ≈ $145 |
The fleet costs more than the database. At this size, the scheduler is almost free to run; what costs is building and owning it.
R1.8 Trade-Offs
Polling a table vs a delay queue
| Poll a jobs table (chosen) | SQS delay queue | Buy: EventBridge Scheduler | |
|---|---|---|---|
| How far ahead | Unlimited | At most 15 minutes (DelaySeconds max 900) | Unlimited; one-time, rate and cron schedules with time zones |
| Cancel / edit | Update the item | Can't find or remove one delayed message | UpdateSchedule / DeleteSchedule |
| Status and history | Ours | None | Invocation metrics; outcomes are the target's business |
| Precision | Poll interval (5 s) | Seconds | 60-second precision: a job for 1:00 runs between 1:00:00 and 1:00:59 |
| Retries, DLQ | Ours | Ours | Built in: up to 185 retries or 24 h max event age, then an SQS DLQ |
| Default quotas | Ours | – | 10M schedules per Region; 1,000 invocations/s; 5,000 CreateSchedule/s in large Regions (all adjustable) |
A delay queue can't hold a reminder set for next week without re-queuing it every 15 minutes, and we couldn't cancel it. It's useful for short retry delays (Round 2), not as the schedule.
Relational vs key-value store
| DynamoDB (chosen) | PostgreSQL on Amazon RDS | |
|---|---|---|
| Claim | Conditional UpdateItem | SELECT … FOR UPDATE SKIP LOCKED in a transaction: workers skip rows another worker has locked |
| "Now" | The worker's clock (no server time in conditions) | now() inside the query: one clock for everyone |
| Due index | Sparse GSI with shards | A partial index on next_run_at WHERE status = 'SCHEDULED' |
| Operations | None to speak of | Patching, failover, connection limits |
| Scale ceiling | Very high, with sharding by us | One writer; fine far beyond Round 1 |
At 100K jobs, Postgres is an equally good choice, and now() in the query makes lease math simpler. We pick DynamoDB because it needs no operating and because Round 2's scale will need it anyway. Either answer is fine in an interview if you name what the other one gives you.
Build or buy? Honestly, at this scale: buy. The app would create one EventBridge Scheduler schedule per reminder (one-time schedules with ActionAfterCompletion: DELETE, since completed one-time schedules still count against the quota until deleted) and one per digest (cron with the user's time zone). 175K runs a day is about 5.3M invocations a month, inside the 14M free invocations a month; after that it's $1.00 per million. The schedule's target would be an SQS queue or a Lambda function that calls the internal endpoint. What we'd give up: 60-second precision (fine here), status history (we'd log it ourselves), and our own misfire policies. The design above is what we'd build if the interviewer says "no managed scheduler", and it's the foundation for Rounds 2 and 3.
R1.9 Failure Modes
| Trigger | What you'd see | How the design responds |
|---|---|---|
| A worker crashes | One fewer poller; its claims stop making progress. | Other workers keep polling. Its claims expire and are reaped (step 1.3). ECS replaces the task. |
| DynamoDB errors or throttles | Polls and claims fail; jobs run late. | Workers back off with jitter and keep trying. Nothing is lost: jobs stay in the table, and the misfire policy (step 1.6) handles lateness. With Postgres instead, a Multi-AZ failover typically takes one to two minutes, with the same effect: jobs run late. |
| A target outage | Errors from one target; retries pile up. | Backoff with jitter spreads retries (step 1.4); after 5 attempts the DLQ alarm fires. Jobs for other targets are unaffected, because each worker call has a timeout. |
| Clock skew on a worker | A worker whose clock runs 10 s fast claims jobs 10 s early, or reaps a live lease early. | DynamoDB conditions can't read the server's time, so every time in a condition is some worker's clock. Workers sync with the Amazon Time Sync Service, which normally keeps skew far below a second, but a paused VM or broken sync can be seconds off. So: (1) the reaper only reclaims leases that ended more than 5 s ago by its clock; (2) a worker refuses to claim a job whose next_run_at is still in the future by its own clock, and we alarm if a worker's clock offset exceeds 1 s. With Postgres, now() in the claim query uses one clock, the database's, and this whole row disappears. |
| A poison job | One job crashes the worker process every time it's claimed. | The attempt counter increases on every claim, so after 5 attempts it goes to the DLQ instead of crashing workers forever. |
R1.10 Pillar Check
| Pillar | What Round 1 covers |
|---|---|
| Reliability | Stateless workers in 3 AZs; conditional claims; leases and a reaper; retries with jitter and a DLQ; misfire policies. REL 5 · REL 11 |
| Performance Efficiency | A sparse index that holds only waiting jobs; the busiest minute derived (50 runs/s) and sized for. PERF 3 |
| Security | Callers authenticate with IAM at the API; workers can only call targets on an allow list; payloads encrypted at rest with KMS; TLS to targets. SEC 2 · SEC 8 · SEC 9 |
| Cost Optimization | About $145 a month; an honest "buy EventBridge Scheduler" answer at this scale. COST 5 |
| Operational Excellence | Light this round: alarms on DLQ depth and on "oldest due job older than 60 s". OPS 8 |
| Sustainability | Skipped this round: the fleet is three small tasks. |
R1.11 Round 1 Rubric and Follow-Ups
What a strong mid-level (L5) answer shows
- Asks about precision, duration, failure handling and time zones before designing.
- Moves the schedule out of one machine into a shared table with stateless workers.
- Claims with a conditional write and explains why read-then-write fails.
- Says plainly that a crashed worker means a possible duplicate, and gives targets an idempotency key.
- Handles DST gaps and overlaps with a stated policy, and knows cron dialects differ.
- Derives the busiest minute instead of using an average, and names the managed alternative.
Follow-up questions
-
"Why not have one leader worker that assigns jobs to the others?" Answer: it works, but we'd need leader election, and the leader becomes a bottleneck and a single point of failure until a new one is elected. Conditional claims let every worker take work directly, with the database as the referee. At Round 1's 50 runs/s the claim conflicts cost almost nothing.
-
"A user edits a reminder from 9:00 to 9:30 while a worker is claiming the 9:00 run. What happens?" Answer: the edit sets
next_run_atto 9:30. If the edit lands first, the claim's conditionnext_run_at = 9:00fails and nothing runs at 9:00. If the claim lands first, the 9:00 run goes ahead and the edit, which is also conditional on the job beingSCHEDULED, gets a409telling the user the reminder is running now. -
"Why compute the next cron time from the scheduled time and not from now?" Answer: if a 02:00 job starts at 02:00:40 and we compute from now, a job like
*/1 * * * *drifts, and a slow run could make us skip an occurrence. From the scheduled time, every occurrence is exactly where the cron expression says; lateness is handled by the misfire policy, separately.
Interview gotchas from this round's wrong answers
| Gotcha | Why it's wrong |
|---|---|
| "Two cron servers for redundancy" | Every job runs twice. |
| "Read the status, then update it" | Two workers read the same value; both win. |
| "The lease guarantees it runs once" | A lease only says who may work on it; a crash after the side effect still causes a repeat. |
| "Convert the cron to UTC once" | Wrong for half the year wherever DST applies. |
| "SQS delay queue as the scheduler" | 15 minutes maximum, and no cancel. |
| "Retry immediately" | Hammers a recovering target, in synchronized waves. |
Round 2 · Senior · "100M Active Jobs and the Top-of-the-Hour Stampede"
~40 min · Senior SDE (L6) · 1 region, 3 AZs · 100M active jobs · 500M runs/day: 5,787/s average, ~36.8K/s for 5 minutes after midnight UTC · dispatch P99 < 1 s · 99.99%
R2.0 Where We Left Off
This is what the candidate says aloud in the first 60 seconds of Round 2. If you're starting here, it's everything you need from Round 1.
Round 1 in 60 seconds. "We built a scheduler for one app: about 100K jobs, one-off reminders and cron jobs with time zones, 175K runs a day and 50 runs a second in the busiest minute. Jobs live in a DynamoDB table with
next_run_atin UTC. Three stateless Fargate workers poll a sparse, sharded index for due jobs and claim each one with a conditional update that sets a lease. If a worker dies, the lease expires and a reaper returns the job, so a job can run twice; every call carries an idempotency key so the target can ignore repeats. Failures retry with exponential backoff and full jitter, then go to an SQS dead-letter queue. The next cron time is computed in the job's time zone from the scheduled time, with explicit rules for DST gaps and overlaps, and a misfire policy says what to do with runs missed during an outage. About $145 a month; honestly, at that size we'd buy EventBridge Scheduler. Open costs: polling one index doesn't scale, and nothing stops one slow target or one crowded minute from delaying everyone."
Architecture v1, compact
Synthesizing vector architecture diagram...
Round 1 in one picture: the table is the only shared state; workers find, claim, run, retry and reap.
Round 1 step summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 1.1 | The cron server died | Jobs table + stateless polling workers | Polling load; lateness up to one poll |
| 1.2 | Two workers took one job | Conditional claim with a lease | A lease to manage |
| 1.3 | A worker died holding a job | Lease expiry + reaper; idempotency keys | Possible duplicates |
| 1.4 | The target was down | Backoff with full jitter, max attempts, DLQ | Late completion |
| 1.5 | DST ran a job twice or never | Cron + IANA zone, next run stored in UTC | Time-zone data to maintain |
| 1.6 | Missed runs after an outage | Misfire policy per job | Semantics per job |
Open costs: every worker polls the same few index shards; all jobs share one pool of workers; a lease is checked only when the job finishes.
R2.1 The Scope Raise
Interviewer: "Every team in the company wants your scheduler. We're at 100 million active jobs and 500 million runs a day. The payments team schedules retries and settlement jobs; the analytics teams schedule thousands of reports. Payments must never wait behind reports."
Interviewer: "Some jobs now run for hours. Most targets are team webhooks, and some of them slow down badly under load. We want a job dispatched within a second of its time, P99. And it has to survive losing an AZ: 99.99%."
A scope raise is not the end of scoping. We ask back, and say what each answer changes.
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| How are the 500M daily runs spread over the day? | About 8M jobs are 0 * * * * (hourly at minute 0). About 10M are daily at minute 0: 2M at midnight UTC, the rest at local midnight across time zones. Everything else is spread out. | Midnight UTC has 10M jobs due at the same instant. We derive the peak from that, not from an average (R2.6), and we need time buckets (step 2.1) and a way to spread the burst (step 2.3). |
| Does every job need its exact second? | No. Payments do. Most hourly jobs are fine running a few minutes late, if we tell the owners. | A dispatch window per job: exact-time jobs are a limited, reserved resource; the rest spread over their window (step 2.3). |
| What does "payments never wait behind reports" mean in practice? | A settlement job must start on time even while 10M report jobs are queued. | Separate queues and worker pools per priority, with reserved capacity for the critical class (step 2.2). |
| How long do the long jobs run, and where? | Minutes to several hours. Those teams run their own worker processes; they just want us to hand them the work and track it. | A pull API with heartbeats and lease renewal (step 2.5). A 12-hour job can't be tracked by SQS visibility alone: 12 hours is its maximum. |
| Is running a job twice acceptable? | Only if the target recognizes the repeat. And our records must never mark a run "succeeded" because of a worker that had already lost it. | Idempotency keys stay; fencing tokens protect our own records from stale workers (step 2.4). |
| What happens when a target slows down? | Today we make it worse: we keep sending at full rate, and the team's service falls over. | Per-target concurrency limits that adapt (AIMD), plus a circuit breaker (step 2.6). |
| Operators need anything new? | Cancel a job while it runs, and re-run everything in the dead-letter queue once a bug is fixed. | Cancel-while-running through the heartbeat; a redrive API (R2.3). |
Scope change
| Round 1 | Round 2 | |
|---|---|---|
| Active jobs | ~100K | 100M |
| Runs | 175K/day | 500M/day: 5,787/s average |
| Peak | 50 runs/s (spread over a minute) | 10M due at 00:00:00 UTC; ~36.8K runs/s for 5 minutes with windows |
| Timeliness | within 60 s | dispatch P99 < 1 s after the (effective) target time |
| Job length | seconds to minutes | up to hours, with heartbeats |
| Priorities | none | CRITICAL, HIGH, DEFAULT, LOW, isolated |
| Targets | two internal services | thousands of team webhooks, some fragile |
| Survives | a worker | an AZ; 99.99% (about 52.6 min a year) |
R2.2 What Breaks in the Round 1 Design
| Round 1 choice | What breaks at the new scope |
|---|---|
| Poll 4 index shards every 5 s | 500M timers a day in 4 partition-key values. Each partition serves up to 1,000 write units a second, and at midnight tens of thousands of jobs a second move in and out of the index. A 5-second poll also can't meet P99 < 1 s. |
| Everything due at once is claimed at once | 10M jobs due at 00:00:00. "Dispatch within 1 s" would mean 10M claims and 10M target calls in one second. |
| One pool of workers for everything | Settlement jobs queue behind report jobs at midnight. |
| Workers call targets from the polling loop | One slow target ties up every worker's slots, and due jobs wait behind it. |
| A lease checked only at completion | A worker that paused for 60 s wakes up, finishes, and writes "succeeded" over a newer attempt's record. Nothing stops a stale heartbeat either, because there are none. |
| Leases sized to the job's timeout | A 6-hour job would need a 6-hour lease: if its worker dies at minute 1, nobody notices for 6 hours. |
| Full-rate calls to every target | A fragile webhook gets 10K calls in a minute and falls over; our retries then hammer it harder. |
The order we fix it in: finding due jobs (2.1), isolating priorities (2.2), the midnight burst (2.3), stale workers (2.4), long jobs (2.5) and slow targets (2.6).
R2.3 New Requirements and API Additions
The job gains a priority and a dispatch window.
httpPOST /v1/jobs HTTP/1.1 Content-Type: application/json Idempotency-Key: 7b2e0c55-1a9d-4f3e-8f0a-3d6c2b1e9a70 { "schedule": { "type": "CRON", "cron": "0 * * * *", "time_zone": "UTC" }, "priority": "DEFAULT", "dispatch_window_seconds": 300, "target": { "type": "HTTP", "target_id": "tgt_reports_rollup", "path": "/v1/rollups", "timeout_seconds": 60 }, "payload": { "report": "hourly_sales" }, "retry_policy": { "max_attempts": 5, "initial_backoff_seconds": 10, "max_backoff_seconds": 900 } }
dispatch_window_seconds: the job may start anywhere in[scheduled time, scheduled time + window). The default for new jobs is 300.0means strict: it needs a CRITICAL allowance (step 2.2).max_backoff_secondsis capped at 900 so every retry fits in an SQS message delay (15 minutes maximum).- The target is a registered target (
target_id), not a free URL. Registering it sets its owner, its base URL and its concurrency limits:
json{ "target_id": "tgt_reports_rollup", "owner_team": "analytics-platform", "base_url": "https://rollups.analytics.internal", "max_concurrency": 200, "initial_concurrency": 50, "latency_slo_ms": 2000, "auth": { "type": "HMAC_SHA256", "secret_arn": "arn:aws:secretsmanager:us-east-1:111122223333:secret:sched-tgt-rollups" } }
Long jobs run on the team's workers, which pull work.
httpPOST /v1/task-queues/ledger-rebuild/tasks:poll HTTP/1.1 Content-Type: application/json { "worker_id": "ledger-worker-12", "wait_seconds": 20 }
json{ "run_id": "job_7Q2M:2026-09-28T00:00:00Z", "epoch": 3, "lease_expires_at": "2026-09-28T00:00:31Z", "heartbeat_every_seconds": 10, "payload": { "ledger": "EU-2" }, "idempotency_key": "job_7Q2M:2026-09-28T00:00:00Z" }
httpPOST /v1/runs/job_7Q2M:2026-09-28T00:00:00Z:heartbeat HTTP/1.1 Content-Type: application/json { "epoch": 3, "progress": { "rows_done": 1200000 } }
json{ "lease_expires_at": "2026-09-28T00:07:41Z", "cancel_requested": false }
A heartbeat with a stale epoch gets 409 Conflict with "reason": "EPOCH_SUPERSEDED", and the worker must stop at once (step 2.4). Completion is POST /v1/runs/{run_id}:complete with the epoch and an outcome.
Cancel while running. POST /v1/runs/{run_id}:cancel sets cancel_requested. A pulling worker learns it from its next heartbeat response. An HTTP call already in flight can't be recalled; if it succeeds anyway, the run is recorded as SUCCEEDED_AFTER_CANCEL, so nobody is misled.
Redrive. After the cause is fixed:
httpPOST /v1/dlq:redrive HTTP/1.1 Content-Type: application/json { "filter": { "target_id": "tgt_reports_rollup", "failed_after": "2026-09-28T00:00:00Z" }, "max_per_second": 200 }
Redriven runs start a new attempt series under a new epoch, at a capped rate so a recovered target isn't flattened again.
R2.4 Design Evolution: Surviving the Stampede
Step 2.1: Scanning 100M Jobs for Due Ones Is Too Slow
The problem: 100M jobs, 500M runs a day, and we must notice each one within a fraction of a second of its time. Round 1's four index shards and 5-second poll can't do it. What would you do? How do we find "what is due right now" cheaply and precisely?
Why two stores can disagree safely. The bucket is a hint; the conditional claim in DynamoDB decides. A member that's stale because its job was cancelled or rescheduled after loading fails the claim's condition (next_run_at = :scheduled) and is dropped. So:
- Cancel updates the table (the truth) and removes the member as a courtesy. If the loader re-adds it a moment later from a stale index read, the claim still fails. Correctness never depends on the bucket being clean.
- Missing members are the dangerous direction: a job in the table but not in any bucket never runs. Three layers stop that:
- The loader writes a marker key
loaded:{017}:202609280000after it finishes a shard-minute. Valkey deletes a sorted set when its last member is removed, so an empty bucket and a never-loaded bucket look the same; the marker tells them apart. When a minute arrives without its marker (a failover lost it), the dispatcher reads that shard-minute straight fromdue-index. - One minute after each shard-minute's windows close, the dispatcher runs a straggler query on
due-indexfor it. Claimed jobs have moved to their next minute's key, so anything still there was missed, and it's dispatched late. That's256 × 1,440 = 368,640small queries a day. - Direct adds by writers cover jobs created after the loader passed their minute.
- The loader writes a marker key
The {017} in the key is a hash tag: Valkey places a key by the part in braces, so all of shard 17's buckets and markers live on the same node, and one dispatcher talks to one node for them.
Sorted-set scores are doubles, exact for integers up to 2^53 ≈ 9.0 × 10^15. Epoch milliseconds are about 1.79 × 10^12, far below that, so scores are exact.
Who owns which shard? Each dispatcher holds a lease on its shards in a small DynamoDB table (one item per shard: owner, epoch, lease end), renewed every 3 seconds with a 10-second lease. If a dispatcher dies, its shards are claimed by others within about 10 seconds. Two dispatchers briefly owning one shard would send some members twice; the conditional claim turns the second copy into a no-op.
Synthesizing vector architecture diagram...
The table is truth, the buckets are a 2-hour cache of it, and the dotted path is how the dispatcher recovers anything the cache lost.
Primitive: Distributed Cache Patterns & Eviction · Database Sharding & Partition Keys
Step 2.2: Payments Wait Behind Reports
The problem: at midnight, millions of report jobs are queued. A settlement job due at 00:00:00 sits behind them and starts at 00:04. The payments team pages us. What would you do? How does critical work stay on time no matter what else is queued?
Which jobs are CRITICAL? Only strict jobs (window 0) from teams with a CRITICAL allowance. The allowance caps how many strict jobs may be due in the same second, platform-wide: 5,000. The API counts them per due second and rejects a strict job that would exceed its team's share (409 STRICT_ALLOWANCE_EXCEEDED, with the nearest second that has room). That cap is what makes "P99 < 1 s" provable for strict jobs: a pool sized for 5,000 runs a second can't be handed 50,000.
Primitive: Circuit Breaker, Bulkhead & Fault Tolerance
Step 2.3: At 00:00, 10M Jobs Fire at Once
The problem: 8M hourly jobs and 2M daily jobs are all due at 00:00:00 UTC. The fact source for this track plans for "5× the average", about 29K runs a second. What would you do? What is the real peak, and how do we survive it?
The peak, stated once: with 10M due at 00:00, the dispatch rate is burst ÷ window:
| Window | Rate for the 10M | Plus background 3,449/s | What it takes |
|---|---|---|---|
| 0 s (everything "exactly on time") | 10M in 1 s | – | Impossible to size |
| 60 s (the fact source's jitter) | 166,667/s | ≈ 170K/s | Still 29× the average |
| 300 s (our default) | 33,333/s | ≈ 36.8K/s | What we size for (R2.6) |
| 3,600 s | 2,778/s | ≈ 6.2K/s | Flat load; see the cost lever in R2.7 |
A 60-second jitter doesn't lower the peak below the "first minute" rate: it only makes the minute uniform. The window length sets the rate.
Why a queue at all, rather than having dispatchers call targets directly? A direct call ties the dispatcher's pace to the slowest target and has nowhere to put a burst: the target either absorbs 33K calls a second or fails them. With a queue in between, dispatchers stay fast and precise, executors drain at the rate targets can take, and executors can be scaled and fail independently. The price is one more hop and one more component to watch (queue age). What happens when an executor dies halfway through a message is in R2.5.
Drill: Message queue order pipeline
Step 2.4: A Paused Worker Finished a Job Someone Else Already Re-Ran
The problem: a long job's worker hits a 45-second stop-the-world garbage-collection pause. Its lease expires; the reaper gives the run to worker B, which starts over. Worker A wakes up, finishes its old work, and writes "succeeded" with its partial results. Worker B later writes "succeeded" too, and the target got two sets of writes. What would you do? How do we stop a worker that has lost its lease from changing anything?
Synthesizing vector architecture diagram...
Worker A's writes to our table fail on the epoch. Its late call to the target is harmless only because the target checks the idempotency key or the token.
Why an epoch and not a worker ID? Round 1's completion was conditional on lease_owner = me. Two holes: a restarted process can come back with the same worker ID, and a worker can hold two claims on the same run over time (claimed, lost, reclaimed by itself). An epoch is unique per claim. And it's a number, so a downstream system can compare "higher wins"; worker IDs have no order.
Why not just optimistic concurrency at the database, with no leases at all? That's partly what we do: every write is a conditional write on a version (the epoch, or next_run_at for the job item). Optimistic concurrency works well when conflicts are rare, and here they are: a few claimers per run. Under heavy contention it would cause a storm of failed writes and retries. What version checks alone can't do is tell anyone that a worker has died: the lease adds the expiry, the reaper acts on it, and the epoch makes the reaper's takeover safe. An external lock service would add a component and a network hop and still need the same fencing token at the resource.
Primitive: Distributed Locks & Leases · Drill: Distributed lock fencing token
Step 2.5: Long Jobs Look Dead and Get Reclaimed
The problem: a 6-hour ledger rebuild runs on a team's worker. Its lease is 30 seconds. It gets reclaimed every 30 seconds, and 12 copies are soon running. What would you do? How does a long job keep its claim, and how do we still notice quickly when its worker dies?
Why the lease index has 16 shards. Each heartbeat changes the index's sort key, which costs a delete and a put in the GSI: 2,000 × 2 = 4,000 index writes a second. In one partition-key value that's four times what one partition serves; over 16 it's 250 each.
An adaptive reaper. A reaper that reclaims everything it finds is dangerous exactly when many leases expire at once, and the most common cause of that isn't 20K dead workers: it's our heartbeat API or its network having an outage. Reclaiming every lease would start a second copy of every running job. So:
| Situation | Reaper behavior |
|---|---|
| Normal: a few expiries a minute | Reclaim each at once |
| More than 1% of running leases expired in the last minute | Stop reclaiming, page on-call; probably our problem, not theirs |
| Our heartbeat API was unavailable from T1 to T2 | Extend every lease by T2 − T1 before resuming |
| Reclaims resume after a pause | At most 200 a second, oldest first, so the task queues and targets aren't flooded |
Step 2.6: A Target Is Slow and We're Making It Worse
The problem: at midnight, 40,000 report jobs call one analytics webhook. Its latency climbs from 200 ms to 8 s, then it returns 503. Our executors time out, retry, and keep sending at full rate. The service can't recover, and every executor slot is stuck waiting on it.
What would you do? How do we send each target only what it can take?
Why not just very short timeouts? A 50 ms timeout on every call would free executor slots quickly, but a healthy service with normal jitter would fail a share of good calls, and each false timeout becomes a retry, which adds load. A timeout sets how long one call may take; the breaker and AIMD decide whether to call at all. We set timeouts per target from its own latency SLO (here 2 s), and let the breaker cut off a target that's truly down.
Why executors run out of slots without this. Each executor holds at most 500 calls in flight. If a target slows from 200 ms to 8 s, each of its calls holds a slot 40 times longer: at 1,000 calls a second to that target, it would hold 1,000 × 8 s = 8,000 slots, the whole default pool, and every other target's jobs would wait behind it. The per-target limit caps that target at L slots, so the slowdown stays its own problem.
Go deeper: SQS fair queues. In a standard SQS queue, if messages carry a MessageGroupId (here, the target), SQS watches how many messages from each group are in flight; when one group has far more than the others, it prefers returning messages from the other groups, which keeps their waiting time low while the heavy group's backlog drains. It doesn't limit any group's rate and doesn't enforce order. It's a useful second line behind AIMD for the same "one target hogs the pool" problem, and Round 3 uses it for tenants.
Primitive: Distributed Rate Limiting · Circuit Breaker, Bulkhead & Fault Tolerance · Drill: Circuit breaker cascading thread stall
Round 2 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 2.1 | Scanning 100M jobs is too slow | due-index by minute and shard; 2-hour Valkey buckets; markers and straggler queries | A loader; two stores to reconcile |
| 2.2 | Payments wait behind reports | Four queues and executor pools; reserved CRITICAL pool; strict allowance of 5,000/s | Pools to size; idle reserve |
| 2.3 | 10M jobs at 00:00 | Forward-only dispatch windows (default 300 s) with hashed offsets; pre-loading; queues; pre-scaling | Jobs start up to 5 minutes late |
| 2.4 | A stale worker overwrote a newer run | Epoch fencing on every write; token and idempotency key downstream | Targets must cooperate |
| 2.5 | Long jobs reclaimed | 10 s heartbeats, 30 s leases; sharded lease index; adaptive reaper | Heartbeat writes |
| 2.6 | A slow target made worse | AIMD concurrency per target; circuit breaker; deferral by re-send | Backlog for that target |
R2.5 Architecture v2
Synthesizing vector architecture diagram...
Read it as three concerns. When: loaders copy the next 2 hours into Valkey buckets; dispatchers move due members into priority queues. Who: executor pools claim runs in DynamoDB with epochs; team workers claim long jobs through the API and heartbeat. What if it fails: retries through delayed messages, the DLQ, and the reaper. Every finished run flows to S3 through the table's stream.
Tables (one DynamoDB table, jobs, holds both item types)
| Item | Key | Main attributes |
|---|---|---|
| Job | PK = JOB#<job_id>, SK = DEF | schedule, priority, dispatch_window_seconds, target_id, payload (≤ ~700 B; larger payloads go to S3 by reference), job_state (ACTIVE, PAUSED, CANCELLED, DONE), next_run_at, due_key |
| Run | PK = JOB#<job_id>, SK = RUN#<scheduled_for> | run_state, epoch, attempt, lease_expires_at, lease_pk (long jobs only), outcome, expires_at (TTL: 3 days after finishing; none while FAILED) |
A one-off job keeps its single run on the job item itself, so it doesn't pay for a second item.
The task state machine (a run)
Synthesizing vector architecture diagram...
Every arrow is a conditional write, and every arrow into RUNNING raises the epoch, which is what fences the previous holder.
The executor's rules (short webhook jobs, timeout at most 5 minutes)
- Receive up to 10 messages with long polling, and process them concurrently. Messages from one receive share one visibility clock that started at the receive; working through them one after another would let the last ones reappear and run twice.
- Take a target slot (step 2.6); if none, delete and re-send with a delay.
- Claim the run first: create or update the run item
IFit doesn't exist, or is inRETRY_WAIT, or isRUNNINGwith a lease that ended more than 5 s ago. The claim raises the epoch and sets the lease tonow + timeout + 60 s. - Then advance the job:
IF next_run_at = :scheduled AND job_state = ACTIVE, set the next occurrence (orDONEfor one-offs), withReturnValues: ALL_NEWto get the target and payload in the same call. If the next occurrence is within 2 hours, add it to its bucket. If a crash happens between steps 3 and 4, the job is still due; the next claimer finds the run item and finishes the advance. If the condition fails, the executor reads the job: if it has already moved past this occurrence (a retry, or a claimer that crashed after advancing), it carries on; if the job was cancelled, paused or rescheduled away from this time, the run is markedCANCELLEDand nothing is called. - Call the target with
Idempotency-KeyandX-Fencing-Token. - Complete
IF epoch = :e, then delete the message. - If the claim failed because another executor holds a live lease, don't delete the message; hide it until just after that lease ends. If that holder dies, the message comes back and the next executor reclaims the run. Delete a message only when the run is finished or we finished it.
An executor that crashes halfway through a message never deletes it, so after the visibility timeout (timeout + 60 s) the message reappears and another executor reclaims the run with epoch + 1. A message that has crashed executors 10 times moves to a poison-message queue (maxReceiveCount = 10), separate from our DLQ of failed runs.
Retries. On a retryable failure the executor sets RETRY_WAIT with attempt + 1 (IF epoch = :e), sends a new message with DelaySeconds equal to the backoff (at most 900 s), and deletes the old one.
Trace 1: the midnight burst
Synthesizing vector architecture diagram...
Across all shards, about 33.3K windowed jobs a second leave the buckets from 00:00:00 to 00:05:00, plus 3,449 a second of background work.
Trace 2: a zombie rejected by fencing is the sequence in step 2.4: the reaper raises the epoch, and every later write by the paused worker fails its condition.
Trace 3: a slow target throttled
Synthesizing vector architecture diagram...
The target sees its load cut in half within a second of slowing down, and nothing at all while its breaker is open. Its own jobs wait in delayed messages; other targets' jobs never notice.
History. Every finished run is a change on the table, so the table's DynamoDB stream carries it. A Lambda function reads the stream, keeps only the changes into a final state (an event filter drops the rest before the function runs), packs them into newline-delimited JSON records of up to about 1,000 KB, and sends them to Amazon Data Firehose, which writes GZIP files to S3. Packing matters: Firehose bills each record rounded up to the next 5 KB, so 1 KB records sent one by one would be billed as 5 KB each. The lake answers "what happened to my job last month" with Athena, and holds 13 months.
R2.6 Numbers and Cost
Targets
| Quality | Target | Why this number |
|---|---|---|
| Dispatch lag | P99 < 1 s from the effective time (scheduled time + offset) to a claimed run | From the scope raise |
| Strict jobs | ≤ 5,000 due in any one second, platform-wide | The CRITICAL allowance (step 2.2) |
| Window | Non-strict jobs start within their window (default 300 s) | Step 2.3 |
| Availability | 99.99% for the API and dispatch | About 52.6 minutes a year |
| Durability | A created job runs, or ends FAILED with a record and an alarm | As in Round 1 |
Where the 500M runs come from (the interviewer's shape, plus our assumptions for the rest)
| Job class | Active jobs | Runs/day |
|---|---|---|
Hourly at minute 0 (0 * * * *) | 8M | 8M × 24 = 192M |
| Daily at minute 0: 2M at 00:00 UTC, 8M at local midnights (at most 1.5M in any other hour) | 10M | 10M |
| Other recurring jobs (mostly daily and weekly, at scattered minutes) | 62M | 48M |
| One-off and immediate tasks: 250M a day, waiting 1.9 h on average | 250M × 0.08 day = 20M | 250M |
| Total | 100M | 500M |
| Item | Math | Result |
|---|---|---|
| Average rate | 500M ÷ 86,400 s | 5,787/s |
| Background (everything not at minute 0) | (500M − 192M − 10M) ÷ 86,400 | 3,449/s |
| Due at 00:00:00 UTC | 8M hourly + 2M daily | 10M |
| Due at any other top of the hour | 8M + at most 1.5M | ≤ 9.5M |
| The "×5" estimate | 5,787 × 5 | 28.9K/s, which would take 346 s to drain 10M |
| Planned peak | 10M ÷ 300 s + 3,449 | ≈ 36.8K runs/s for 5 minutes |
| Strict burst | CRITICAL allowance | ≤ 5,000 in the first second, in their own pool |
Dispatch latency budget (steps happen one after another, so they add up; each figure is an assumption to confirm in load tests)
| Step | P99 |
|---|---|
| Wait for the next dispatcher pass (every 250 ms) | 250 ms |
ZRANGEBYSCORE on the bucket | 2 ms |
SendMessageBatch to SQS | 50 ms |
| Executor's long-poll receive returns | 20 ms |
| Claim the run, then advance the job (two conditional writes) | 2 × 15 ms |
| Total | ≈ 352 ms, leaving ~650 ms for pauses, retries and queueing |
This holds only if a free executor is waiting when the message arrives, which is why pools are sized for the peak, not the average.
Valkey buckets
| Item | Math | Result |
|---|---|---|
| Members in 2 hours, average | 500M × 2 ÷ 24 = 5,787/s × 7,200 s | 41.67M |
| Memory, average, at ~64 B a member (the fact source's estimate) | 41.67M × 64 B | ≈ 2.7 GB (the fact source's 2.66 GB rounds to 41.6M) |
| Worst 2-hour window | background 3,449 × 7,200 = 24.8M, plus 00:00 (10M) and 01:00 (≤ 9.5M) | ≈ 44.3M |
| Memory, worst, at a more conservative ~100 B a member (skip list node, hash table entry, a ~30-byte member string) | 44.3M × 100 B | ≈ 4.4 GB |
| Cluster | 3 shards, each a primary and a replica in another AZ, on cache.r7g.large (13.07 GiB) | ~1.5 GB per shard at worst |
| Commands at midnight | 256 buckets × 4 passes/s ≈ 1,000 range reads + 1,000 batched removes a second; ~26.7K re-arm adds a second | ~30K commands/s over 3 shards |
Executors (assumption to load-test: one 1-vCPU, 2 GB Graviton task sustains 500 runs a second of I/O-bound HTTP work)
| Pool | Peak rate | Tasks needed | With AZ headroom (× 1.5) | Off-peak |
|---|---|---|---|---|
| CRITICAL | 5,000/s (allowance) | 10 | 15, always on | 15 |
| HIGH (20% of windowed) | 7.4K/s | 15 | 23 | 3 |
| DEFAULT (60%) | 22.1K/s | 45 | 68 | 6 |
| LOW (20%) | 7.4K/s | 15 | 23 | 3 |
| Windowed total | 36.8K/s | 75 | 114 | 12 |
Why × 1.5: losing one AZ of three leaves two-thirds of the tasks. The pools are separate bulkheads, so the math must hold for each pool on its own, rounding up: HIGH and LOW 15 × 1.5 = 22.5 → 23 (23 × 2/3 ≈ 15.3, still ≥ 15; 22 would leave 14.7), DEFAULT 45 × 1.5 = 67.5 → 68 (68 × 2/3 ≈ 45.3 ≥ 45). Off-peak, the background 3,449/s needs 7 tasks; 12 keeps one per AZ per pool with room. The burst comes every hour (≤ 9.5M at other hours: 9.5M ÷ 300 + 3,449 ≈ 35.1K/s), so scheduled scaling runs the 114 from :57 to :08, 11 minutes of every hour.
Long jobs: 20K running at once (assumption), heartbeats every 10 s = 2,000/s.
DynamoDB writes (one write unit per 1 KB; moving an item in a GSI costs 2)
| Kind | Writes each | Per day |
|---|---|---|
Recurring run: run item 1, advance job 1 + due-index move 2, complete 1 | 5 | 250M × 5 = 1.25B |
| One-off: create 1 + index 1, claim 1 + index removal 1, complete 1 | 5 | 250M × 5 = 1.25B |
| Retries: 2% of runs retry once (assumption) | 2 | 10M × 2 = 0.02B |
Heartbeats: table 1 + lease-index move 2 | 3 | 2,000/s × 3 × 86,400 = 0.52B |
| Total | ≈ 3.04B a day, 91.2B a month |
Average: 3.04B ÷ 86,400 ≈ 35.2K writes a second. At the midnight peak: 33.3K × 5 (the burst's runs) + 3,449 × 5 (background) + 6,000 (heartbeats) ≈ 190K writes a second. We set the table's and indexes' warm throughput above that, so DynamoDB has the partitions ready before the first midnight. On-demand tables also have a default maximum of 40,000 write request units a second per table and per GSI (adjustable), so we request a quota increase to cover ~190K before launch.
On-demand or provisioned? Provisioned capacity costs $0.00065 per write unit per hour, which is $0.00065 ÷ 3,600 = $0.18 per million writes if every unit is used every second; on-demand costs $0.625 per million. Provisioned wins only above 0.18 ÷ 0.625 ≈ 29% average utilization. Provisioned flat at ~190K, we'd use 35.2K on average, 18.5%: far worse. Provisioned that follows the spike is cheaper on paper: ~190K for the 11 minutes around each top of the hour and ~40K otherwise averages (190K × 11 + 40K × 49) ÷ 60 ≈ 67K write units, 67,000 × $0.00065 × 730 ≈ $32,000 a month, well below on-demand's $57,000. DynamoDB allows enough decreases for that (four at the start of a day, then one per hour, up to 27 a day). The catch is the increase: before :57 every hour, the table and every GSI must go from ~40K to ~190K, a large increase isn't instant and isn't guaranteed to finish on time, and if it's late the burst is throttled and the 1-second SLO is broken for everyone at the worst moment. Provisioned wins only with pre-scaling that reliably completes every hour. We take on-demand, with warm throughput for the spike, pay about $25K a month for that certainty, and revisit when longer windows flatten the load (R2.7).
Storage
| Item | Math | Result |
|---|---|---|
| Active jobs | 100M × 1 KB | 100 GB |
| Finished one-offs, kept 3 days | 250M × 3 × 1 KB | 750 GB |
| Recurring run items, kept 3 days | 250M × 3 × 0.5 KB | 375 GB |
| Indexes | ~100M × ~150 B projected | ~15 GB |
| Total | TTL deletes typically within a few days, so plan up to ~2 TB | ≈ 1.25–2 TB |
| History lake, raw | 500M × 1 KB × 30 days | 15 TB/month |
| History lake, stored | GZIP at an assumed 5:1 | ~3 TB/month; ~39 TB once 13 months are held |
S3 is designed for 99.999999999% (11 nines) durability of objects. That figure is S3's, for the history files; it says nothing about whether a job runs, which is what the durability target above is about.
Rough monthly cost (us-east-1 list prices; check the AWS Pricing Calculator before quoting)
| Line | Math | ≈ Monthly |
|---|---|---|
| DynamoDB writes | 91.2B × $0.625 per million | $57,000 |
| DynamoDB storage | ~1.6 TB × $0.25/GB | $400 |
| DynamoDB reads | status reads (50M/day, assumed) + loader, straggler and reaper queries ≈ 35M read units/day × 30 × $0.125 per million | $130 |
| SQS | ~0.5 requests per run with batching (send, receive, delete in 10s, plus delays and polls): 7.5B × $0.40 per million | $3,000 |
| ElastiCache for Valkey | 6 × cache.r7g.large at ~$0.175/h (Redis OSS lists at $0.219/h; Valkey is about 20% less) × 730 h | $770 |
| Executors | Fargate Graviton, 1 vCPU and 2 GB at $0.0395/h: CRITICAL 15 × 730 h + off-peak 12 × 730 h + burst 102 × 11/60 × 730 h ≈ 33,360 task-hours | $1,320 |
| Dispatchers and loaders | 16 tasks × 730 h × $0.0395 | $460 |
| API | 12 tasks of 2 vCPU and 4 GB at $0.079/h × 730 h | $690 |
| Application Load Balancer | $16 base + ~25 capacity units (about 25 GB processed an hour) × $0.008 × 730 h | $165 |
| History: Lambda + Firehose | Lambda ~$100; Firehose 15 TB × $0.029/GB = $435 (unpacked 1 KB records would bill 75 TB: $2,175) | $535 |
| S3 history | ~39 TB at 13 months' retention × $0.023/GB | $900 |
| Network to targets | ~2 KB per call, 30 TB/month; if targets sit in other VPCs behind Transit Gateway, $0.02/GB processed (an assumption about the company's network) | $600 |
| CloudWatch | metrics, logs, alarms (estimate) | $1,000 |
| Total | ≈ $67,000 = about $4.50 per million runs |
Three lessons:
- DynamoDB writes are 85% of the bill. Every run costs about five writes, and the index moves are two of them. The biggest levers are fewer writes per run and a flatter load (R2.7).
- The fleet is cheap because it's I/O-bound and scaled to the hour. Running the 114 windowed executors all day would cost
114 × 730 × $0.0395 ≈ $3,290instead of about $890. - Firehose's 5 KB rounding would have made its line five times bigger. Packing records is a one-line design decision worth about $1,700 a month.
R2.7 Trade-Offs
Where the timers live
| Valkey time buckets (chosen) | DynamoDB due-index only | SQS delay messages | |
|---|---|---|---|
| Precision | 250 ms passes over exact scores | A query per shard every 250 ms: 256 × 4 = 1,024 queries a second, and the index trails the table by a moment | Seconds, but delays max out at 15 minutes |
| How far ahead | 2 hours loaded; the table holds the rest | Unlimited | 15 minutes |
| Cancel and reschedule | Stale members are harmless (conditional claim) | Native | A delayed message can't be found or removed |
| Cost and parts | A cache cluster and a loader, ~$770 + tasks | No extra parts; ~$170 more in reads | None extra |
| When the cache is lost | Fall back to reading the index (the markers tell us when) | – | – |
Honestly, the index-only design would work at this scale, and it's simpler. We keep the buckets because they turn every dispatcher pass into one in-memory range read with exact per-job scores (offsets included), and they keep the hot midnight read path off DynamoDB, which is already absorbing ~190K writes a second. The index remains the fallback for any bucket we lose. SQS delays are right for what they're good at: retry backoff up to 15 minutes.
Push to webhooks vs workers pulling
| Push: our executors call HTTP targets | Pull: team workers poll our API | |
|---|---|---|
| Suits | Short calls (≤ 5 minutes) | Long jobs, and work that needs the team's code and data |
| Liveness | SQS visibility timeout + lease | Heartbeats + lease + reaper |
| Load control | We must protect the target (AIMD) | The worker takes work when it has room |
| Team effort | Expose an endpoint | Run a worker fleet |
We support both, because the scope asked for both.
The size of the window
| Window | Peak | Cost effect |
|---|---|---|
| 60 s | ≈ 170K runs/s | Executors, DynamoDB and warm throughput sized ~4.6× larger than at 300 s |
| 300 s | ≈ 36.8K runs/s | On-demand DynamoDB at ~$57K |
| 3,600 s | ≈ 6.2K runs/s | Load is almost flat (~40K writes/s): provisioning ~45K write units costs 45,000 × $0.00065 × 730 ≈ $21,400 a month, less than half the on-demand bill |
The window is a product decision and the biggest cost lever we have. A good next step: default LOW-priority jobs to a 30-minute window, and show teams what their window costs.
Buy: EventBridge Scheduler at this scale? List price: 15B invocations a month × $1.00 per million ≈ $15,000, far less than our $67,000. What stops us: 60-second precision (we need P99 < 1 s), no priority classes or reserved capacity, no heartbeat-and-lease model for long jobs, and no per-target AIMD. Its default quotas (10M schedules, 1,000 invocations a second) are adjustable, AWS says to billions of schedules and tens of thousands of invocations a second, so quotas alone wouldn't decide it. If the business dropped the 1-second requirement, buying would be the right call for most of the jobs. COST 11
R2.8 Failure Modes
| Trigger | What you'd see | How the design responds |
|---|---|---|
| A Valkey shard fails over | The replica is promoted; writes from the last moments before the failure may be missing, because replication is asynchronous. | Missing markers send dispatchers to due-index for those shard-minutes; straggler queries catch any member lost from a loaded bucket. Jobs run late by up to about a minute; none are lost. |
| A Valkey shard and its replica are both lost | Every bucket on that shard is gone. | Dispatchers for its shards run from due-index directly (1,024 small queries a second at most) while the loader rebuilds the next 2 hours at 20K members a second. |
| An AZ goes down | A third of executors, dispatchers and API tasks vanish. | Each pool was sized ×1.5 on its own, so the surviving 15 HIGH, 45 DEFAULT and 15 LOW executors cover the 36.8K/s peak and the 10 surviving CRITICAL executors cover the 5,000/s allowance. Dead dispatchers' shard leases expire in ~10 s and are taken by survivors: jobs in those shards are late by up to ~10 s once, so the P99 target can be missed for that minute. Valkey promotes replicas in other AZs. DynamoDB and SQS are Regional services and keep working. |
| A stampede bigger than planned | A team adds 5M hourly jobs with a 60 s window. | Creation-time checks: the API tracks due-minute load and refuses windows that would push a minute past the planned rate, suggesting a longer window. Without that, queues absorb the excess and it drains late; the "queue age" alarm fires. |
| A reaper storm | Our heartbeat API had a 3-minute outage; 20K leases expired at once. | The adaptive reaper pauses above 1% expiries a minute, extends leases by the outage length and pages us (step 2.5). |
| The DLQ grows silently | Hundreds of failed runs a day nobody looks at, until someone asks why a report never ran. | Alarms on DLQ depth per priority (any CRITICAL message pages) and on the age of its oldest message. The DLQ holds messages for at most 14 days, so the run item is the real record: a FAILED run has no TTL until it's redriven or acknowledged. A daily report per team lists its failed runs. |
Primitive: Message Queues vs Event Streams
R2.9 Production Gotchas
1. All timers in one sorted set
- Symptom: one Valkey node at 100% CPU at midnight while the others idle; dispatch lag in minutes.
- Cause: a single key holding every timer; one key lives on one node and one thread.
- Fix: per-minute, per-shard buckets with hash tags (step 2.1).
2. No fencing on heartbeats
- Symptom: a job reclaimed after a pause keeps both workers heartbeating "successfully", and two results are recorded.
- Cause: heartbeats conditional only on the run ID.
- Fix: every heartbeat and completion is
IF epoch = :e AND run_state = RUNNING(step 2.4).
3. Synchronous webhook calls inside the dispatcher
- Symptom: one slow target delays thousands of unrelated jobs by seconds.
- Cause: the dispatcher waits on HTTP calls between bucket passes.
- Fix: dispatchers only enqueue; executors call targets under AIMD limits (steps 2.3, 2.6).
4. Non-idempotent jobs
- Symptom: after an AZ event, some customers got two invoices.
- Cause: the target ignored the
Idempotency-Key. - Fix: make the key part of the target's contract, and test it: our staging environment deliberately delivers 1% of calls twice.
5. Deferring with ChangeMessageVisibility
- Symptom: healthy jobs for a throttled target end up in the poison-message queue.
- Cause: each deferral is another receive, counted toward
maxReceiveCount. - Fix: defer by deleting and re-sending with a delay (step 2.6).
6. Random jitter instead of hashed offsets
- Symptom: after a bucket reload, some jobs run twice in one hour: once at their old random time, once at the new one.
- Cause: a new random offset on every load gives the same run two dispatch times.
- Fix:
offset = hash(job_id) mod window, stable across reloads and the same every hour (step 2.3); and the conditional claim rejects the second copy anyway.
R2.10 Pillar Check
| Pillar | What Round 2 adds |
|---|---|
| Reliability | Bulkheads per priority; pools sized for AZ loss; epoch fencing; heartbeats and an adaptive reaper; buckets rebuildable from the table; markers and straggler queries. REL 5 · REL 10 · REL 11 |
| Performance Efficiency | Time buckets with 250 ms passes; hashed dispatch windows; the peak derived from the cron distribution; a latency budget that sums to ~352 ms. PERF 3 |
| Security | Registered targets only; HMAC-signed calls; least-privilege roles per component (below). SEC 3 · SEC 5 · SEC 9 |
| Cost Optimization | ≈ $67K a month, 85% DynamoDB writes; on-demand justified by the 18.5% utilization; windows as a cost lever; packed Firehose records. COST 5 · COST 7 |
| Operational Excellence | Alarms with first actions (below); scheduled pre-scaling; a redrive API with a rate cap. OPS 8 · OPS 10 |
| Sustainability | Light this round: executors run at burst size 11 minutes an hour instead of all day; longer windows flatten load and need less standing capacity. SUS 2 |
Security in detail. A job can only name a registered target, and registering one requires the owning team's approval, so a team can't point jobs at another team's admin endpoint or at the instance metadata address. Every call is signed with an HMAC of the body, the timestamp and the idempotency key using the target's secret from Secrets Manager, so the target can reject calls that didn't come from us or that were replayed later. Executors have one IAM role: receive and delete on their own pool's queue, send to the delay path, conditional updates on the jobs table, read the target secrets. Dispatchers can send to queues but can't touch the table except their shard leases. The API authenticates callers by IAM identity and scopes each team to its own jobs.
Operations in detail.
| Alarm | Threshold | Severity | First action |
|---|---|---|---|
| Dispatch lag P99, per priority | > 1 s for 3 min (CRITICAL: any minute) | P1 | Which shards? Check dispatcher leases and Valkey health |
| Oldest message age, per queue | > window + 60 s | P2 | Is the pool scaled? Is one target holding slots? |
| Straggler jobs found per hour | > 0.01% of runs | P3 | Buckets are losing members: check failovers and loader errors |
| Lease expiries per minute | > 1% of running long jobs | P1 | Reaper has paused: is our heartbeat API healthy? |
| DLQ depth, CRITICAL | > 0 | P1 | Look at the failed run and its target |
| DynamoDB throttled requests | > 0 for 1 min | P2 | Check warm throughput and hot partitions |
R2.11 Round 2 Rubric and Follow-Ups
What a strong senior (L6) answer adds over L5
- Derives the peak from how cron jobs cluster, rejects the uniform multiplier, and shows the window sets the rate.
- Keeps the durable table as the truth and treats the in-memory buckets as a rebuildable cache, with a way to detect what the cache lost.
- Uses bulkheads, not a priority field, to protect critical work, and caps strict jobs so the SLO is provable.
- Explains fencing honestly: it protects our records; external effects need the target to check the key or the token.
- Designs heartbeats and a reaper that won't turn our own outage into 20K duplicate jobs.
- Protects targets with AIMD and a breaker, and knows the SQS receive-count trap.
- Finds the cost driver (writes per run) and the lever (the window).
Follow-up questions
-
"Can a job run early?" Answer: not by design. Offsets only move jobs later, and an executor won't claim a run whose scheduled time is still ahead by its own clock. The remaining risk is a worker whose clock runs fast; we alarm on clock offset above 1 s, and a job that must never be early by even a second should be strict, where the only delay is dispatch.
-
"How would you add a 'run at most once' option for a team that would rather miss a run than repeat it?" Answer: mark the run
STARTEDwith a conditional write before calling the target, and never reclaim aSTARTEDrun: if its executor dies, the run becomesUNKNOWNand goes to the team's DLQ instead of being retried. That's at-most-once: some runs will be lost, and the team must accept that in writing. -
"Your executors call a target that takes 4 minutes. What visibility timeout do you set, and why does it matter?" Answer: the job's timeout plus a margin, here 5 minutes. Shorter, and the message reappears while the first executor is still working; the second executor's claim fails on the live lease, and it hides the message until after that lease ends (rule 7), so nothing breaks, but we waste receives. Longer, and a crashed executor's run waits longer to be retried.
Interview gotchas from this round's wrong answers
| Gotcha | Why it's wrong |
|---|---|
| "Peak is 5× the average" | 10M jobs share one instant; the window, not a multiplier, sets the rate. |
| "60 s of jitter solves the stampede" | It still means ~170K runs a second for that minute. |
| "Fencing tokens prevent duplicate execution" | Only for systems that check them. Webhooks need idempotency keys. |
| "A 24-hour lease for long jobs" | A dead worker goes unnoticed for 24 hours. |
| "Reclaim every expired lease immediately" | Our own heartbeat outage becomes a duplicate of every running job. |
| "Just make the message visible again to defer it" | Each receive counts toward maxReceiveCount. |
| "EventBridge Scheduler meets P99 < 1 s" | It invokes with 60-second precision. |
Round 3 · Architect · "Workflows, Tenants, and Regions"
~45 min · Principal (L7) · 3 regions · 5,000 tenants · 1B job runs + 500M workflow steps a day · 1.7B pending timers · survives a region with duplicates bounded by replication lag
R3.0 Where We Left Off
This is what the candidate says aloud in the first 60 seconds of Round 3. If you're starting here, it's everything you need from Round 2.
Round 2 in 60 seconds. "We run the company's scheduler: 100M active jobs, 500M runs a day, 5,787 a second on average, in one region across three AZs. The worst moment is midnight UTC, when 10M jobs are due at once; we derived the peak from that, not from a multiplier. Non-strict jobs get a forward-only dispatch window, 5 minutes by default, with an offset hashed from the job ID, which turns the burst into about 36.8K runs a second. Strict jobs need a CRITICAL allowance capped at 5,000 a second. DynamoDB is the source of truth; a
due-indexkeyed by minute and shard feeds 2-hour Valkey buckets, with markers and straggler queries to catch anything the cache loses. Dispatchers move due jobs into four priority queues with separate executor pools. Every run has an epoch: every write is conditional on it, and targets get it plus an idempotency key. Long jobs run on team workers with 10-second heartbeats and an adaptive reaper. AIMD limits protect each target. About $67K a month, 85% of it DynamoDB writes. Open costs: jobs have no memory of each other, every team shares one pool, and it all lives in one region."
Architecture v2, compact
Synthesizing vector architecture diagram...
Round 2 in one picture: the table is the truth, the buckets make time precise, the queues and pools isolate priorities, and epochs fence stale workers.
Round 2 step summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 2.1 | Scanning 100M jobs | due-index by minute and shard; Valkey buckets; markers, straggler queries | A loader; two stores |
| 2.2 | Payments behind reports | Queues and pools per priority; reserved CRITICAL; strict allowance | Idle reserve |
| 2.3 | 10M at 00:00 | Forward-only windows with hashed offsets; pre-load; pre-scale | Jobs up to 5 minutes late |
| 2.4 | Stale workers | Epoch fencing; key and token downstream | Targets must cooperate |
| 2.5 | Long jobs reclaimed | Heartbeats, sharded lease index, adaptive reaper | Heartbeat writes |
| 2.6 | Slow targets | AIMD per target; circuit breaker | Backlog per target |
Open costs: a multi-step process is a chain of unrelated jobs; one team's flood competes with everyone's; a region outage stops every timer.
R3.1 The Scope Raise
Interviewer: "We're turning the scheduler into a product that outside customers pay for. They don't just want jobs: they want workflows. Charge a card, reserve stock, ship; if shipping fails, release the stock and refund the card. Some workflows wait days for a signal, like 'the courier picked it up'."
Interviewer: "We'll have thousands of tenants. One of them already has 50 million jobs, almost all at midnight, and it must not delay anybody else. We're going into three regions. If a region fails, timers must still fire, and when the region comes back we can't have a flood of duplicate charges. And finance wants to know what each tenant costs us."
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| How many workflows, and how long do they live? | About 50M started a day, 10 steps each on average. Many sleep for weeks: renewals, reminders, return windows. Average life is about 30 days. | 500M steps a day, and 50M × 30 = 1.5B open workflows, each waiting on a timer or a signal. A durable engine with event histories (step 3.1). |
| What happens when a step fails halfway through a workflow? | Undo the earlier steps: refund, release. Some undo steps can fail too. | Sagas with compensations that retry durably, and a human path when they can't finish (step 3.1). |
| How uneven are tenants? | 5,000 tenants. The biggest has 50M jobs, 80% of them at 00:00 UTC. | Per-tenant quotas, fair dispatch, and cells for the largest tenants (step 3.2). |
| Which regions, and what must survive? | us-east-1, eu-west-1 and ap-northeast-1; about 60%, 25% and 15% of the load. Losing one must not stop timers. | Each workflow has a home region; ownership moves on failure under a new epoch (step 3.3). |
| Sales has been promising "exactly once". Can we? | That's what we're asking you. | A precise contract: which transitions are exactly-once, which calls are at-least-once (step 3.4). |
| Do some tenants need their data kept in their home region? | Some do. Legal will tell us who. | A per-workflow region policy, including one that never replicates (R3.3). We keep legal claims general. |
| How do we charge? | Finance wants cost per tenant, not an even split. | Metering of timers, executions, history bytes and compute (step 3.5). |
Scope change
| Round 2 | Round 3 | |
|---|---|---|
| Customers | internal teams | 5,000 external tenants, one with 50M jobs |
| Work | 500M job runs/day | 1B job runs + 500M workflow steps a day |
| State per unit | one run | event history per workflow, open for ~30 days |
| Pending timers | ~100M | ~1.7B (1.5B workflows + 200M jobs) |
| Footprint | 1 region, 3 AZs | 3 regions: us-east-1 (60%), eu-west-1 (25%), ap-northeast-1 (15%) |
| Survives | an AZ | a region, without a flood of duplicates when it returns |
| SLOs | one for everyone | per tenant: dispatch lag within quota; workflow step latency |
| Cost | one bill | per tenant |
R3.2 What Breaks in the Round 2 Design
| Round 2 choice | What breaks at the new scope |
|---|---|
| Independent jobs | A 3-step order becomes three jobs whose callbacks chain them. When step 2 fails, nothing knows step 1 happened, and nothing can undo it. Wait-for-a-signal-for-3-days has nowhere to live. |
| One shared set of pools and queues | The 50M-job tenant fills every bucket and queue at midnight; a small tenant's job waits behind it. |
| Priorities by class, not by tenant | Two CRITICAL tenants still compete, and a tenant can mark everything CRITICAL. |
| One region | A region outage stops 1B timers. |
| State in several places (jobs, runs, queues, target databases) | No single record says where a multi-step process is, so recovery after failover means guessing. |
| One bill | No way to tell which tenant drives the DynamoDB writes. |
The order we fix it in: durable workflows (3.1), tenant isolation (3.2), regions (3.3), the contract (3.4), metering (3.5), and build or buy (3.6).
R3.3 New Requirements and API Additions
A workflow definition (declarative; the engine interprets it)
yamlworkflow: order_fulfillment version: 7 region_policy: HOME_ASYNC steps: - id: charge_card activity: { task_queue: payments, type: CHARGE } timeout: 30s retry: { max_attempts: 5, initial_backoff: 2s, max_backoff: 5m } compensate_with: refund_card - id: reserve_stock activity: { task_queue: inventory, type: RESERVE } timeout: 30s retry: { max_attempts: 5, initial_backoff: 2s, max_backoff: 5m } compensate_with: release_stock - id: wait_for_pickup wait_for_signal: { name: picked_up, timeout: 3d } on_timeout: compensate - id: ship activity: { task_queue: shipping, type: CREATE_SHIPMENT } timeout: 60s retry: { max_attempts: 8, initial_backoff: 5s, max_backoff: 15m } compensations: - id: release_stock activity: { task_queue: inventory, type: RELEASE } retry: { max_attempts: unlimited, initial_backoff: 10s, max_backoff: 1h } - id: refund_card activity: { task_queue: payments, type: REFUND } retry: { max_attempts: unlimited, initial_backoff: 10s, max_backoff: 1h }
Start a workflow (the caller picks the ID, so a retried start is recognized)
httpPOST /v1/workflows HTTP/1.1 Content-Type: application/json X-Tenant-Id: tnt_4411 { "workflow_id": "order-88213", "definition": "order_fulfillment", "version": 7, "input": { "order_id": 88213, "amount_cents": 4599 } }
json{ "workflow_id": "order-88213", "run_id": "wr_01JA0M3", "status": "RUNNING", "home_region": "us-east-1" }
Send a signal (deduplicated by signal_id)
httpPOST /v1/workflows/order-88213/signals/picked_up HTTP/1.1 Content-Type: application/json X-Tenant-Id: tnt_4411 { "signal_id": "courier-evt-5521", "payload": { "courier": "ACME", "at": "2026-09-30T14:02:00Z" } }
Per-tenant quotas
json{ "tenant_id": "tnt_4411", "cell": "use1-shared-3", "quotas": { "job_runs_per_second": 2000, "workflow_steps_per_second": 1000, "strict_jobs_per_second": 50, "pending_timers": 20000000, "api_requests_per_second": 500 }, "default_region_policy": "HOME_ASYNC", "home_region": "us-east-1" }
Region policies (chosen per workflow or job; the price differs)
| Policy | Where state lives | If the home region fails | Write latency | Storage and writes |
|---|---|---|---|---|
SINGLE_REGION | Home only | Stops until the region returns; restored from backups if it doesn't | Local | × 1 |
HOME_ASYNC (default) | Home, replicated asynchronously to the other two (DynamoDB global tables, eventual mode) | Another region takes over in about a minute; work from the last seconds may repeat | Local | × 3 |
STRICT | Three regions, strongly consistent (DynamoDB multi-Region strong consistency) | Another region takes over; nothing acknowledged is lost | Every write waits for another region | × 3 |
Status codes gain 409 WORKFLOW_ID_IN_USE, 422 NONDETERMINISTIC_DEFINITION_CHANGE (a new version must not change steps that running workflows already recorded), and 429 with "scope": "tenant" when a tenant is over its quota.
R3.4 Design Evolution: Durable Workflows Across Tenants and Regions
Step 3.1: A 3-Step Order Workflow Failed at Step 2
The problem: the card was charged (step 1). Reserving stock (step 2) failed after its retries. The customer paid for something we won't ship. Separately, workflows that wait 3 days for a "picked up" signal keep their state in a tenant's own database and are lost when that database has a bad day. What would you do? How does a multi-step process remember what it did, resume after crashes, and undo what it must?
Writing an event safely. The history uses items keyed PK = WF#<workflow_id>, SK = EV#<epoch>#<sequence>, both zero-padded (EV#0007#00000042) so string order equals number order. A state item (SK = STATE) holds next_seq, the status, pending timers, the owning region's epoch, and the sequence at which each epoch took over. The epoch is in the key so that a region returning after a failover can't occupy sequence numbers the new owner needs: its late appends land under the old epoch's prefix, and replay reads each epoch only up to the point where the next one took over. An engine host appends event 42 with a PutItem conditional on attribute_not_exists(SK), then updates the state item IF next_seq = 42. If it crashes between the two, event 42 exists and next_seq is still 42; the next host sees that, finishes the state update, and moves on. It's the same "write the record, then advance the pointer, both idempotent" pattern as Round 2's run-then-job claim. Two hosts racing to append event 42 can't both succeed: the key is unique.
Determinism. Replay only works if the same history always produces the same decisions. So a definition's choices may depend only on recorded values: step results, signal payloads, the recorded start time. Never on the current wall clock, a random number, or a fresh read from another service; those must be recorded as events first. A code-based engine (Temporal, for example) enforces the same rule on workflow code. And a new version of a definition must not change steps that running workflows have already recorded; running workflows keep the version they started with.
When a compensation fails too. A refund can fail: the payment provider is down. Compensations retry with backoff for as long as it takes (max_attempts: unlimited, backoff capped at 1 hour), and each carries an idempotency key (order-88213:refund_card) so a repeat can't refund twice. If a compensation still hasn't succeeded after a tenant-set deadline (say 24 hours), the workflow moves to COMPENSATION_STUCK: it goes to the tenant's DLQ, pages their on-call, and appears on a reconciliation report that compares our history with the provider's records. Money isn't lost silently: the history says exactly which charge has no refund yet.
Why a saga and not two-phase commit? Two-phase commit (2PC) needs every participant to lock its data, vote, and wait for the coordinator's decision. If the coordinator fails after the vote, participants stay locked until it recovers. Across services owned by different companies, over days, with a payment provider that doesn't support 2PC at all, that blocking is unacceptable. A saga never holds locks across steps; the price is that other readers can see the in-between state (card charged, stock not reserved), so each step's effect must be acceptable to see, and each must be undoable.
Synthesizing vector architecture diagram...
Every decision comes from the history. A crash anywhere just means another engine host replays the same events and continues from the first step without a result.
Primitives: Two-Phase Commit & Saga Orchestration · Event Sourcing & CQRS · Drill: Saga orchestration compensating failure
Step 3.2: One Tenant's Flood Delays Everyone
The problem: tenant T schedules 50M jobs, 40M of them at 00:00 UTC. At midnight, T's jobs fill the buckets, the queues and the executors in the cell it shares with 800 other tenants. A small tenant's 00:00 job starts at 00:14. What would you do? How does each tenant get what it pays for, whatever the others do?
Moving a tenant to another cell works like a region takeover (step 3.3) on a small scale: the tenant's timers are copied into the new cell's tables; the old cell stops dispatching for that tenant at a set instant under the tenant's current epoch; the new cell starts under epoch + 1; conditional claims reject anything the old cell still sends.
What about a query across all tenants, like a platform-wide "failed runs by region" report? We never scatter a query across cells' live tables. Every cell streams its finished runs and workflow events to the S3 history lake (R2.5), which we query with Athena. The report is minutes behind and costs nothing on the serving path.
Primitives: Circuit Breaker, Bulkhead & Fault Tolerance · Database Sharding & Partition Keys · Drill: Sharding tenant hotspot
Step 3.3: A Region With 1B Pending Timers Failed
The problem: us-east-1 holds about 1B of our 1.7B pending timers. The region has a serious outage. Timers there must still fire, and when the region comes back, it must not fire them again. What would you do? Where do timers live, and who fires them when their region is gone?
Why the leases use our clocks, not DynamoDB's. DynamoDB conditions can't read the server's time, and conditions on MREC tables are checked against the local replica only. So the safety rule is built from our own clocks with margins on both sides: the old owner stops 5 s before the lease ends by its clock, and the new owner starts 5 s after by its clock. With clock skew under 5 s (we alarm at 1 s), there is a gap of at least a few seconds when nobody owns the group, never an overlap. For STRICT workflows we get more: their state tables are MRSC too, so every conditional write, including the epoch check, is evaluated against the latest version in all regions, and a stale owner's write simply fails.
Why not make every workflow STRICT? Because it's paid for on every write, all day: each write waits for a round trip to another region (tens to a couple of hundred milliseconds, depending on the regions), MRSC tables don't support TTL, so cleanup becomes our job, and there's still no escaping duplicates of external calls that were in flight during the failure. We offer it for the workflows that need zero lost state (payments), and keep HOME_ASYNC as the default.
How many duplicates? (numbers from R3.6)
| Case | Window | us-east-1 actions (jobs + steps, ~10.4K/s) that may repeat |
|---|---|---|
| Region crashes; replication lag ~1 s (AWS says MREC replicates typically within a second or less) | ~1 s | ~10K |
| Lag had grown to our 5 s alarm threshold | 5 s | ~52K |
| Region is isolated but alive: it works until two renewals fail | ≤ 10 s | ≤ ~104K |
Each repeat reaches a tenant with the same idempotency key as the first, so a tenant that honors the key sees no double effect. When us-east-1 returns, its unreplicated writes replicate out. Its events are keyed under the old epoch (EV#0007#…), so they never collide with the new owner's appends (EV#0008#…), and replay ignores anything under epoch 7 past the takeover point; a reconciliation job lists, per tenant, every step that ran under the old epoch after the takeover point, so the tenant can check it.
Synthesizing vector architecture diagram...
The old owner stops before the lease ends; the new one starts after it. The 15-second gap costs us late timers, not duplicate ones. What us-east-1 did between 12:00:02 and 12:00:10 may not have replicated, and that is the duplicate window.
Primitive: Cloud Disaster Recovery & Multi-Region Active-Active · Drill: Replication multi-region consistency
The drill's two questions, answered here. A booking visible in one region and not another is replication lag plus a claim checked against a local copy: exactly why we never let two regions own the same group, and why the takeover condition runs in the strongly consistent ownership table. Why not replicate every write synchronously? Because every write would pay a cross-region round trip and stop when another region is unreachable; we scope that cost to the STRICT workflows that need it and to the tiny ownership table.
Step 3.4: "Exactly Once?"
The problem: a tenant's architect asks: "Your sales team says every step runs exactly once. Is that true?" What would you do? What can we honestly promise?
AWS draws the same line. Step Functions Standard workflows follow an exactly-once model where tasks and states are never run more than once unless you configure Retry; asynchronous Express workflows are at-least-once, and synchronous Express workflows at-most-once. Even "exactly-once" there means the workflow won't start a task twice; a task that calls an outside API and times out still leaves the same question for that API.
Primitive: Database Isolation Levels, ACID & Concurrency Anomalies
Step 3.5: What Does Each Tenant Cost?
The problem: finance splits the $465K monthly bill evenly across 5,000 tenants: $93 each. The 50M-job tenant pays $93; a small tenant with 200 jobs pays $93 too. What would you do? How do we know what each tenant really costs?
The meter records are counted once per engine host per minute and keyed (tenant, region, host, minute), so a record re-sent after a crash overwrites itself in the daily rollup instead of adding to it.
Step 3.6: Build or Buy?
The problem: leadership asks: "Why are we building a workflow engine? AWS has two services for this, and there are hosted engines." What would you do? How do you decide, without dogma either way?
The recommendation, with its conditions. We are selling a scheduling product, the volume is large, and we need three-region failover with per-tenant fairness that none of the managed options gives us directly. Our list cost of about $9 per million steps is well below the managed alternatives at this volume, even after a team costing a few million dollars a year. So: build, and consider building on the open-source Temporal server rather than from scratch. For an internal platform at a twentieth of this volume, the same table says buy Step Functions Standard, whose exactly-once semantics and one-year limit fit most business workflows, and spend the team elsewhere. COST 11
Round 3 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 3.1 | A failed step 2 left a charge | Durable engine: event history, replay, timers and signals as events, sagas | An engine; deterministic definitions |
| 3.2 | One tenant floods everyone | Quotas, fair dispatch, SQS fair queues, cells, dedicated cells for large tenants | Placement and migration |
| 3.3 | A region with 1B timers failed | Home region per placement group; MRSC ownership leases; self-fencing by time; epoch + 1 on takeover | Duplicates bounded by lag; late timers |
| 3.4 | "Exactly once?" | A two-part contract; idempotency keys | Tenants must do their part |
| 3.5 | Cost per tenant | Metering by cost driver | A billing-grade pipeline |
| 3.6 | Build or buy? | Compared on semantics, regions, price, people | A team, or a vendor's limits |
R3.5 Global Architecture
Synthesizing vector architecture diagram...
Any region's API accepts a request and forwards it to the region that currently owns the workflow's placement group. Each region runs its own cells; data for HOME_ASYNC and STRICT workflows exists in all three regions, but only the owner of a group acts on it. The small ownership table is the only thing every region must agree on synchronously.
Where each piece lives
| Component | Scope | Notes |
|---|---|---|
| API | Every region | Authenticates the tenant, applies API quotas, forwards to the owning region |
| Engine (workflow tasks, replay) | Per cell | Keeps recently used histories in memory, so most decisions need no history read |
| Timer service | Per cell | The Round 2 machinery: due-index, loaders, Valkey buckets, dispatchers |
| Executors and task queues | Per cell | Four priority queues as SQS fair queues keyed by tenant; pools as in Round 2 |
| History, state and job tables | Per cell, global tables | MREC for HOME_ASYNC, MRSC for STRICT, single-region tables for SINGLE_REGION |
| Ownership | Global, MRSC | One item per placement group; renewals every 5 s |
| Placement | Global | Tenant → cell; changes rarely |
| Metering and lake | Per region to S3 | Athena for reports and bills |
Trace 1: a saga with compensation is the sequence in step 3.1.
Trace 2: a region failover is the sequence in step 3.3, followed by recovery:
Synthesizing vector architecture diagram...
Trace 3: a tenant flood contained
Synthesizing vector architecture diagram...
T's midnight never touches the shared cell. Inside a shared cell, a tenant that bursts past its quota is pushed later in its own buckets and ranked behind quieter tenants in the fair queues.
R3.6 Numbers and Cost
Targets
| Quality | Target | Why this number |
|---|---|---|
| Dispatch lag, per tenant | P99 < 1 s for strict jobs; within the window for the rest, while the tenant is within quota | Round 2's SLO, now measured per tenant |
| Workflow step latency | P99 < 2 s from a step's completion to the next step's task being available (HOME_ASYNC) | Tenants chain steps; engine overhead must stay small |
| Region failover | Timers back on schedule within 2 minutes; duplicates bounded by replication lag | From the scope raise |
| Availability | 99.99% per region for the API and dispatch, and a failover plan instead of a bigger number | DynamoDB global tables carry a 99.999% availability SLA (single-Region tables: 99.99%), but that is one dependency's number; ours also depends on SQS, Valkey, compute and our own code, so we don't claim five nines |
Volumes (the interviewer's figures; the rest are our assumptions)
| Item | Math | Result |
|---|---|---|
| Job runs | 1B a day ÷ 86,400 | 11,574/s average |
| Workflow starts | 50M a day ÷ 86,400 | 579/s |
| Workflow steps | 50M × 10 | 500M a day, 5,787/s |
| Open workflows | 50M a day × 30 days (Little's law) | 1.5B |
| Pending timers | 1.5B workflows waiting on ~1 timer each + 200M active jobs | ~1.7B |
| Timers per region | 60% / 25% / 15% of 1.7B | us-east-1 1.02B, eu-west-1 425M, ap-northeast-1 255M |
| Timers firing | 1B job runs + ~300M workflow timers (sleeps, signal deadlines, step timeouts) a day | 1.3B a day; us-east-1 fires 780M a day, 9,028/s |
| Actions in us-east-1 | 60% × (11,574 + 5,787) | ≈ 10.4K/s |
History growth
| Item | Math | Result |
|---|---|---|
| Per workflow | 10 steps × 3 events × ~400 B + ~1 KB for start, signals and close | ~13 KB |
| Written per day | 50M × 13 KB | 650 GB/day (us-east-1: 390 GB) |
| Open workflows' histories | 1.5B × ~6.5 KB (halfway through their steps, on average) | 9.75 TB |
| Closed, kept 7 days | 50M × 7 × 13 KB | 4.55 TB |
| Workflow state items | 1.5B × ~1 KB | 1.5 TB |
| Job tables | Round 2's ~1.25 TB × 2 | ~2.5 TB |
| Timer indexes | 1.7B × ~100 B | ~0.2 TB |
| Per home copy | ≈ 18.5 TB | |
| With replicas | 80% of it (HOME_ASYNC + STRICT) has 2 more copies: 18.5 × (1 + 0.8 × 2) | ≈ 48 TB |
A workflow that runs for months would grow its history without limit, so a definition can continue as new: close this execution and start a fresh one with the current state as input, and an empty history. Hosted engines cap histories for the same reason: Temporal, for example, warns at 10,240 events or 10 MB and terminates at 51,200 events or 50 MB. We cap ours at 10,000 events or 10 MB (our choice).
DynamoDB writes (home region first; replicated writes are billed as replicated write units, at the same list price as ordinary writes since November 2024)
| Kind | Writes each | Home writes per day |
|---|---|---|
| Job runs (as in Round 2) | 5 | 1B × 5 = 5.0B |
| Heartbeats: 40K long runs, every 10 s | 3 | 4,000/s × 3 × 86,400 = 1.04B |
| Workflow steps: 3 events + 3 state updates + ~2 timer-index changes | ~8 | 500M × 8 = 4.0B |
| Starts, closes, signals | 2 | (100M + 100M) × 2 = 0.4B |
| Home total | 10.44B a day | |
| Replicated copies | 80% × 2 more regions = 1.6 × home | 16.7B a day |
| All regions | ≈ 27.1B a day |
Replication volume: 16.7B replicated writes a day of items averaging about half a kilobyte, roughly 8 TB a day crossing regions. DynamoDB doesn't charge data transfer for global-table replication; the cost is the replicated write units. If we replicated ourselves instead, the far side would still pay the same writes, plus inter-Region transfer, which is $0.02/GB out of us-east-1 and differs by Region (it's higher out of several Asia Pacific Regions), so check each pair.
Failover headroom. When a region fails, its groups are split between the two survivors. Each region's compute must cover its own load plus half of the largest region it might inherit:
| Region | Own | Plus half of | Must handle |
|---|---|---|---|
| us-east-1 | 60% | eu-west-1 (25%) | 72.5% |
| eu-west-1 | 25% | us-east-1 (60%) | 55% |
| ap-northeast-1 | 15% | us-east-1 (60%) | 45% |
| Total provisioned | 172.5% of the global load |
That 72.5% overhead applies to what we must pre-provision: compute and Valkey. DynamoDB on-demand follows the load (with warm throughput set ahead), and the survivors already hold replicas of the failed region's data. On-demand tables default to a maximum of 40,000 write request units a second per table and per GSI; our busiest job tables need far more (~190K at the top of the hour, and more in a survivor after a failover), so we raise that quota for every table, in every replica Region, before launch.
Failover timing (from step 3.3): a failure just after a renewal leaves 20 s on the lease, and takeover waits 5 s more: ~25 s before a survivor owns the groups. Then, per survivor: overdue timers ~25 s × 9,028/s ÷ 2 ≈ 113K, plus the next 10 minutes 9,028 × 600 ÷ 2 ≈ 2.7M timers loaded at 50K/s ≈ 55 s. Timers are back on schedule about 1.5 minutes after the failure, inside the 2-minute target.
Rough monthly cost (us-east-1 list prices applied everywhere for simplicity; Tokyo and Ireland prices differ, so check the AWS Pricing Calculator before quoting)
| Line | Math | ≈ Monthly |
|---|---|---|
| DynamoDB writes, job tables (on-demand: spiky) | home 6.04B + replicated 9.66B = 15.7B a day × 30 × $0.625 per million | $294,000 |
| DynamoDB writes, workflow tables (provisioned: smooth, ~60% utilization assumed) | home 4.4B + replicated 7.04B = 11.44B a day × 30 = 343B × $0.18 ÷ 0.6 = $0.30 per million | $103,000 |
| DynamoDB storage | 48 TB × $0.25/GB | $12,000 |
| DynamoDB reads | cache misses on history (10% of workflow tasks, ~2 read units each) + job-side reads ≈ 170M units/day | $640 |
| SQS | ~1B requests a day (jobs 0.5 per run, workflow tasks 2 per step at ~0.5) × 30 = 30B; every queue is a fair queue (messages carry MessageGroupId), which adds $0.10 per million: $0.50 per million | $15,000 |
| ElastiCache for Valkey | one cluster per cell (3 shards × 2 cache.r7g.large); assume 12 cells (7 in us-east-1, 3 in eu-west-1, 2 in ap-northeast-1, dedicated cells included): 72 nodes × ~$0.175/h × 730 h. Memory needs are small; the count comes from isolation | $9,200 |
| Compute | Round 2's ~$2,470 × 4 (twice the job runs, plus steps) × 1.25 (engine services) × 1.725 (failover headroom) | $21,000 |
| Ownership table (MRSC) | 768 groups × 1 renewal / 5 s ≈ 154 writes/s × 3 copies ≈ 1.2B a month × $0.625 per million | $750 |
| History lake | Firehose 1.65 TB/day × 30 × $0.029/GB ≈ $1,440; S3 ~129 TB at 13 months' retention × $0.023 ≈ $2,960 | $4,400 |
| Metering, CloudWatch, Route 53 | estimate | $5,500 |
| Total | ≈ $465,000 |
Cost per unit
| Unit | What we allocate | Result |
|---|---|---|
| Per million workflow steps | workflow writes $103K + ~85% of storage $10K + half of SQS and compute $18K + share of lake $1.7K and Valkey $2.1K ≈ $135K ÷ 15,000M steps | ≈ $9.00 |
| Per million job runs | job writes $294K + the rest of the direct lines (storage, SQS, compute, Valkey, lake) ≈ $324K ÷ 30,000M runs | ≈ $10.80 (Round 2: $4.50) |
Three lessons:
- Replication is the bill. 62% of all writes are replicas. A job costs 2.4× more per run than in Round 2 almost entirely because it's written in three regions.
- The levers are per policy. Replicating to one partner instead of two cuts replicated writes in half; jobs that can wait out a regional outage can be
SINGLE_REGIONand cost a third; longer windows let job tables move to provisioned capacity at ~$0.30 instead of $0.625 per million. Reserved capacity can't be bought for replicated write capacity, so it only helps the home-region share. - Tenants must see this. A
STRICTworkflow and aSINGLE_REGIONone are different products at different prices; metering (step 3.5) makes that visible.
R3.7 Trade-Offs
| Decision | Why we chose it | What we concede |
|---|---|---|
| A workflow engine over job chains | The history is one durable record of a multi-step process; replay makes any host able to continue it; compensation has the facts it needs. | Deterministic definitions, version rules, and an engine to run. |
| Home region with epoch takeover over multi-writer timers REL 13 | Only one region acts on a timer at a time; duplicates are limited to the replication lag, not doubled every day. | A failover gap of ~25 s plus loading, and a duplicate window after a crash. |
| Region policy per workflow | Tenants pay for the recovery they need: zero loss (STRICT), seconds (HOME_ASYNC), or none (SINGLE_REGION). | Three behaviors to explain, test and price. |
| Cells, with dedicated cells for the largest tenants REL 10 | A tenant's flood or a bad deploy is contained to one cell; each cell stays at a size we've run. | Placement, tenant migration, and more stacks to operate. |
| Build over buy (at this volume) COST 11 | About $9 per million steps vs $25+ managed at list; three-region failover and tenant fairness built in. | A platform team, on call around the clock. Below about a twentieth of this volume, buy. |
Closing the loop on the loop's question. "How do we run every job at the right time, and make it safe when machines fail and a job runs twice?"
- Round 1: at the right time within a minute, with at-least-once execution and an idempotency key.
- Round 2: at the right time within a second, where "right time" is honestly redefined for 10M simultaneous jobs as "within your window", and a stale worker can't touch our records.
- Round 3: our own state changes happen exactly once, within an owner and across regions for
STRICT; everything we deliver to tenant code happens at least once, with a key that makes it safe; and after a region failure the repeats are bounded by replication lag, counted, and reported.
R3.8 Failure Modes
| Trigger | What you'd see | How the design responds | Drill |
|---|---|---|---|
| A region outage | Renewals from one region stop; its tenants' API calls fail over by DNS. | Self-fence, takeover at lease end + 5 s, overdue timers first (step 3.3). Survivors were sized for it (R3.6: 172.5%). | Replication multi-region consistency |
| A replay storm after failover | The new owner has cold caches; every open workflow's next decision needs its history read, and read units spike. | Workflows are woken in order of their next due time, at a capped rate per cell; histories are read with eventually consistent reads (half the cost); warm throughput for reads is raised in every region during a game day, so we know it holds. | – |
| A poison workflow | One workflow's decision fails on every replay (a bad definition version, corrupt input); an engine host retries it forever. | After 5 failed workflow tasks in a row it's marked STUCK, its timers are paused, and the tenant is alerted. It never blocks other workflows, because decisions for different workflows run independently. | – |
| History grows without bound | A subscription workflow that has run for two years has 60,000 events; each replay is slow and the item collection is huge. | Continue-as-new at 10,000 events or 10 MB (R3.6). The API rejects new events beyond the hard limit with a clear error rather than slowing the whole cell. | – |
| Clock skew in one region | A region's hosts drift 3 s ahead: it would stop late and take over early. | The 5 s margins on both sides of each lease absorb it; the clock-offset alarm fires at 1 s and blocks that region from taking over any group until fixed. | – |
| A compensation that never succeeds | A refund keeps failing for a day. | COMPENSATION_STUCK: tenant DLQ, page, reconciliation report (step 3.1). | Saga orchestration compensating failure |
R3.9 Runbook
Signals, per cell and per tenant OPS 8 · REL 6
| Signal | Alarm | Severity | First action |
|---|---|---|---|
| Dispatch lag: actual start minus effective time, P99, per priority | > 1 s for 3 min | P1 | Which cell and shards? Dispatcher leases, Valkey health, queue age |
| Bucket backlog: members due more than 2 s ago | > 10K for 1 min | P2 | Dispatchers behind or a Valkey shard slow |
| Claim conflicts per second | > 5% of claims | P3 | Duplicate sends: two dispatchers on one shard, or a reload storm |
| Lease expiries per minute | > 1% of running long jobs | P1 | Reaper paused; check our heartbeat API first |
| DLQ inflow and oldest message age, per tenant | CRITICAL: any; others: 10× normal | P1/P2 | Look at the failed runs and their targets |
| Per-tenant lag | a within-quota tenant's lag > its SLO for 5 min | P2 | Is its cell overloaded, or is a neighbor over quota? |
| Replication latency, per table and region pair | > 5 s for 5 min | P2 | The duplicate window is growing; check the region |
| Ownership renewal failures, per region | 2 in a row for any group | P1 | Is the region isolated? Expect takeover in ~25 s |
| Clock offset per host | > 1 s | P2 | Replace the host; block takeovers from its region |
Procedure: a stampede bigger than planned
- Confirm it's a known top-of-hour burst: bucket sizes for the coming minutes, per tenant.
- Check the executor pools were pre-scaled and that warm throughput covers the new peak; raise both if not.
- If a tenant is over quota, its members are already being pushed later; if the burst is inside quotas and still too big, lengthen LOW windows for that minute (forward only) and tell the affected tenants.
- Never shorten a window or run anything early.
Procedure: redrive after a fix
- Confirm the cause is fixed with a handful of manual redrives.
- Redrive the rest through the redrive API with
max_per_secondwell below the target's normal rate; AIMD will raise it if the target is healthy. - For the poison-message queue (messages that crashed executors), redrive with SQS's message move task once the executor bug is fixed.
Procedure: evacuate a region (planned, for a degraded region)
- Stop placing new tenants in it. 2. Move its placement groups one at a time: the owner drains, waits for replication latency to reach its normal level, then hands the lease over under epoch + 1. 3. Watch dispatch lag in the receiving regions after each group. 4. Leave the region passive until the cause is understood; fail back the same way.
Game days. Monthly, AWS Fault Injection Service pauses global-table replication for one region (it supports isolating a replica for both MREC and MRSC tables) and we confirm self-fencing, takeover and the reconciliation report. Quarterly, a full evacuation during business hours. Every surprise gets a blameless Correction of Error with owners and dates. REL 12 · OPS 11
Go deeper: CLI playbook
Plain commands an on-call engineer runs, one at a time. Replace the names with real ones.
text# 1. Depth, in-flight and delayed counts for one priority queue in one cell aws sqs get-queue-attributes --queue-url https://sqs.us-east-1.amazonaws.com/111122223333/cell3-jobs-default --attribute-names ApproximateNumberOfMessages ApproximateNumberOfMessagesNotVisible ApproximateNumberOfMessagesDelayed # 2. Replication latency from us-east-1 to eu-west-1 for one cell's history table aws cloudwatch get-metric-statistics --region us-east-1 --namespace AWS/DynamoDB --metric-name ReplicationLatency --dimensions Name=TableName,Value=cell3-history Name=ReceivingRegion,Value=eu-west-1 --start-time 2026-09-28T00:00:00Z --end-time 2026-09-28T00:15:00Z --period 60 --statistics Average Maximum # 3. Redrive the poison-message queue back to its source queue at a capped rate aws sqs start-message-move-task --source-arn arn:aws:sqs:us-east-1:111122223333:cell3-jobs-poison --max-number-of-messages-per-second 200 # 4. Watch that redrive aws sqs list-message-move-tasks --source-arn arn:aws:sqs:us-east-1:111122223333:cell3-jobs-poison # 5. Raise a table's warm throughput before a known burst (indexes are set separately; # above the 40,000 write units/s on-demand default, raise the per-table quota first, in every replica Region) aws dynamodb update-table --table-name cell3-jobs --warm-throughput ReadUnitsPerSecond=20000,WriteUnitsPerSecond=250000 # 6. Scheduled scale-out of an executor pool at :57 every hour aws application-autoscaling put-scheduled-action --service-namespace ecs --resource-id service/cell3/executor-default --scalable-dimension ecs:service:DesiredCount --scheduled-action-name top-of-hour-out --schedule "cron(57 * * * ? *)" --scalable-target-action MinCapacity=68,MaxCapacity=90
Without --destination-arn, the move task sends messages back to the queue they originally came from, which is what we want for the poison-message queue.
R3.10 Pillar Check
| Pillar | What Round 3 adds |
|---|---|
| Reliability | Durable workflows with sagas; region failover by MRSC ownership leases and epochs; failover headroom sized and priced (172.5%); cells; game days that isolate a region. REL 10 · REL 12 · REL 13 |
| Performance Efficiency | Latency-based routing to the nearest API; engines cache hot histories so most decisions need no read; STRICT latency only where it's bought. PERF 3 · PERF 4 |
| Security | Tenant isolation by cell and by IAM-scoped access to each tenant's items; per-tenant signing secrets for task delivery; SINGLE_REGION for tenants whose data must stay in one region; history encrypted with KMS keys per region. SEC 3 · SEC 7 · SEC 8 |
| Cost Optimization | ≈ $465K a month derived; replication named as the driver; unit costs per step and per run; per-tenant metering; build vs buy compared at list prices. COST 3 · COST 5 · COST 11 |
| Operational Excellence | Signals per cell and per tenant with first actions; procedures for stampedes, redrives and evacuations; COEs. OPS 8 · OPS 10 · OPS 11 |
| Sustainability | Region policies avoid storing and writing three copies of data that doesn't need them; continue-as-new and 7-day history retention in DynamoDB, with GZIP in S3 after; executors scaled to the hour. SUS 2 · SUS 4 |
R3.11 Round 3 Rubric and Follow-Ups
What an architect (L7) answer adds over L6
- Replaces job chains with a durable history and replay, and knows why decisions must be deterministic.
- Chooses sagas over two-phase commit for a reason, and plans for a compensation that fails.
- Isolates tenants with quotas, fair dispatch and cells, and has a plan for the tenant that outgrows a shared cell.
- Makes region failover safe without trusting DynamoDB's clock or local conditions: strongly consistent ownership, self-fencing by time with margins, and an honest duplicate window.
- States the exactly-once contract precisely, the way AWS does for Step Functions.
- Prices the system, finds the driver (replicated writes), and turns it into per-policy, per-tenant prices.
- Treats build vs buy as a volume-and-team decision, with the crossover point named.
Follow-up questions
-
"A tenant sends the
picked_upsignal at the same moment the 3-day timer fires. Which wins?" Answer: whichever event is appended first. Both try to append the next event number to the history; the key is unique, so one succeeds and the other retries on the next number. When the loser's host replays, it sees the winner's event and acts on that: if the signal won, the timer's event is recorded as ignored; if the timer won, the late signal is recorded and reported as arriving after the timeout. Either way the decision is made once, from the history. -
"Why not run the ownership leases on DynamoDB TTL, so a dead region's lease item just disappears?" Answer: TTL deletes items in the background, typically within a few days of expiry, never at a deadline, so a lease could outlive its owner by days. MRSC tables don't support TTL anyway. Lease expiry must be a timestamp that readers compare against their own clocks, with margins.
-
"Your biggest tenant asks for its own region." Answer: it gets a dedicated cell in any of our regions today; a new region is a business decision priced from R3.6: its own cells, a partner region for replication, and headroom for failover. If it only wants its data kept in one region,
SINGLE_REGIONin its home region does that for the workflows it chooses, at the cost of waiting out that region's outages.
Interview gotchas from this round's wrong answers
| Gotcha | Why it's wrong |
|---|---|
| "Chain jobs with callbacks" | No record of earlier steps, no compensation, and a crash between links loses the rest. |
| "2PC across the services" | Payment providers don't take part, and locks can't be held for days. |
| "Global FIFO is fair" | The biggest, earliest tenant goes first; everyone else waits. |
| "Fire timers in two regions to be safe" | Every timer fires twice, every day. |
| "Global tables will stop a double claim" | Conditions on eventually consistent global tables are checked locally; last writer wins. |
| "Leases with TTL" | TTL deletes typically within days, and MRSC tables have no TTL. |
| "We guarantee exactly once" | Only for our own state; delivery to tenant code is at least once. |
Loop Closer: Interview Strategy for All Three Rounds
How to Run Each 60-Minute Round
| Time | Round 1 | Round 2 | Round 3 |
|---|---|---|---|
| 0–5 min | Scoping questions: precision, duration, failures, time zones | Restate Round 1 in 60 seconds | Restate Round 2 in 60 seconds |
| 5–15 min | Requirements + API | Scope raise → what breaks | Scope raise → what breaks |
| 15–40 min | Steps 1.0–1.6 | Steps 2.1–2.6 | Steps 3.1–3.6 |
| 40–50 min | Numbers (derive the busiest minute) + trade-offs | Numbers (derive the midnight peak), cost, trade-offs | Numbers, failover headroom, cost, build vs buy |
| 50–60 min | Failures + pillar check | Failures + pillar check | Failures, runbook, pillar check |
For how to spend a single 45-minute round, see the 45-minute interview blueprint. Related loops: the distributed message queue goes deeper on queues, and the distributed rate limiter on protecting shared targets.
The Two Sentences That Matter Most
- Opening a round: "Before I design, let me ask a few scoping questions."
- When the scope is raised: "Here's what breaks in the current design, and here's the order I'll fix it in."
And the one sentence specific to this system: "We run at least once, make our own state changes exactly once with conditional writes, and give every target an idempotency key so a repeat does no harm."
Well-Architected Review Sheet
| Pillar | Question you'll hear | One-sentence answer | Round | Backed by |
|---|---|---|---|---|
| Reliability | "What if a worker dies mid-job?" (REL 11) | Its lease expires, the job is reclaimed under a higher epoch, and the target's idempotency key absorbs the repeat. | 1–2 | Steps 1.3, 2.4 |
| "What if an AZ goes down?" (REL 10) | Pools are sized ×1.5, dispatcher shards move in ~10 s, Valkey promotes replicas, and DynamoDB and SQS are Regional. | 2 | R2.6, R2.8 | |
| "What if a region goes down?" (REL 13) | Its placement groups move to the survivors under a new epoch in ~25 s, timers are back in ~1.5 minutes, and repeats are bounded by replication lag. | 3 | Step 3.3, R3.6 | |
| "How do you protect your dependencies?" (REL 5) | Backoff with full jitter, AIMD per target, circuit breakers, and a reaper that won't turn our outage into duplicates. | 1–2 | Steps 1.4, 2.5, 2.6 | |
| "How do you test it?" (REL 12) | Staging delivers 1% of calls twice; game days isolate a region's replicas and evacuate it. | 2–3 | R2.9, R3.9 | |
| Performance | "How do you find due jobs fast?" (PERF 3) | A sparse index by minute and shard, with the next 2 hours in Valkey buckets scanned every 250 ms. | 1–2 | Steps 1.1, 2.1 |
| "What's your peak?" (PERF 2) | 10M due at midnight UTC; with 300 s windows that's ~36.8K runs a second, derived from the cron shape, not a multiplier. | 2 | Step 2.3, R2.6 | |
| Cost | "What does it cost, and what drives it?" (COST 5) | ≈ $145, ≈ $67K and ≈ $465K a month; DynamoDB writes drive it, and in Round 3 replicated writes. | 1–3 | R1.7, R2.6, R3.6 |
| "Why not provisioned capacity?" (COST 7) | Below ~29% utilization on-demand is cheaper; longer windows flatten the load and change the answer. | 2 | R2.6, R2.7 | |
| "Should we just buy it?" (COST 11) | At one app's scale, yes; at the company's with a 1 s SLO, no; as a product at our volume, build, possibly on open-source Temporal. | 1–3 | R1.8, R2.7, step 3.6 | |
| "What does each tenant cost?" (COST 3) | Meter runs, steps, timer-hours, history bytes by region policy, and executor time, priced from unit costs. | 3 | Step 3.5 | |
| Operations | "How do you know it's healthy?" (OPS 8) | Dispatch lag against the effective time, bucket backlog, lease expiries, DLQ age, per-tenant lag and replication latency, each with a first action. | 2–3 | R2.10, R3.9 |
| "How do you recover failed work?" (OPS 10) | Failed runs keep their records; a redrive API re-runs them at a capped rate once the cause is fixed. | 2–3 | R2.3, R3.9 | |
| Security | "Can a job call anything it wants?" (SEC 5) | Only registered targets, approved by their owners, with signed requests. | 2 | R2.10 |
| "Who can see a tenant's jobs?" (SEC 3) | Only callers with that tenant's IAM-scoped identity; cells add a second boundary for the largest tenants. | 2–3 | R2.10, R3.10 | |
| "Can a tenant's data stay in one region?" (SEC 7) | Yes, SINGLE_REGION, at the price of waiting out that region's outages. | 3 | R3.3 | |
| Sustainability | "How do you avoid waste?" (SUS 2) | Executors scaled to the hour, longer windows for flexible jobs, and replicas only where the policy asks for them. | 2–3 | R2.10, R3.10 |
Rubric Across Levels
| Dimension | L5 (Round 1) | L6 (Round 2) | L7 (Round 3) |
|---|---|---|---|
| When | Jobs table with a sparse due index; cron in the job's zone with DST policies; misfire policies. | Minute-and-shard index plus Valkey buckets; the peak derived from the cron shape; forward-only hashed windows. | Timers and signals as history events; failover that loads overdue timers first. |
| Who | Conditional claims with a lease; knows read-then-write fails. | Epoch fencing on every write; heartbeats; an adaptive reaper; knows fencing only protects systems that check it. | Region ownership by MRSC leases, self-fencing by time with margins, epoch on takeover. |
| What if it fails | Backoff with full jitter, max attempts, DLQ; idempotency keys. | Bulkheads per priority; AIMD and circuit breakers per target; the SQS receive-count trap. | Sagas with compensations that retry and escalate; a bounded, reported duplicate window. |
| Numbers and cost | Derives the busiest minute; honest "buy" at small scale. | Rejects the ×5 multiplier; costs writes per run; on-demand vs provisioned break-even; the window as a cost lever. | Failover headroom priced; replication as the driver; unit costs per step and run; build vs buy crossover. |
| Evolving under new scope | Builds from one crontab, one problem at a time. | Opens with "what breaks", fixes in a stated order, keeps the table as truth. | Changes the contract (exactly-once where it's true), the operating model (cells, policies, metering) and the product. |