NEW: ML Mock & Coaching now available

Questions

Job Scheduler

MetaGoogleOpenAI

Design Job Scheduler - This system design covers distributed scheduling, fault tolerance, delivery semantics, and high-throughput job execution at scale.

50 min read

Staff+ Engineer: 12+ years of experience in distributed scheduling and workflow orchestration. Tech lead for job execution platforms processing millions of tasks daily. Engineering Manager: Over 12 years of industry experience with 6 years in engineering leadership, building and scaling tech infrastructure teams that deliver end-to-end large-scale distributed systems.

Challenge

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:

  1. Support immediate job execution (run now).
  2. Support future and recurring schedules (cron, delayed).
  3. 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)
NameDescription
1. Immediate ExecutionSupport running jobs immediately (run now).
2. Future & Recurring ExecutionSupport jobs at specific future time or recurring cadence (cron).
3. ObservabilityProvide visibility into queue depth, success/error rates, latency, and worker status.
Non-Functional Requirements (NFRs)
NameDescription
1. High ScalabilitySupport up to 10,000 job executions per second.
2. Low LatencyExecute jobs within 2 seconds of scheduled time.
3. CorrectnessAt-least-once execution with configurable delivery semantics.
4. High AvailabilityFavor availability over strict consistency.
5. Fault ToleranceTolerate 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_at or scheduled_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
FieldTypeKeyDescription
job_idUUIDPKLogical job identifier (the 'intent').
typeTEXTNOT NULLJob kind, e.g., email.send.
payloadJSONBNOT NULLParameters for the job.
schedule_typeTEXTNOT NULLimmediate / at / cron.
next_fire_atTIMESTAMPTZNULLABLENext due time (used for at/cron).
priorityINTNOT NULL DEFAULT 0Scheduling priority.
statusTEXTNOT NULL DEFAULT activeLifecycle control.
created_atTIMESTAMPTZNOT NULL DEFAULT now()Auditing/slicing.
Job_runs Table
FieldTypeKeyDescription
run_idUUIDPKPhysical execution attempt ID.
job_idUUIDFK → jobs(job_id)Which job this run belongs to.
run_numberINTNOT NULL1,2,3… (attempt/order).
statusENUMNOT NULLqueued/running/succeeded/failed/canceled.
queued_atTIMESTAMPTZNOT NULL DEFAULT now()When the run was enqueued.
available_atTIMESTAMPTZNOT NULL DEFAULT now()When it becomes eligible to run (backoff).
claimed_byTEXTNULLABLEWorker ID that leased it.
claimed_atTIMESTAMPTZNULLABLELease start.
lease_expires_atTIMESTAMPTZNULLABLEWhen others may safely steal it.
started_atTIMESTAMPTZNULLABLEExecution start.
finished_atTIMESTAMPTZNULLABLEExecution end.
error_codeTEXTNULLABLECategorical failure reason.
error_messageTEXTNULLABLEDebuggable failure info.

High Level Design

FR1: Scheduler component to support immediate job execution

job-scheduler-FR1

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 jobs table to describe what to run
  • Inserts a new row into the job_runs table to represent this specific execution attempt, with status = queued and available_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, and lease_expires_at
  • Change job_run.status from queued → running

This atomic operation ensures:

  1. No two workers can claim the same job simultaneously
  2. 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):

  1. No double execution: prevents two workers from running the same attempt at once.
  2. 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:

  1. Claim: worker pulls a message, then atomically flips the run to running and sets fields like claimed_by, claimed_at, lease_expires_at = now() + ttl.
  2. Hold/Renew: long jobs heartbeat—each heartbeat extends lease_expires_at before it lapses.
  3. Finish/Release: on success or a terminal failure, the worker writes the final status and clears/releases the lease.
  4. 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_runs row (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 new job_run with an increased backoff (available_at in the future), and re-enqueues.

    This preserves the original intent while keeping each attempt auditable.

Tip

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.

Tip

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 LOCKED enables 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-ordered job_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.

Sign in to continue reading

"Job Scheduler" requires a free account to access.

Sign in to continue