Design a Distributed Job Scheduler & Task Execution Queue
1. Problem Statement & Scope Clarification
System Mission
Design an enterprise-grade, distributed, fault-tolerant Job Scheduling and Task Execution Platform (equivalent to Temporal, Quartz Scheduler, or AWS EventBridge Scheduler) capable of scheduling, orchestrating, and executing billions of delayed, recurring (cron), and one-off asynchronous workloads with at-least-once dispatch and fenced, effectively-once execution (a task can be dispatched twice after a crash, but fencing tokens guarantee only one worker's side effects are ever committed), dynamic worker auto-scaling, conditional write lease fencing, and end-to-end execution observability.
Functional Requirements
- Diverse Schedule Types (
POST /v1/jobs):- One-Off Delayed Tasks: Execute once at an exact future epoch timestamp (e.g., in 45 minutes).
- Recurring Cron Workloads: Execute repeatedly based on standard 5-to-6-field cron expressions (e.g.,
0 0 * * *). - Instant Asynchronous Tasks: Enqueue immediate background executions with priority routing.
- Prioritized Workload Dispatch: Segregate jobs into priority tiers (
CRITICAL,HIGH,DEFAULT,LOW) to guarantee that critical payment workflows never suffer head-of-line blocking behind bulk report generation. - Task Lifecycle Tracking: Real-time tracking of task states (
SCHEDULED,CLAIMED,DISPATCHED,RUNNING,SUCCEEDED,FAILED,CANCELLED). - Resilient Retries & Dead-Letter Queues (DLQ): Configurable exponential backoff with full randomized jitter and automatic DLQ routing after max attempts.
- Worker Heartbeating & Lease Fencing: Long-running jobs renew heartbeats every 10 seconds. If a worker dies mid-execution, its lease expires, and the task is safely reclaimed by another node without split-brain duplicate execution.
Non-Functional Requirements (SLAs & SLOs)
- High Availability: ("five nines") uptime SLA for the scheduling ingestion and dispatching plane.
- Timing Accuracy: Dispatched within of the target execution time ().
- Task Durability: (11 9s) zero task loss guarantee prior to execution.
- Throughput & Scale: Support active scheduled jobs, daily completed task runs, and peak dispatch bursts of .
2. Capacity & Scale Estimation (Back-of-the-Envelope Math)
Ingestion & Dispatch Scale
- Total Active Scheduled Tasks: ().
- Daily Completed Workloads: ().
- Average Dispatch Throughput:
- Peak Dispatch Throughput ( Top-of-Hour Multiplier):
Recurring cron jobs spike massively at the start of hours (
00:00,01:00):
Storage Footprint (Active State vs Historical Lake)
- Task Definition Size:
job_id(),tenant_id(),target_url(),payload_json(),cron_expression(),next_run_epoch_ms(),retry_policy(), metadata pointers () .
- Active Task Store (Amazon DynamoDB):
- Monthly Execution History Lake (Amazon S3 via Kinesis Firehose):
Redis Time-Bucket Memory Sizing
To prevent a single giant Redis Sorted Set (ZSET) from bottlenecking memory and CPU, scheduled tasks are partitioned into 1-Minute Time Buckets (schedule_bucket:YYYYMMDDHHMM):
- Jobs due within the next 2 hours are buffered in Redis:
- Each ZSET entry stores
member: job_id() andscore: epoch_ms(). - With Redis jemalloc and skip-list pointer overhead : An Amazon ElastiCache Redis Cluster with replicas easily hosts this dataset in memory.
3. High-Level Architecture & AWS Component Mapping
Synthesizing vector architecture diagram...
Why Two-Phase Scheduling (Redis ZSET + DynamoDB Fencing)?
Direct database polling (WHERE run_at <= NOW()) causes crippling table scan bottlenecks at high scale. Sharding upcoming tasks into discrete 1-minute Redis Sorted Sets enables in-memory range queries. Atomic DynamoDB conditional updates with monotonically increasing fencing tokens eliminate duplicate executions even across scanner crashes and network partitions.
Data Flow Walkthrough
- Job Scheduling Ingress: A client issues
POST /v1/jobswith a schedule and target payload. The Job Management Service generates a unique 128-bitjob_id, persists the master record to DynamoDB, and (if the execution time is within the next 2 hours) inserts(score=run_at_epoch_ms, member=job_id)into the appropriate 1-minute Redis Sorted Set. - Time-Slice Scanning: The Scanner Fleet uses consistent hashing to assign 1-minute time buckets across scanner nodes. Every , scanners execute
ZRANGEBYSCOREfor tasks withscore <= now_ms. - Two-Phase Claim & Fencing: For each due job, the scanner executes a conditional DynamoDB update (
status = 'DISPATCHED', version = version + 1). Only the scanner that successfully increments the version token claims the job, preventing duplicate dispatches. - Priority Queueing: Claimed jobs are dispatched to Amazon SQS queues segmented by priority tier (
CRITICAL,HIGH,DEFAULT). - Execution & Heartbeat Loop: Workers pull messages from SQS, transition the job state to
RUNNINGin DynamoDB, and start a background heartbeat thread that renews the 30-second execution lease every 10 seconds. Upon successful completion, the task status transitions toSUCCEEDED.
Concrete Step-by-Step Request Walkthrough: Tracing a Scheduled Delayed Task
| Step # | Event / Action | Component State | Distributed Transition | Output / Response |
|---|---|---|---|---|
| 1 | Client calls POST /v1/jobs(Target execution delay: 15 mins) | HTTP request with idempotency token and HMAC signature | Gateway validates token; checks tenant quota and rate limiting bucket | Ingress accepted; HTTP 202 Accepted with assigned job_id |
| 2 | Master task record persisted | DynamoDB JobSchedulerTablereceives master execution item | status = 'SCHEDULED',next_run_epoch_ms = Now + 900,000 | Record committed in ; DynamoDB Streams event published |
| 3 | Indexing into Redis time bucket | Time bucket computed:schedule:202609161215 | ZADD schedule:2026091612151767226500000 job_9981 | Redis returns 1 (Stored in memory);4-hour TTL assigned to bucket |
| 4 | Clock advances to target minute | Scanner Node 02 assigned bucketschedule:202609161215 | Scanner runs ZRANGEBYSCORE0 (now_epoch_ms) LIMIT 0 500 | Task job_9981 identifiedas ready for immediate dispatch |
| 5 | Atomic claim with fencing token | Scanner executes conditional update:status == 'SCHEDULED' | Sets status = 'DISPATCHED',fencing_epoch = 104, version = 2 | Lock acquired; task deleted from active Redis ZSET |
| 6 | Priority SQS dispatch & worker claim | Scanner pushes task payload to SQS_HIGH_PRIORITY | Worker consumes message; starts dedicated 10-second background heartbeat thread | Worker invokes webhook:billing.internal/invoices |
| 7 | Completion commit & lakehouse sync | Target API returns HTTP 200 OK with confirmation receipt | Worker updates DynamoDB: status = 'SUCCEEDED';heartbeat timer cleanly cancelled | Execution archived to S3 via Firehose in background |
4. API Interface Design & Wire Protocol
1. Schedule a Job (POST /v1/jobs)
httpPOST /v1/jobs HTTP/1.1 Host: scheduler.production.aws.internal Authorization: Bearer eyJhbGciOi... Idempotency-Key: idemp_job_8492019481 Content-Type: application/json { "tenant_id": "tenant_enterprise_01", "job_name": "monthly_billing_aggregation", "schedule_type": "CRON", "cron_expression": "0 2 1 * *", "timezone": "America/New_York", "priority": "HIGH", "target": { "type": "HTTP_WEBHOOK", "endpoint": "https://billing.internal.aws/invoices/generate", "method": "POST", "timeout_seconds": 120, "headers": { "Content-Type": "application/json" } }, "payload": { "billing_cycle": "2026-09", "dry_run": false }, "retry_policy": { "max_retries": 3, "initial_interval_sec": 10, "backoff_multiplier": 2.0, "max_interval_sec": 300 } } Response: 201 Created Content-Type: application/json { "job_id": "job_9847291048", "tenant_id": "tenant_enterprise_01", "status": "SCHEDULED", "schedule_type": "CRON", "next_run_time_epoch_ms": 1767225600000, "next_run_time_iso": "2026-10-01T02:00:00.000Z", "created_at": "2026-09-16T12:00:00.000Z" }
2. Inspect Job Execution Status (GET /v1/jobs/{job_id})
httpGET /v1/jobs/job_9847291048 HTTP/1.1 Host: scheduler.production.aws.internal Accept: application/json Response: 200 OK Content-Type: application/json { "job_id": "job_9847291048", "status": "RUNNING", "last_heartbeat_at": "2026-09-16T12:00:10.000Z", "current_lease_worker_id": "worker-ecs-task-04", "current_attempt": 1, "fencing_token": 104 }
5. Data Models & Storage Architecture
1. Amazon DynamoDB Single-Table Design (JobSchedulerMasterTable)
PK (Partition Key) | SK (Sort Key) | Attributes | Description |
|---|---|---|---|
JOB#<job_id> | METADATA | tenant_id, schedule_type, cron_expr, target_config, payload, priority, created_at | Task master definition |
JOB#<job_id> | EXECUTION#current | status, next_run_ms, worker_lease_id, lease_expiry_ts, fencing_token, retry_count, version | The single mutable execution row (one active run per job); completed runs are archived to S3 via DynamoDB Streams, not kept as extra rows |
BUCKET#<epoch_minute> | JOB#<job_id> | next_run_ms, priority, status | Secondary index for database fallback scanning |
Global Secondary Indexes (GSIs)
GSI1(Reaper Index):GSI1-PK:status#shardwhereshard = hash(job_id) % 16(e.g.RUNNING#07)GSI1-SK:lease_expiry_ts- Allows background reaper daemons to discover abandoned tasks where
lease_expiry_ts < now_mswith 16 targeted range queries (one per shard, run in parallel). - Why the shard suffix: a GSI keyed on bare
statuswould put everyRUNNINGjob (millions of items, tens of thousands of writes per second as leases renew) into a single GSI partition, and DynamoDB caps any one partition at / . Write-sharding the key spreads the heartbeat load across 16 partitions; the reaper simply fans out its query.
2. Redis Time-Bucket Sorted Sets
- Key Pattern:
schedule:<epoch_minute>(e.g.,schedule:202609161205) - Data Structure:
ZSET- Score:
epoch_ms(exact execution timestamp) - Member:
job_id
- Score:
- Expiry: Key TTL set to 4 hours (
EXPIRE 14400) after the bucket minute completes, ensuring automatic memory cleanup.
Unlock Complete Architecture & Production Runbooks
You have explored the free architectural preview (~35%). Spend 1 Coin to unlock the remaining 6 production deep-dive sections for a full 24 hours.