System Design: Job Scheduler (10M Jobs/day, DAG Dependencies, Effectively-Once Execution)
Goal
A distributed job scheduler.
Scale:
- 10M jobs/day
- 100K concurrent in flight
- Sub-second dispatch, 99.99% availability
Features:
- Cron scheduling
- DAG dependencies
- Priority queues with aging
- Effectively-once execution
1. Final Architecture
🔒 Premium section
2. Problem Statement
A scheduler sounds easy (tasks, times, dependencies) and then breaks in the same three places at scale.
- Effectively-once is hard. Distributed locks expire mid-execution when a worker GCs or stalls on IO. A TTL alone cannot stop two workers from writing conflicting results. Even with perfect state dedup in the database, a downstream HTTP call may already have charged a card before the worker crashed.
- DAG failures cascade. One bad upstream stalls hundreds of downstream dependents. Retry budgets and condition edges must be explicit.
- Priority without starvation. A flood of low-priority analytics cannot block critical billing. The answer is separate queues with aging, not a FIFO with a priority field.
Scale targets.
- 10M jobs/day (~116/sec avg, 500/sec peak)
- 100K concurrent in flight
- 50K DAGs, avg depth 5, max depth 20
- Avg job 30 s, max 2 h
- Submit-to-worker-start latency <1 s p99 (warm pool; see NFR-03 for the full breakdown)
- 99.99% scheduler availability
What not to do.
- In-memory-only state loses every in-flight job on restart.
sleep(next_run - now)drifts, and it scales to one job per thread.- Unbounded retries flood logs and starve the worker pool.
- A single Kafka topic causes priority inversion inside partitions.
- Active-active schedulers without sharded ownership dispatch the same job twice.
- Recursive foreign keys on DAG edges lock up the DB for any 20-node resolution.
- Timeouts alone hide fast crashes. A 2-hour timeout will miss a 1-minute crash for 119 minutes.
- Parsing cron on every tick wastes CPU, and 10M rows amplify the waste.
- Wall-clock scheduling without a timezone breaks twice a year on DST boundaries.
- Infrastructure alone cannot guarantee idempotency. The scheduler has no way of knowing the job charged a card.
3. Functional Requirements
| ID | Requirement | Priority |
|---|---|---|
| FR-01 | Submit a job for immediate execution | P0 |
| FR-02 | Schedule a job for a specific future timestamp | P0 |
| FR-03 | Recurring jobs via cron (5-field + 6-field with seconds, timezone-aware) | P0 |
| FR-04 | DAG dependencies (on_success, on_failure, always) | P0 |
| FR-05 | Effectively-once execution (exactly-once state commit + idempotency-key handoff for side effects) | P0 |
| FR-06 | Four priorities (critical, high, normal, low) with aging | P0 |
| FR-07 | Retries with exponential / fixed / linear backoff | P0 |
| FR-08 | Cancel pending or running jobs | P1 |
| FR-09 | Timeout enforcement | P0 |
| FR-10 | Dead letter queue with inspection and manual replay | P1 |
| FR-11 | Status tracking (PENDING → SCHEDULED → QUEUED → RUNNING → terminal) | P0 |
| FR-12 | DAG-level progress view | P1 |
| FR-13 | Webhooks on completion / failure / timeout | P1 |
| FR-14 | Tags + metadata for filtering | P2 |
| FR-15 | Per-tenant rate limits | P1 |
4. Non-Functional Requirements
The three latency NFRs below (03a/b/c) measure different slices:
- NFR-03a is the scheduler's internal decide-to-produce path and is the page-on signal.
- NFR-03b measures end-to-end submit-to-start with a warm worker pool.
- NFR-03c is cold-start, bounded by Kubernetes pod startup.
| ID | Requirement | Target |
|---|---|---|
| NFR-01 | Throughput | 10M/day, 500/sec peak |
| NFR-02 | Concurrent in flight | 100K |
| NFR-03a | Scheduler decide → Kafka produce | <50 ms p99 |
| NFR-03b | Submit → worker starts (end-to-end, warm pool) | <1 s p99 |
| NFR-03c | Submit → worker starts (cold-start, pod spin-up) | <30 s p99 |
| NFR-04 | Execution guarantee | Effectively-once (state commit + target-side idempotency) |
| NFR-05 | Scheduler availability | 99.99% (52 min/year) |
| NFR-06 | Worker failure detection | <30 s via heartbeats |
| NFR-07 | DAG resolution for 20-node DAG | <100 ms |
| NFR-08 | Cron drift | <1 s from intended fire time |
| NFR-09 | Retention | Hot 7 d, warm 90 d, cold 1 y (S3) |
| NFR-10 | Scheduler RTO | <15 s |
| NFR-11 | Job RPO | 0 (durable after 202 Accepted) |
| NFR-12 | Max DAG depth | 20 |
| NFR-13 | Max DAG width per level | 100 |
5. Design Assumptions
🔒 Premium section
6. Expanded Architecture
🔒 Premium section
7. Technology Selection
🔒 Premium section
8. Back-of-the-Envelope
🔒 Premium section
9. Data Model
🔒 Premium section
10. API Design
🔒 Premium section
11. Effectively-Once Execution
🔒 Premium section
12. Leader Election
🔒 Premium section
13. DAG Resolution
🔒 Premium section
14. Cron Parsing & Clock Skew
🔒 Premium section
15. Priority Queues & Multi-Tenant Fairness
🔒 Premium section
16. Multi-Region Scheduling
🔒 Premium section
17. Bottlenecks
🔒 Premium section
18. Backpressure, Retries, Webhooks
🔒 Premium section
19. Failure Scenarios
🔒 Premium section
20. Operational Playbook
🔒 Premium section
21. SLOs and Error Budgets
🔒 Premium section
22. Security
🔒 Premium section
23. Key Takeaways
🔒 Premium section
24. Appendix
🔒 Premium section
Explore the Technologies
🔒 Premium section