Think Beyond the Happy Path
Real production systems aren't just about running short tasks on schedule — they're about resilience, complex coordination, and scaling under pressure.
Before diving into the material, take a moment to ask yourself:
- Do you know how to safely schedule 10k jobs per second, like Slack fanout notifications?
- Do you know how scheduler and queue design mitigate tiny gap between scheduled time and actual execution — critical for ad auctions or stock trades at market open?
- Do you know how to guarantee the right delivery semantics — at-least-once vs. exactly-once — in systems like Youtube email campaigns (avoiding duplicates) or Stripe payment retries (avoiding double charges)?
- Do you know how to handle dependencies and multi-step pipelines, such as Airflow ETL workflows in data warehouses or payment processing flows that must debit before credit?
- Do you know what happens if a 6-hour job crashes halfway through, like ML model training or a large-scale analytics backfill?
You don't need to answer all of these right away. A Senior/Staff+ Engineer doesn't stop at basic functionality — they anticipate edge cases, design for resilience, and push for production-grade reliability.
But great systems start with great questions. What would you ask next?
Problem Statement
Design a distributed job scheduler that supports immediate, delayed, and recurring job execution at scale. The system should:
- Support immediate job execution (run now).
- Support future and recurring schedules (cron, delayed).
- Provide observability into system health and job status. And the scale is up to 10,000 job executions per second.
Functional Requirements
FR1 – Scheduling (Immediate Execution)
The system must support running jobs immediately (run now).
FR2 – Scheduling (Future & Recurring Execution)
The system must support running jobs at a specific future time or on a recurring cadence (e.g., every Monday at 8:00 AM).
FR3 – Observability
The system must provide visibility into overall health, including queue depth, success/error rates, latency, and worker status.
Non-Functional Requirements
NFR1 – High Scalability
The system should support up to 10,000 job executions per second.
NFR2 – Low Latency
Jobs should be executed within 2 seconds of their scheduled time.
NFR3 – Correctness
Ensure at-least-once execution, with delivery semantics configurable (e.g., at-most-once, exactly-once) based on use case.
NFR4 – High Availability
The system must be highly available, favoring availability over strict consistency.
NFR5 – Fault Tolerance
The system must tolerate worker crashes, network partitions, and partial failures, with mechanisms for automatic retries, failover, and recovery without data loss.
Requirement Summary
| Functional Requirements (FRs) | |
|---|---|
| Name | Description |
1. Immediate Execution | Support running jobs immediately (run now). |
2. Future & Recurring Execution | Support jobs at specific future time or recurring cadence (cron). |
3. Observability | Provide visibility into queue depth, success/error rates, latency, and worker status. |
| Non-Functional Requirements (NFRs) | |
|---|---|
| Name | Description |
1. High Scalability | Support up to 10,000 job executions per second. |
2. Low Latency | Execute jobs within 2 seconds of scheduled time. |
3. Correctness | At-least-once execution with configurable delivery semantics. |
4. High Availability | Favor availability over strict consistency. |
5. Fault Tolerance | Tolerate crashes, partitions, and failures with automatic recovery. |
Core Entity
To tie these requirements together into a coherent architecture, we begin with two foundational building blocks: the Job and the Job Run.
Job — The Blueprint of Work
A Job is the reusable definition of what needs to be done. It defines:
- The action or task to perform (e.g., "send email", "run backup")
- Optional input parameters or configuration
- Its intended schedule (e.g., immediate, cron)
- Priority or policy metadata
A Job is not tied to a specific execution. Think of it as a template — it can be reused across different times, contexts, or users.
Job Run — A Specific Execution of That Job
A Job Run is a concrete instance of a job being executed. It defines:
- When the job is to be run (
available_atorscheduled_time) - The actual input payload or overrides
- Execution status (
queued,running,succeeded,failed) - Retry count, error codes, logs
If the Job is the blueprint, the Job Run is the actual build based on that blueprint.
For example:
- The Job is "generate sales report."
- A Job Run would be "generate sales report for Q2 2025 at 6:00 AM Monday."
Why the Separation Job vs Job Run Matters?
This separation between Job and Job Run is not just semantic — it's architectural. It enables:
- Reuse of logic across many scheduled times
- Clear tracking of retries and failures per execution
- Durable, auditable history of what was run and when
- Simplified scaling, since Jobs are relatively static and Runs drive most traffic
Interviewer Tip
Different companies might use slightly different terms for these concepts:
- "Task" and "Attempt"
- "Schedule" and "Execution"
- "Action" and "Instance"
That's fine. The important thing is that you clearly explain the distinction:
- The definition of the work to be done
- The specific execution of that definition, tracked and managed over time
In this lesson, we'll stick to Job and Job Run as our canonical terms.
| Jobs Table | |||
|---|---|---|---|
| Field | Type | Key | Description |
job_id | UUID | PK | Logical job identifier (the 'intent'). |
type | TEXT | NOT NULL | Job kind, e.g., email.send. |
payload | JSONB | NOT NULL | Parameters for the job. |
schedule_type | TEXT | NOT NULL | immediate / at / cron. |
next_fire_at | TIMESTAMPTZ | NULLABLE | Next due time (used for at/cron). |
priority | INT | NOT NULL DEFAULT 0 | Scheduling priority. |
status | TEXT | NOT NULL DEFAULT active | Lifecycle control. |
created_at | TIMESTAMPTZ | NOT NULL DEFAULT now() | Auditing/slicing. |
| Job_runs Table | |||
|---|---|---|---|
| Field | Type | Key | Description |
run_id | UUID | PK | Physical execution attempt ID. |
job_id | UUID | FK → jobs(job_id) | Which job this run belongs to. |
run_number | INT | NOT NULL | 1,2,3… (attempt/order). |
status | ENUM | NOT NULL | queued/running/succeeded/failed/canceled. |
queued_at | TIMESTAMPTZ | NOT NULL DEFAULT now() | When the run was enqueued. |
available_at | TIMESTAMPTZ | NOT NULL DEFAULT now() | When it becomes eligible to run (backoff). |
claimed_by | TEXT | NULLABLE | Worker ID that leased it. |
claimed_at | TIMESTAMPTZ | NULLABLE | Lease start. |
lease_expires_at | TIMESTAMPTZ | NULLABLE | When others may safely steal it. |
started_at | TIMESTAMPTZ | NULLABLE | Execution start. |
finished_at | TIMESTAMPTZ | NULLABLE | Execution end. |
error_code | TEXT | NULLABLE | Categorical failure reason. |
error_message | TEXT | NULLABLE | Debuggable failure info. |
High Level Design
FR1: Scheduler component to support immediate job execution

Click to expand
When designing a distributed job scheduler, it's easier to start by supporting immediate job execution —that is, jobs that should run as soon as they're created. This lays down the core scheduling loop without needing to deal with cron parsing or future timing calculations. By focusing on the now, we can build the foundation for more complex features like delayed and recurring jobs.
The essential flow for immediate jobs can be summarized as:
Record Intent → Enqueue → Lease → Execute → Record Outcome → Retry if Needed
Why start with immediate job?
Immediate jobs strip away the complexities of time-driven orchestration and allow us to build and test the fundamental execution flow. No clocks, no delays — just a direct path from request to execution. Once that path is stable, the same model can be extended to handle future or recurring work with minimal changes to the core loop.
Step-by-step data flow:
Step 1: Receive the Job Request
The client or upstream service issues a request to run a job immediately — no run_at , no cron , no future schedule. This defines the intent to execute. Nothing actually runs yet.
Step 2: Validate and Normalize
The scheduler authenticates the request and generates a unique job_id . It marks the job with schedule_type = immediate so that all downstream components know this is eligible to run right away.
Step 3: Atomically Write Job and Run
In a single transaction, the system:
- Inserts a new row into the
jobstable to describe what to run - Inserts a new row into the
job_runstable to represent this specific execution attempt, withstatus = queuedandavailable_at = now()
This atomic write ensures that either both rows exist (and the job can run) or neither (no inconsistent state).
Step 4: Enqueue the Job for Execution
Once the job is recorded and eligible, the scheduler publishes it to the job queue. This decouples job creation from execution and allows workers to pull jobs as needed — without polling the database.
Step 5: Worker Claims and Executes
When a worker begins processing a job, it must first secure exclusive rights to that job to prevent concurrent execution by other workers. This is achieved through a lease-based ownership model.
5.1: Pull the job: The worker retrieves the next available job — either by pulling from a queue or receiving it through a push-based delivery mechanism.
5.2: Lease acquisition: The worker calls into the job run metadata to claim a lease on the specific run_id. Before running the job, the worker must claim it in the job run metadata store. This step prevents two workers from executing the same job simultaneously.
5.3 Atomic state flip: The lease acquisition and job status update happen in the same transaction:
- Set
claimed_by,claimed_at, andlease_expires_at - Change
job_run.statusfromqueued→running
This atomic operation ensures:
- No two workers can claim the same job simultaneously
- The job's lifecycle state is always consistent with its lease ownership
5.4 Execute the job: Once the lease is secured, the worker begins execution.
Execution must be idempotent wherever possible — meaning that if the job is retried (due to crash, timeout, or failover), repeating it won't cause unintended side effects (e.g., sending duplicate emails, charging twice). Details of idempotent design are covered in Deep Dive 4 (DD4)
5.5 Renew the lease: For long-running jobs, the worker must periodically extend the lease (heartbeat) before it expires.
If the worker fails to renew:
- The lease expires
- Another worker may safely steal and re-run the job
The specifics of heartbeat intervals and renewal mechanisms are discussed in Deep Dive 6 (DD6).
What a lease means in a job scheduler
Definition: A lease is a time-bound, exclusive claim that a worker takes on a single job_run. While the lease is valid, no other worker should execute that same attempt. When the lease expires (or is explicitly released), the run can be taken by someone else. In queue terms, this is the same idea as a visibility timeout.
Why it exists (the two problems it solves):
- No double execution: prevents two workers from running the same attempt at once.
- Automatic recovery: if a worker crashes or stalls, the lease expires and a sweeper (or the queue) can safely requeue the run.
How it works:
- Claim: worker pulls a message, then atomically flips the run to
runningand sets fields likeclaimed_by,claimed_at,lease_expires_at = now() + ttl. - Hold/Renew: long jobs heartbeat—each heartbeat extends
lease_expires_atbefore it lapses. - Finish/Release: on success or a terminal failure, the worker writes the final status and clears/releases the lease.
- Expire/Rescue: if
lease_expires_at ≤ now()(no heartbeat), a sweeper requeues the run (status →queued) or spawns the next attempt with backoff.
Where it lives:
- In DB-backed schedulers: on the
job_runsrow (claimed_by,lease_expires_at,lease_token). - In message queues: as the broker's visibility timeout; un-acked messages reappear after the timeout.
Step 6: Execute the Job
With the lease secured, the worker performs the task. If the job is long-running, the worker sends heartbeats to extend the lease before expiration.
Step 7: Complete or Retry
-
On success, the worker marks the run
succeeded. -
On retryable failure, the worker records
failed, creates a newjob_runwith an increased backoff (available_atin the future), and re-enqueues.This preserves the original intent while keeping each attempt auditable.
The exact same loop works for future/recurring jobs—only available_at changes (now vs. later). That's why starting with immediate cleanly sets the foundation.
DB schema design
When you design a scheduler component (whether for immediate, or maybe later to support delayed, or recurring execution), having two separate tables — jobs and job_runs — isn't just a nice-to-have, it's the backbone for traceability, retries, and scaling execution safely.
Why do you need both jobs table and job_runs table?
Separating jobs and job_runs brings critical advantages: - Durable history of all executions - Support for retries, backoffs, and audits - Clean separation of logic (definition vs. instance)
Queues are ephemeral — once a message is consumed, it's gone. But scheduling systems need traceability. By keeping execution logs in a relational DB, we gain observability, consistency, and strong recovery semantics.
SQL vs NoSQL for the Core Store
A relational database (e.g., PostgreSQL or MySQL) is a strong choice for the system of record. It offers:
- ACID + constraints: atomic write of
job+job_run, foreign keys for lineage, unique/partial indexes for idempotency (e.g.,UNIQUE(dedupe_key) WHERE dedupe_key IS NOT NULL). - Worker-safe concurrency:
SELECT … FOR UPDATE SKIP LOCKEDenables safe claiming with high parallelism; leases map cleanly to row updates. - Right indexes for hot paths: partial indexes on status subsets (e.g.,
status='queued') keep ready-pick scans tiny; BRIN/partitioning handle large, time-orderedjob_runs.
NoSQL systems (like Redis or DynamoDB) are great for coordination — leases, heartbeats, rate limits — but fall short for long-term tracking, history, and consistency guarantees. They lack:
- Multi-row transactions
- Foreign key enforcement
- Efficient time-based queries
By establishing immediate job execution first, we create a resilient, debuggable, and extensible foundation — one that scales naturally to support future and recurring jobs, and powers reliable infrastructure at scale.