Interview Setup
Interview Prompt
Design a distributed job scheduler that runs recurring (cron) and one time jobs across a cluster of workers. Support job priorities, dependencies (DAGs), at least once delivery with effectively once side effects for idempotent handlers, and deduplication of duplicate submissions.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| Exactly once, at least once, or at most once? | End to end exactly once behavior across arbitrary external side effects requires transactional coordination. At least once delivery with idempotent handlers and deduplication is the production default. |
| How many jobs per second and max concurrent executions? | Handling 10,000 jobs per second drives sharded scheduler metadata storage and worker pool sizing. |
| Do jobs have dependencies (DAG) or are they independent? | DAGs require topological sort, dependency tracking, and failure propagation. |
| What's the max job duration and retry policy? | Long running jobs require periodic lease renewal. Short jobs can use faster retries with exponential backoff. |
Scope
In scope
- Cron scheduling with timezone support
- Effectively once side effects through submission deduplication and execution leases
- Priority based job queuing
- Distributed leases for single active execution
- Job DAGs with dependency resolution
- Worker assignment and lease management
Out of scope (state explicitly)
- Full workflow engine with sagas and compensation, which is compared with Temporal instead
- Building a container orchestrator because Kubernetes handles that layer
- Real time streaming job processing
Functional Requirements
Clarify whether the system must support one time execution, recurring cron schedules, retry policies, and priority tiers. In senior interviews, probe early into DAG dependencies and effectively once execution guarantees.
- Job Submission: Submit one time jobs for a specific future timestamp or recurring jobs on a cron schedule.
- One Time Execution: Execute a task once at a precise timestamp, such as delivering a reminder email at a specified time.
- Recurring Schedules: Execute jobs repeatedly on fixed intervals or complex cron expressions with timezone support.
- Priority Tiers: Support configurable priority tiers across critical, high, normal, and low workloads.
- DAG Dependencies: Enforce directed acyclic graph dependencies so child jobs execute only after their parent jobs complete.
- Configurable Retries: Support automatic retry policies with exponential backoff, jitter, and maximum retry limits.
- Lifecycle Tracking: Track execution status across pending, scheduled, dispatched, running, completed, failed, and cancelled states.
- Job Cancellation: Allow clients to cancel pending or scheduled one time jobs before dispatch and allow clients to cancel recurring schedule definitions.
- Effectively Once Processing: Combine at least once dispatch with distributed execution leases and idempotent workers so duplicate deliveries produce the same side effects.
Non-Functional Requirements
Interviewers focus primarily on zero job loss and high scheduling precision. Emphasize high availability for the scheduler and worker idempotency early to demonstrate production readiness.
- Zero Job Loss: Scheduled jobs must persist durably and remain recoverable after worker crashes, scheduler failover, or temporary broker outages.
- High Scalability: Support millions of registered job definitions with tens of thousands of concurrent executions.
- Low Scheduling Latency: Dispatch jobs within one second of their scheduled fire time under normal operating load.
- High Availability: Achieve 99.99% availability because scheduler outages block time sensitive business workflows.
- Partition Fault Tolerance: Maintain correctness through network partitions, node crashes, and leader elections.
- Idempotent Side Effects: Guarantee that duplicate task deliveries or retries do not result in duplicate state mutations.
- Deterministic Ordering: Respect configured priority tiers and DAG dependency constraints while enforcing fairness and starvation protection.
Capacity Estimations
| Metric | Value |
|---|---|
| Total scheduled jobs | 1M (active definitions + instances) |
| Average job executions / minute | 100K |
| Average / peak job executions / sec | ~1,700 / 10K |
| Avg job duration | 30 seconds |
| Concurrent running jobs | ~50K |
| Job metadata size | 1 KB |
| Storage | 1M x 1 KB ≈ 1 GB (hot), with history preserved in PostgreSQL cold store |
Scale Insight: Only jobs due within the next 24 hours reside in Redis timer buckets. Future jobs remain in the PostgreSQL durable store, while a rolling loader transfers jobs entering the upcoming 24 hour window into active Redis sorted sets.
Architecture Diagram
In the interview, compare Redis sorted set timer buckets with hierarchical timing wheels before selecting the primary storage engine.
The architecture separates durable persistence, time based triggering, and worker execution into distinct tiers. Jobs enter through the API Gateway, are validated by the Job Store, and persist in PostgreSQL. Jobs due within the next 24 hours are indexed in Redis sorted set timer buckets. A leader elected Schedule Manager scans due buckets every second, claims jobs atomically, records the dispatch state durably, and publishes work to priority based Kafka topics. A reconciliation loop republishes any job left in dispatched state without a confirmed Kafka publication. Workers pull tasks, acquire execution leases with fencing tokens, enforce idempotency checks, and record completion.
Component Deep Dives
Each architectural component addresses a specialized responsibility across persistent storage, time indexed scheduling, and distributed worker execution.
1. Job Store (CRUD API)
The Job Store serves as the durable source of truth where every submission lands before scheduling. It validates the job payload, enforces submission idempotency using a client key, persists metadata to PostgreSQL, computes the target execution window, and inserts the job ID into the Redis sorted set timer bucket when the occurrence falls within the upcoming 24 hour window.
2. Time Bucket Strategy (Redis ZSETs) ⭐
Partitioning scheduled jobs into granular time buckets prevents schedulers from scanning millions of future records on every tick. Each bucket is implemented as a Redis sorted set with the key timer:bucket:{minute_timestamp}:{shard_id}. In the baseline deployment, shard_id is 0. At scale, the shard suffix distributes timer buckets across Redis nodes. Job IDs are scored by their exact epoch execution second, so the scheduler can query only the active bucket instead of performing expensive full table scans.
3. Schedule Manager (Scanner Loop) ⭐
The Schedule Manager executes on a leader node elected through etcd or ZooKeeper. Every second, the scanner loop queries the current and previous minute buckets with ZRANGEBYSCORE to capture due jobs and boundary stragglers. It atomically claims each due job with ZREM, acquires a secondary dispatch lock, records the dispatched state and timestamp in PostgreSQL, and publishes the execution event to Kafka. A reconciliation loop republishes stale dispatched jobs when publication was not confirmed.
4. Job Dispatcher (Priority Based Routing)
Priority based dispatch isolates workloads into dedicated Kafka topics for critical, high, normal, and low execution tiers. This prevents heavy batch workloads from causing head of line blocking for interactive jobs and enables independent autoscaling of specialized worker pools.
5. Worker Pool & Executor Lifecycle
Workers subscribe to priority topics through dedicated consumer groups. Each worker follows a structured execution lifecycle to maintain safety during consumer rebalancing and process crashes.
- Consume the job payload from the designated priority Kafka topic.
- Acquire a bounded distributed execution lease with a fencing token to prevent concurrent runs if Kafka rebalances partitions.
- Perform a database idempotency check by reading the current status and skipping execution if the record is already marked as completed.
- Update the database status to running and record the execution start timestamp.
- Execute the actual payload by invoking the target HTTP endpoint, script, or containerized task.
- Upon successful completion, conditionally transition the job to completed, release the execution lease, trigger callbacks, and notify dependent DAG nodes. Only the successful state transition may decrement dependency counters.
- Upon execution failure, increment the retry counter and schedule a future retry using exponential backoff with jitter, or push the task to the dead letter queue if retries are exhausted.
6. Recurring Job Scheduler
The recurring job engine calculates future executions from the cron schedule rather than from the actual completion time of the previous run. It accounts for daylight saving time transitions in the configured timezone, applies the recurring job's misfire and concurrency policies, writes the next execution instance to PostgreSQL, and places it in a Redis timer bucket only when the execution window falls within the upcoming 24 hours. Cancelling a recurring definition marks it inactive and removes its future timer entries. Any already running instance follows the configured execution policy.
7. DAG Dependency Resolver ⭐
The dependency resolver coordinates multi stage workflows where child tasks wait for upstream parent jobs to finish. When a job with dependencies enters the system, edge records persist in the database and an atomic Redis counter initializes to the number of prerequisite parents. Each successful parent completion decrements the child counter. When the counter reaches zero, all prerequisites are satisfied and the child job can be dispatched immediately.
8. Dead Letter Queue (DLQ)
Jobs that exhaust their configured retry budget are routed to the jobs.dead_letter topic. A dedicated dead letter consumer logs the complete error stack trace, increments failure metrics, and triggers on call alerts for critical workloads. Administrative interfaces allow engineers to inspect failed payloads, modify faulty parameters, and re-enqueue tasks into active processing queues.
9. Callback Service
The callback service delivers execution outcomes to upstream caller systems asynchronously via HTTP webhooks. It uses a stable callback delivery ID derived from the completed execution, performs three retries with exponential backoff, and sends that ID with every attempt so receivers can deduplicate repeated webhook deliveries. If callback attempts fail completely, the underlying job remains safely marked as completed because webhook delivery is decoupled as a best effort operational notification.
10. Effectively Once Execution
True end to end exactly once behavior across arbitrary external side effects requires transactional coordination. The scheduler instead targets effectively once side effects by combining at least once message dispatch with three defense layers: atomic ZREM claims on the Redis timer bucket, distributed dispatch locks using monotonic fencing tokens, and worker execution locks coupled with database idempotency checks that skip processing when records are already marked completed.
Event Bus Design (Kafka)
The event bus decouples the schedule manager from worker execution pools and absorbs sudden dispatch bursts. Topics are partitioned by priority to maintain strict workload isolation.
# Kafka Topic Architecture and Partition Topology
topics:
jobs.critical:
priority: 1
partitions: 16
consumers: "dedicated critical worker pool"
use_cases: ["billing", "fraud alerts", "payment reconciliation"]
jobs.high:
priority: 2
partitions: 32
consumers: "shared high priority workers"
jobs.normal:
priority: 3
partitions: 64
partition_key: "job_id"
consumers: "default worker pool"
producers: "Schedule Manager after atomic ZREM claim and dispatch record"
jobs.low:
priority: 4
partitions: 32
consumers: "idle only batch workers"
jobs.dead_letter:
partitions: 8
producers: "Worker Pool after max retries are exhausted"
consumers: ["DLQ inspector", "on-call alerts", "manual retry console"]
configuration:
replication_factor: 3
min_insync_replicas: 2
producer_acks: "all"
producer_idempotence: true
dispatch_path: "Redis ZREM claim -> PG status=dispatched with dispatched_at -> Kafka publish, while stale dispatched jobs are reconciled and republished idempotently"
execution_path: "worker consume -> exec lease -> run payload -> PG completed/failed"API Design
Client API Signatures
The scheduler exposes strongly typed interfaces for job submission, recurring schedule registration, lifecycle tracking, and cancellation.
// Core Domain Types
export type JobId = string;
export type RecurringJobId = string;
export type IsoTimestamp = string;
export type JobPriority = "critical" | "high" | "normal" | "low";
export type JobStatus =
| "pending"
| "scheduled"
| "dispatched"
| "running"
| "completed"
| "failed"
| "cancelled";
export interface RetryPolicy {
maxRetries: number;
backoff: "exponential" | "fixed" | "linear";
initialDelaySec: number;
maxDelaySec?: number;
jitterRatio?: number;
}
export type MisfirePolicy = "fire_all" | "fire_once" | "skip";
export type ConcurrencyPolicy = "allow" | "forbid" | "replace";
export interface SubmitOneTimeJobRequest {
jobType: string;
executeAt: IsoTimestamp; // ISO 8601 UTC timestamp
priority?: JobPriority;
payload: Record<string, unknown>;
retryPolicy?: RetryPolicy;
callbackUrl?: string;
idempotencyKey?: string;
}
export interface SubmitRecurringJobRequest {
jobType: string;
cronExpression: string; // Standard five field cron format
timezone: string; // IANA timezone, e.g., "America/New_York"
priority?: JobPriority;
payload: Record<string, unknown>;
retryPolicy?: RetryPolicy;
callbackUrl?: string;
idempotencyKey?: string;
misfirePolicy?: MisfirePolicy;
concurrencyPolicy?: ConcurrencyPolicy;
}
export interface SubmitOneTimeJobResponse {
jobId: JobId;
status: JobStatus;
executeAt: IsoTimestamp;
}
export interface SubmitRecurringJobResponse {
recurringId: RecurringJobId;
nextExecuteAt: IsoTimestamp;
}
export interface JobExecutionRecord {
jobId: JobId;
status: JobStatus;
executeAt: IsoTimestamp;
startedAt?: IsoTimestamp;
completedAt?: IsoTimestamp;
retries: number;
callbackDeliveryId?: string;
result?: Record<string, unknown>;
errorMessage?: string;
}
export interface CancelResponse {
status: "cancelled";
}
// Client Interface for Distributed Job Scheduling
// Public job APIs use canonical UTC timestamps and idempotency keys for safe retries.
export interface JobSchedulerClient {
submitOneTimeJob(req: SubmitOneTimeJobRequest): Promise<SubmitOneTimeJobResponse>;
submitRecurringJob(req: SubmitRecurringJobRequest): Promise<SubmitRecurringJobResponse>;
getJobStatus(jobId: JobId): Promise<JobExecutionRecord>;
cancelJob(jobId: JobId): Promise<CancelResponse>;
cancelRecurringJob(recurringId: RecurringJobId): Promise<CancelResponse>;
}Submit One-Time Job
POST /api/v1/jobs
Content-Type: application/json
{
"job_type": "send_email",
"execute_at": "2026-03-14T15:00:00Z",
"priority": "normal",
"payload": {
"to": "user@example.com",
"template": "welcome_email",
"vars": {"name": "Alice"}
},
"retry_policy": {
"max_retries": 3,
"backoff": "exponential",
"initial_delay_sec": 60
},
"callback_url": "https://myservice.com/job-callback"
}
Response: 201 Created
{
"job_id": "job-uuid-12345",
"status": "scheduled",
"execute_at": "2026-03-14T15:00:00Z"
}Submit Recurring Job
POST /api/v1/jobs/recurring
Content-Type: application/json
{
"job_type": "generate_report",
"cron_expression": "0 0 * * *",
"timezone": "America/New_York",
"payload": {"report_type": "daily_sales"},
"priority": "high",
"idempotency_key": "daily_sales_schedule_v1",
"misfire_policy": "fire_once",
"concurrency_policy": "forbid"
}
Response: 201 Created
{
"recurring_id": "rec-uuid-67890",
"next_execute_at": "2026-03-15T00:00:00-04:00"
}Cancel Recurring Job
DELETE /api/v1/jobs/recurring/rec-uuid-67890
Response: 200 OK
{
"status": "cancelled"
}Get Job Status
GET /api/v1/jobs/job-uuid-12345
Response: 200 OK
{
"job_id": "job-uuid-12345",
"status": "completed",
"execute_at": "2026-03-14T15:00:00Z",
"started_at": "2026-03-14T15:00:01Z",
"completed_at": "2026-03-14T15:00:05Z",
"retries": 0,
"callback_delivery_id": "cb-uuid-54321",
"result": {"email_sent": true}
}Cancel Job
DELETE /api/v1/jobs/job-uuid-12345
Response: 200 OK
{
"status": "cancelled"
}Common Error Responses
400 Bad Request: invalid input, missing required fields, or malformed JSON payload
401 Unauthorized: missing or invalid authentication token or API key
403 Forbidden: authenticated caller lacks required permissions for this resource
404 Not Found: requested resource ID does not exist
409 Conflict: duplicate write or version conflict, retry with a unique idempotency key
422 Unprocessable Entity: syntactically valid request failed semantic business validation
429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header
500 Internal Error: unexpected server failure, retry safely with an idempotency key
503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff
202 Accepted: asynchronous job queued successfully, poll GET /jobs/{id} for completion status
408 Request Timeout: background job is still executing, continue polling status endpointData Model
PostgreSQL Schema
Relational tables in PostgreSQL provide durable ACID guarantees for job definitions and DAG dependencies, while Redis key value schemas power low latency timers and distributed leases.
CREATE TABLE jobs (
job_id UUID PRIMARY KEY,
job_type VARCHAR(64) NOT NULL,
status VARCHAR(20) NOT NULL, -- pending, scheduled, dispatched, running, completed, failed, cancelled
priority SMALLINT DEFAULT 3, -- 1=critical, 2=high, 3=normal, 4=low
payload JSONB,
execute_at TIMESTAMPTZ NOT NULL,
started_at TIMESTAMPTZ,
dispatched_at TIMESTAMPTZ,
completed_at TIMESTAMPTZ,
retry_count INT DEFAULT 0,
max_retries INT DEFAULT 3,
backoff_strategy VARCHAR(20),
next_retry_at TIMESTAMPTZ,
idempotency_key VARCHAR(128),
dedup_key VARCHAR(128),
result JSONB,
error_message TEXT,
callback_url TEXT,
callback_delivery_id UUID,
created_by VARCHAR(128),
created_at TIMESTAMPTZ,
updated_at TIMESTAMPTZ
);
CREATE UNIQUE INDEX uq_jobs_idempotency
ON jobs (idempotency_key)
WHERE idempotency_key IS NOT NULL;
CREATE UNIQUE INDEX uq_jobs_dedup_key
ON jobs (dedup_key)
WHERE dedup_key IS NOT NULL;
CREATE INDEX idx_execute
ON jobs (status, execute_at);
CREATE INDEX idx_type
ON jobs (job_type, status);
CREATE TABLE recurring_jobs (
recurring_id UUID PRIMARY KEY,
job_type VARCHAR(64) NOT NULL,
cron_expression VARCHAR(64) NOT NULL,
timezone VARCHAR(64) NOT NULL,
payload JSONB,
priority SMALLINT NOT NULL DEFAULT 3,
idempotency_key VARCHAR(128),
misfire_policy VARCHAR(20) NOT NULL DEFAULT 'fire_once',
concurrency_policy VARCHAR(20) NOT NULL DEFAULT 'allow',
is_active BOOLEAN DEFAULT TRUE,
last_executed TIMESTAMPTZ,
next_execute_at TIMESTAMPTZ,
created_at TIMESTAMPTZ NOT NULL
);
CREATE UNIQUE INDEX uq_recurring_idempotency
ON recurring_jobs (idempotency_key)
WHERE idempotency_key IS NOT NULL;
CREATE INDEX idx_recurring_due
ON recurring_jobs (is_active, next_execute_at);
CREATE TABLE job_deps (
parent_job_id UUID NOT NULL,
child_job_id UUID NOT NULL,
PRIMARY KEY (parent_job_id, child_job_id)
);
CREATE INDEX idx_child
ON job_deps (child_job_id);Redis Key Schemas
# Time bucket ZSET for due jobs in a given minute
timer:bucket:{minute_epoch}:{shard_id}:
type: ZSET
member: job_id
score: exact_timestamp_seconds
purpose: "Partitions jobs into minute-based buckets, where shard_id is 0 in the baseline deployment and distributes buckets across Redis nodes at scale"
# Distributed lock to guarantee single dispatch by the scheduler
job:dispatch:{job_id}:
type: String
value: "1"
ttl_seconds: 60
purpose: "Belt-and-suspenders lock preventing duplicate worker queue enqueue"
# Distributed execution lease held by worker during job processing
job:exec:{job_id}:
type: String
value: "{worker_id}:{fencing_token}"
ttl_seconds: "min(job_timeout * 2, 300) with periodic renewal, ensuring the lease remains bounded independently of the full job duration"
purpose: "Prevents concurrent execution after Kafka rebalances, while fencing tokens reject stale workers"
# DAG dependency counter tracking unmet upstream parent completions
deps:{child_job_id}:
type: String (integer counter)
value: remaining_parent_count
purpose: "Atomic DECR on parent success, dispatching the child job immediately when the counter reaches 0"Kafka Topics Catalog
jobs.critical: Critical priority queue reserved for billing, payments, and fraud alerts.jobs.high: High priority queue for user facing interactive asynchronous tasks.jobs.normal: Default task pipeline for standard background processing.jobs.low: Low priority queue reserved for bulk analytics and batch processing.jobs.dead_letter: Quarantined tasks that exhausted all configured retry attempts.
Fault Tolerance
Scheduler leader failure, duplicate execution, and stuck jobs each need explicit handling.
| Concern | Solution |
|---|---|
| Scheduler leader crash | Standby scheduler promoted via etcd leader election in under 5 seconds. |
| Worker crash mid execution | The execution lease expires, allowing the scheduler to re-dispatch the job. Fencing tokens and idempotent handlers protect downstream side effects. |
| Scheduler down during job times | Upon recovery, the scanner sweeps past due scheduled jobs from PostgreSQL, rebuilds Redis timer entries, and dispatches them according to the configured misfire policy. |
| Duplicate executions | Atomic ZREM claims, dispatch locks, execution locks, and idempotent application handlers. |
| Kafka brokers offline | Keep the durable dispatched record in PostgreSQL and let the reconciliation loop retry publication with exponential backoff. Scheduler local memory is not the source of truth. |
| Clock skew on nodes | Synchronize nodes with strict NTP. One-minute time buckets tolerate minor millisecond clock variations. |
Retry Backoff Jitter Policy
Retries follow an exponential backoff formula: delay = initial_delay * 2^attempt + random_jitter. Adding random jitter decorrelates retry attempts and prevents thundering herds from overwhelming downstream services.
Additional Considerations
Comparison with Existing Systems
Job dependencies, priority queues, execution leases, and dead letter handling are senior topics.
| System | Type | Primary Use Case | Trade-off |
|---|---|---|---|
| Cron | Single Node | Simple server scripts | Single point of failure, no scaling. |
| Celery | Task Queue | Asynchronous task execution | No robust cron or recurring scheduler built in. |
| Airflow | Workflow Engine | Batch ETL and data workflows | Higher latency and heavier operational model than a lightweight scheduler. |
| Quartz | Distributed Java Lib | JVM-based application schedules | Coupled to the Java ecosystem. |
| Temporal | Stateful Workflows | Complex long lived workflows and sagas | Heavier operational model than a focused scheduler. |
System Metrics to Monitor
- Lag Time: Difference between scheduled and actual start execution time. Target < 1s.
- Success Rate: Completed vs failed jobs. Target > 99.9%.
- DLQ Depth: Backlog of permanently failed tasks (alert if growing).
- Resource Utilization: Redis CPU/RAM, worker CPU utilization.
Related Problems and Concepts
Explore related architectures and distributed systems foundations that connect directly to this design:
- Distributed Message Broker: High throughput partitioned commit logs that decouple schedule dispatch from worker pools.
- Distributed Stream Processing: Real time event streaming and stateful stream joins for continuous data pipelines.
- Message Queues Fundamentals: Queue patterns, backpressure, dead letter routing, and consumer group mechanics.
- Consistent Hashing: Partitioning Redis time buckets and worker clusters across horizontal nodes.
- Replication, Failover, and Leader Election: Consensus protocols and etcd leader leases ensuring high availability for the schedule manager.
- System Design Interview Patterns: Standard architectural blueprints and communication frameworks for staff level interviews.
Interview Walkthrough
- 25-minute pacing guide
Focus on the core scheduling loop first, then expand to priority queues and fault tolerance if interviewing at staff level.
- Requirements and effectively once side effect guarantees (3 min)
- Redis sorted set timer buckets vs hierarchical timing wheels (7 min)
- Leader elected Schedule Manager for cron dispatch (6 min)
- Idempotency keys and distributed execution locks (5 min)
- Worker pool execution, backoff retries, and dead letter queues (4 min)
- Compare this design with Temporal for durable long lived workflows and Airflow for batch ETL pipelines. Use this lightweight scheduler for high throughput, low latency cron tasks, while moving durable business sagas to Temporal and large data workflows to Airflow.
- Recommend Redis sorted set timer buckets over hierarchical timing wheels when durable persistence and horizontal sharding are required. Timing wheels excel for transient in memory timers with constant time expiration, whereas Redis sorted sets provide crash recovery and cluster partitionability.
- Frame the core problem around decoupling schedule timing from worker execution. The schedule manager decides when a job becomes due, whereas decoupled worker pools determine who executes the task.
- Explain how a leader elected schedule manager dispatches due jobs into a Distributed Message Broker to balance load across worker consumer groups.
- Highlight job deduplication using client idempotency keys and lease based claiming so that worker crashes never result in dropped tasks or unhandled duplicate mutations.
- Detail retry policies featuring exponential backoff with randomized jitter, dead letter queues, and maximum attempt thresholds to isolate poison pill jobs.
- Address DAG dependencies by explaining how atomic Redis counters and directed workflow graphs orchestrate child job triggering without flattening graph structures into simple queues.
- Address the common pitfall of database polling. Avoid having workers poll relational database tables every second for due jobs because index contention degrades throughput. Use time bucketed Redis sets or native delay queues for the hot scheduling path, while PostgreSQL remains the durable recovery source.
Engineering Trade-offs
1. End-to-End Execution Flow
Your interviewer will push on the worker dispatch model and queue backend. Walk through push vs pull, then explain which choice you would use for scheduling a million jobs.
Consider a recurring job scheduled for execution at 15:00:00 UTC:
- T minus 60 minutes (Submission): The client submits the job specification. PostgreSQL stores the canonical record. Because the execution time falls within the 24 hour active window, the job ID is added to the Redis sorted set timer bucket for 15:00:00.
- T minus 0 seconds (Dispatch): The leader scheduler scans the active bucket, atomically claims the job using
ZREM, updates PostgreSQL status todispatched, and publishes an execution event to the designated Kafka topic. - T plus 100 milliseconds (Worker Execution): An available worker in the consumer group pulls the message from Kafka, obtains a distributed execution lease with a fencing token, and begins processing the job payload.
- T plus 5 seconds (Completion): The worker completes the task, updates the status to
completedin PostgreSQL, releases the execution lease, triggers downstream callbacks, and evaluates dependent DAG nodes.
2. Race Conditions in Distributed Scheduling
When multiple scheduler instances experience a split brain network partition and scan the same timer bucket, atomic Redis ZREM operations ensure that only one instance can claim the job ID. The other instance receives zero and skips the claim. Fencing and idempotent execution still protect downstream side effects if a stale worker survives beyond its lease.
Boundary race conditions can occur when a job scheduled at 15:00:59.9 is evaluated by a scanner running at 15:00:59.5. If the scanner never revisits the previous bucket, that record could remain stranded. Inspecting both the current minute and the previous minute buckets reduces that boundary window, while the durable PostgreSQL recovery sweep provides the final safety net.
3. Scale Out Timer Buckets
When managing tens of millions of jobs, a single Redis instance becomes bottlenecked by memory consumption and single-threaded CPU execution. We resolve this by sharding timer buckets across Redis nodes using Consistent Hashing. Keys follow the partitioned schema timer:bucket:{minute}:{hash(job_id) % 16}, enabling multiple scheduler shards to concurrently scan distinct bucket partitions without contention.
4. Timer Approaches Comparison
Distributed schedulers typically evaluate three primary mechanisms for tracking upcoming execution timers:
- Redis Sorted Sets: Provide an optimal balance of low latency reads, memory efficiency, and snapshot persistence, while naturally supporting horizontal sharding across cluster nodes.
- Hierarchical Timing Wheels: Deliver constant time insertion and expiration operations, but introduce substantial complexity when persisting state across process restarts and distributing timers across a cluster.
- Relational Database Polling (FOR UPDATE SKIP LOCKED): Offers total ACID durability, but frequent high frequency polling can place heavy read contention on database indexes. It functions best as an initial baseline or emergency fallback.
5. Delivery Guarantees
The system implements at least once delivery coupled with idempotent application handlers. If an executor crashes mid-task, its visibility lease expires and the scheduler re-dispatches the work. Application handlers must verify deduplication keys or rely on atomic database upserts to guarantee that duplicate deliveries produce identical side effects.
Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.