Designing a Task Scheduler
Difficulty: Intermediate Patterns: Strategy, Observer, Thread Pool Asked at: Google, Amazon, Microsoft, Uber
Functional Requirements
- Schedule tasks at specific times - submit a task to run at a given future timestamp
- Recurring tasks - support tasks that repeat at fixed intervals (every N seconds/minutes)
- Priority queue - higher priority tasks execute first when multiple tasks are due
- Cancel tasks - cancel a scheduled or recurring task by its ID before execution
- Thread pool execution - execute tasks concurrently using a configurable thread pool
- Delayed execution - schedule a task to run after a specified delay from now
Non-Functional Requirements
- Thread-safety - concurrent scheduling/cancelling must not corrupt internal state
- Efficient polling - use min-heap so next task is always O(1) peek, O(log n) extract
- Graceful shutdown - complete in-progress tasks, discard pending ones on shutdown
Core Entities
| Entity | Description |
|---|---|
Task |
Encapsulates a unit of work with an ID, priority, and execution logic |
TaskType |
Enum: ONE_TIME, RECURRING, DELAYED |
Priority |
Enum: LOW, MEDIUM, HIGH, CRITICAL |
ScheduledTask |
Wraps a Task with scheduling metadata (next run time, interval) |
TaskScheduler |
Central orchestrator - schedules, cancels, dispatches tasks |
ExecutionStrategy |
Interface for how tasks are executed (thread pool, single thread) |
ThreadPoolExecutor |
Executes tasks using a fixed thread pool |
TaskObserver |
Interface notified on task lifecycle events (started, completed, failed) |
Class Diagram
classDiagram
class Priority {
<<enumeration>>
LOW
MEDIUM
HIGH
CRITICAL
}
class TaskType {
<<enumeration>>
ONE_TIME
RECURRING
DELAYED
}
class Task {
-String id
-String name
-Priority priority
-Runnable action
+execute()
+getId() String
+getPriority() Priority
}
class ScheduledTask {
-Task task
-long nextRunTimeMs
-long intervalMs
-TaskType type
-boolean cancelled
+cancel()
+isCancelled() boolean
+getNextRunTime() long
+reschedule()
}
class ExecutionStrategy {
<<interface>>
+execute(Task task)
+shutdown()
}
class ThreadPoolStrategy {
-ExecutorService pool
+execute(Task task)
+shutdown()
}
class TaskObserver {
<<interface>>
+onTaskStarted(Task task)
+onTaskCompleted(Task task)
+onTaskFailed(Task task, Exception e)
+onTaskCancelled(Task task)
}
class TaskScheduler {
-PriorityQueue~ScheduledTask~ taskQueue
-Map~String, ScheduledTask~ taskRegistry
-ExecutionStrategy executor
-List~TaskObserver~ observers
-boolean running
+schedule(Task, long delayMs) String
+scheduleAt(Task, long timestampMs) String
+scheduleRecurring(Task, long intervalMs) String
+cancel(String taskId) boolean
+start()
+shutdown()
}
TaskScheduler --> ExecutionStrategy
TaskScheduler --> ScheduledTask
TaskScheduler --> TaskObserver
ScheduledTask --> Task
Task --> Priority
ScheduledTask --> TaskType
ExecutionStrategy <|.. ThreadPoolStrategy
Design Patterns
| Pattern | Where | Why |
|---|---|---|
| Strategy | ExecutionStrategy interface with ThreadPoolStrategy |
Swap execution model (single-thread, thread pool, async) without changing scheduler logic |
| Observer | TaskObserver notified on task start/complete/fail/cancel |
Decouple logging, metrics, alerting from core scheduling |
| Thread Pool | ThreadPoolStrategy wraps a fixed-size executor |
Bound concurrency, reuse threads, prevent resource exhaustion |
| Command | Task encapsulates action as an object |
Tasks are first-class objects that can be queued, cancelled, rescheduled |
Data Structures
| Component | Structure | Why |
|---|---|---|
| Task queue | PriorityQueue<ScheduledTask> (min-heap by next run time + priority) |
O(1) peek next task, O(log n) insert/extract |
| Task registry | ConcurrentHashMap<String, ScheduledTask> |
O(1) lookup for cancel by ID |
| Observers | CopyOnWriteArrayList<TaskObserver> |
Thread-safe iteration during notification |
| Thread pool | FixedThreadPool(n) |
Bounded concurrency for task execution |
How It All Fits Together
Hereโs what happens when a recurring task is scheduled:
- Client calls
scheduleRecurring(task, intervalMs) - Scheduler wraps the task in a
ScheduledTaskwith type=RECURRING and computesnextRunTime = now + intervalMs - ScheduledTask is added to the min-heap (ordered by nextRunTime, then priority)
- ScheduledTask is registered in the taskRegistry map by ID
- Schedulerโs dispatcher thread wakes up, peeks the heap
- If top taskโs
nextRunTime <= now, it extracts the task - Checks if cancelled โ if yes, discard and notify observers
- Passes the task to
ExecutionStrategy.execute() - ThreadPoolStrategy submits the task to the thread pool
- On completion, observers are notified; if RECURRING, task is rescheduled (nextRunTime += interval) and re-inserted into the heap
Complete Code
Task and Priority
The Task class is an immutable command object. Priority determines execution order when multiple tasks are due simultaneously.
import java.util.*;
import java.util.concurrent.*;
import java.util.concurrent.atomic.*;
import java.util.concurrent.locks.*;
// โโโ Enums โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
enum Priority {
LOW(0), MEDIUM(1), HIGH(2), CRITICAL(3);
private final int level;
Priority(int level) { this.level = level; }
public int getLevel() { return level; }
}
enum TaskType {
ONE_TIME, RECURRING, DELAYED
}
// โโโ Task โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class Task {
private final String id;
private final String name;
private final Priority priority;
private final Runnable action;
public Task(String name, Priority priority, Runnable action) {
this.id = UUID.randomUUID().toString().substring(0, 8);
this.name = name;
this.priority = priority;
this.action = action;
}
public void execute() {
action.run();
}
public String getId() { return id; }
public String getName() { return name; }
public Priority getPriority() { return priority; }
@Override
public String toString() {
return "Task[" + name + " | " + priority + " | id=" + id + "]";
}
}
// โโโ ScheduledTask โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class ScheduledTask implements Comparable<ScheduledTask> {
private final Task task;
private long nextRunTimeMs;
private final long intervalMs;
private final TaskType type;
private volatile boolean cancelled;
public ScheduledTask(Task task, long nextRunTimeMs, long intervalMs, TaskType type) {
this.task = task;
this.nextRunTimeMs = nextRunTimeMs;
this.intervalMs = intervalMs;
this.type = type;
this.cancelled = false;
}
public void cancel() { this.cancelled = true; }
public boolean isCancelled() { return cancelled; }
public Task getTask() { return task; }
public long getNextRunTimeMs() { return nextRunTimeMs; }
public TaskType getType() { return type; }
public long getIntervalMs() { return intervalMs; }
public void reschedule() {
if (type == TaskType.RECURRING) {
this.nextRunTimeMs = System.currentTimeMillis() + intervalMs;
}
}
@Override
public int compareTo(ScheduledTask other) {
int timeCompare = Long.compare(this.nextRunTimeMs, other.nextRunTimeMs);
if (timeCompare != 0) return timeCompare;
// Higher priority first (descending)
return Integer.compare(other.task.getPriority().getLevel(),
this.task.getPriority().getLevel());
}
}
// โโโ Observer Interface โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
interface TaskObserver {
void onTaskStarted(Task task);
void onTaskCompleted(Task task);
void onTaskFailed(Task task, Exception e);
void onTaskCancelled(Task task);
}
// โโโ Execution Strategy Interface โโโโโโโโโโโโโโโโโโโโโโโ
interface ExecutionStrategy {
void execute(Task task, List<TaskObserver> observers);
void shutdown();
boolean isShutdown();
}
// โโโ Thread Pool Strategy โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class ThreadPoolStrategy implements ExecutionStrategy {
private final ExecutorService pool;
private volatile boolean shutdown;
public ThreadPoolStrategy(int poolSize) {
this.pool = Executors.newFixedThreadPool(poolSize);
this.shutdown = false;
}
@Override
public void execute(Task task, List<TaskObserver> observers) {
if (shutdown) return;
pool.submit(() -> {
observers.forEach(o -> o.onTaskStarted(task));
try {
task.execute();
observers.forEach(o -> o.onTaskCompleted(task));
} catch (Exception e) {
observers.forEach(o -> o.onTaskFailed(task, e));
}
});
}
@Override
public void shutdown() {
shutdown = true;
pool.shutdown();
try {
pool.awaitTermination(5, TimeUnit.SECONDS);
} catch (InterruptedException e) {
pool.shutdownNow();
}
}
@Override
public boolean isShutdown() { return shutdown; }
}
// โโโ Logging Observer โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class LoggingObserver implements TaskObserver {
@Override
public void onTaskStarted(Task task) {
System.out.println(" [STARTED] " + task);
}
@Override
public void onTaskCompleted(Task task) {
System.out.println(" [COMPLETED] " + task);
}
@Override
public void onTaskFailed(Task task, Exception e) {
System.out.println(" [FAILED] " + task + " - " + e.getMessage());
}
@Override
public void onTaskCancelled(Task task) {
System.out.println(" [CANCELLED] " + task);
}
}
// โโโ Task Scheduler โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class TaskScheduler {
private final PriorityQueue<ScheduledTask> taskQueue;
private final ConcurrentHashMap<String, ScheduledTask> taskRegistry;
private final ExecutionStrategy executor;
private final List<TaskObserver> observers;
private final ReentrantLock lock;
private final Condition taskAvailable;
private volatile boolean running;
private Thread dispatcherThread;
public TaskScheduler(ExecutionStrategy executor) {
this.taskQueue = new PriorityQueue<>();
this.taskRegistry = new ConcurrentHashMap<>();
this.executor = executor;
this.observers = new CopyOnWriteArrayList<>();
this.lock = new ReentrantLock();
this.taskAvailable = lock.newCondition();
this.running = false;
}
public void addObserver(TaskObserver observer) {
observers.add(observer);
}
public String schedule(Task task, long delayMs) {
long runTime = System.currentTimeMillis() + delayMs;
ScheduledTask st = new ScheduledTask(task, runTime, 0, TaskType.DELAYED);
enqueue(st);
return task.getId();
}
public String scheduleAt(Task task, long timestampMs) {
ScheduledTask st = new ScheduledTask(task, timestampMs, 0, TaskType.ONE_TIME);
enqueue(st);
return task.getId();
}
public String scheduleRecurring(Task task, long intervalMs) {
long runTime = System.currentTimeMillis() + intervalMs;
ScheduledTask st = new ScheduledTask(task, runTime, intervalMs, TaskType.RECURRING);
enqueue(st);
return task.getId();
}
public boolean cancel(String taskId) {
ScheduledTask st = taskRegistry.get(taskId);
if (st == null || st.isCancelled()) return false;
st.cancel();
observers.forEach(o -> o.onTaskCancelled(st.getTask()));
return true;
}
private void enqueue(ScheduledTask st) {
lock.lock();
try {
taskQueue.offer(st);
taskRegistry.put(st.getTask().getId(), st);
taskAvailable.signal();
} finally {
lock.unlock();
}
}
public void start() {
running = true;
dispatcherThread = new Thread(this::dispatch, "scheduler-dispatcher");
dispatcherThread.setDaemon(true);
dispatcherThread.start();
}
private void dispatch() {
while (running) {
lock.lock();
try {
while (taskQueue.isEmpty() && running) {
taskAvailable.await(100, TimeUnit.MILLISECONDS);
}
if (!running) break;
ScheduledTask top = taskQueue.peek();
if (top == null) continue;
long now = System.currentTimeMillis();
if (top.getNextRunTimeMs() <= now) {
taskQueue.poll();
if (top.isCancelled()) {
taskRegistry.remove(top.getTask().getId());
continue;
}
executor.execute(top.getTask(), observers);
// Reschedule recurring tasks
if (top.getType() == TaskType.RECURRING && !top.isCancelled()) {
top.reschedule();
taskQueue.offer(top);
} else {
taskRegistry.remove(top.getTask().getId());
}
} else {
// Wait until next task is due
long waitMs = top.getNextRunTimeMs() - now;
taskAvailable.await(waitMs, TimeUnit.MILLISECONDS);
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
} finally {
lock.unlock();
}
}
}
public void shutdown() {
running = false;
lock.lock();
try {
taskAvailable.signalAll();
} finally {
lock.unlock();
}
executor.shutdown();
}
public int getPendingCount() {
lock.lock();
try {
return taskQueue.size();
} finally {
lock.unlock();
}
}
public int getActiveTaskCount() {
return taskRegistry.size();
}
}
// โโโ Main Demo โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
public class TaskSchedulerDemo {
public static void main(String[] args) throws InterruptedException {
System.out.println("โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ");
System.out.println(" TASK SCHEDULER - LLD DEMO ");
System.out.println("โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ\n");
// Build scheduler with thread pool of 4
ExecutionStrategy executor = new ThreadPoolStrategy(4);
TaskScheduler scheduler = new TaskScheduler(executor);
scheduler.addObserver(new LoggingObserver());
scheduler.start();
// โโโ Schedule One-Time Delayed Tasks โโโโโโโโโโโโโโโโ
System.out.println("--- Scheduling Delayed Tasks ---");
Task t1 = new Task("SendEmail", Priority.HIGH, () -> {
System.out.println(" >> Sending email...");
});
Task t2 = new Task("GenerateReport", Priority.LOW, () -> {
System.out.println(" >> Generating report...");
});
Task t3 = new Task("CriticalAlert", Priority.CRITICAL, () -> {
System.out.println(" >> CRITICAL ALERT FIRED!");
});
scheduler.schedule(t1, 500); // run after 500ms
scheduler.schedule(t2, 1000); // run after 1s
scheduler.schedule(t3, 500); // same delay but higher priority
System.out.println("Pending tasks: " + scheduler.getPendingCount());
Thread.sleep(2000); // wait for tasks to execute
// โโโ Schedule Recurring Task โโโโโโโโโโโโโโโโโโโโโโโโ
System.out.println("\n--- Scheduling Recurring Task ---");
Task heartbeat = new Task("Heartbeat", Priority.MEDIUM, () -> {
System.out.println(" >> heartbeat ping at " + System.currentTimeMillis());
});
String heartbeatId = scheduler.scheduleRecurring(heartbeat, 600);
Thread.sleep(2000); // let it run a few times
// โโโ Cancel Recurring Task โโโโโโโโโโโโโโโโโโโโโโโโโโ
System.out.println("\n--- Cancelling Heartbeat ---");
boolean cancelled = scheduler.cancel(heartbeatId);
System.out.println("Cancelled: " + cancelled);
Thread.sleep(1000); // confirm no more heartbeats
// โโโ Schedule Task That Fails โโโโโโโโโโโโโโโโโโโโโโโ
System.out.println("\n--- Task That Throws Exception ---");
Task failingTask = new Task("FailingJob", Priority.HIGH, () -> {
throw new RuntimeException("Something went wrong!");
});
scheduler.schedule(failingTask, 100);
Thread.sleep(500);
// โโโ Shutdown โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
System.out.println("\n--- Shutting Down Scheduler ---");
scheduler.shutdown();
System.out.println("\nโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ");
System.out.println(" DEMO COMPLETE ");
System.out.println("โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ");
}
}
import threading
import heapq
import time
import uuid
from enum import Enum
from typing import Callable, Optional
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
# โโโ Enums โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class Priority(Enum):
LOW = 0
MEDIUM = 1
HIGH = 2
CRITICAL = 3
class TaskType(Enum):
ONE_TIME = "ONE_TIME"
RECURRING = "RECURRING"
DELAYED = "DELAYED"
# โโโ Task โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class Task:
def __init__(self, name: str, priority: Priority, action: Callable):
self._id = uuid.uuid4().hex[:8]
self._name = name
self._priority = priority
self._action = action
def execute(self):
self._action()
@property
def id(self) -> str:
return self._id
@property
def name(self) -> str:
return self._name
@property
def priority(self) -> Priority:
return self._priority
def __str__(self) -> str:
return f"Task[{self._name} | {self._priority.name} | id={self._id}]"
# โโโ ScheduledTask โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class ScheduledTask:
def __init__(self, task: Task, next_run_time: float, interval: float, task_type: TaskType):
self._task = task
self._next_run_time = next_run_time
self._interval = interval
self._type = task_type
self._cancelled = False
def cancel(self):
self._cancelled = True
@property
def is_cancelled(self) -> bool:
return self._cancelled
@property
def task(self) -> Task:
return self._task
@property
def next_run_time(self) -> float:
return self._next_run_time
@property
def task_type(self) -> TaskType:
return self._type
@property
def interval(self) -> float:
return self._interval
def reschedule(self):
if self._type == TaskType.RECURRING:
self._next_run_time = time.time() + self._interval
def __lt__(self, other: "ScheduledTask") -> bool:
if self._next_run_time != other._next_run_time:
return self._next_run_time < other._next_run_time
# Higher priority first (higher enum value = higher priority)
return self._task.priority.value > other._task.priority.value
# โโโ Observer Interface โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class TaskObserver(ABC):
@abstractmethod
def on_task_started(self, task: Task): pass
@abstractmethod
def on_task_completed(self, task: Task): pass
@abstractmethod
def on_task_failed(self, task: Task, error: Exception): pass
@abstractmethod
def on_task_cancelled(self, task: Task): pass
# โโโ Execution Strategy Interface โโโโโโโโโโโโโโโโโโโโโโโ
class ExecutionStrategy(ABC):
@abstractmethod
def execute(self, task: Task, observers: list[TaskObserver]): pass
@abstractmethod
def shutdown(self): pass
# โโโ Thread Pool Strategy โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class ThreadPoolStrategy(ExecutionStrategy):
def __init__(self, pool_size: int):
from concurrent.futures import ThreadPoolExecutor
self._pool = ThreadPoolExecutor(max_workers=pool_size)
self._shutdown = False
def execute(self, task: Task, observers: list[TaskObserver]):
if self._shutdown:
return
def run():
for obs in observers:
obs.on_task_started(task)
try:
task.execute()
for obs in observers:
obs.on_task_completed(task)
except Exception as e:
for obs in observers:
obs.on_task_failed(task, e)
self._pool.submit(run)
def shutdown(self):
self._shutdown = True
self._pool.shutdown(wait=True)
# โโโ Logging Observer โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class LoggingObserver(TaskObserver):
def on_task_started(self, task: Task):
print(f" [STARTED] {task}")
def on_task_completed(self, task: Task):
print(f" [COMPLETED] {task}")
def on_task_failed(self, task: Task, error: Exception):
print(f" [FAILED] {task} - {error}")
def on_task_cancelled(self, task: Task):
print(f" [CANCELLED] {task}")
# โโโ Task Scheduler โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class TaskScheduler:
def __init__(self, executor: ExecutionStrategy):
self._task_queue: list[ScheduledTask] = [] # min-heap
self._task_registry: dict[str, ScheduledTask] = {}
self._executor = executor
self._observers: list[TaskObserver] = []
self._lock = threading.Lock()
self._condition = threading.Condition(self._lock)
self._running = False
self._dispatcher_thread: Optional[threading.Thread] = None
def add_observer(self, observer: TaskObserver):
self._observers.append(observer)
def schedule(self, task: Task, delay_seconds: float) -> str:
run_time = time.time() + delay_seconds
st = ScheduledTask(task, run_time, 0, TaskType.DELAYED)
self._enqueue(st)
return task.id
def schedule_at(self, task: Task, timestamp: float) -> str:
st = ScheduledTask(task, timestamp, 0, TaskType.ONE_TIME)
self._enqueue(st)
return task.id
def schedule_recurring(self, task: Task, interval_seconds: float) -> str:
run_time = time.time() + interval_seconds
st = ScheduledTask(task, run_time, interval_seconds, TaskType.RECURRING)
self._enqueue(st)
return task.id
def cancel(self, task_id: str) -> bool:
st = self._task_registry.get(task_id)
if st is None or st.is_cancelled:
return False
st.cancel()
for obs in self._observers:
obs.on_task_cancelled(st.task)
return True
def _enqueue(self, st: ScheduledTask):
with self._condition:
heapq.heappush(self._task_queue, st)
self._task_registry[st.task.id] = st
self._condition.notify()
def start(self):
self._running = True
self._dispatcher_thread = threading.Thread(
target=self._dispatch, daemon=True, name="scheduler-dispatcher"
)
self._dispatcher_thread.start()
def _dispatch(self):
while self._running:
with self._condition:
while not self._task_queue and self._running:
self._condition.wait(timeout=0.1)
if not self._running:
break
if not self._task_queue:
continue
top = self._task_queue[0]
now = time.time()
if top.next_run_time <= now:
heapq.heappop(self._task_queue)
if top.is_cancelled:
self._task_registry.pop(top.task.id, None)
continue
self._executor.execute(top.task, self._observers)
if top.task_type == TaskType.RECURRING and not top.is_cancelled:
top.reschedule()
heapq.heappush(self._task_queue, top)
else:
self._task_registry.pop(top.task.id, None)
else:
wait_time = top.next_run_time - now
self._condition.wait(timeout=wait_time)
def shutdown(self):
self._running = False
with self._condition:
self._condition.notify_all()
self._executor.shutdown()
@property
def pending_count(self) -> int:
with self._lock:
return len(self._task_queue)
# โโโ Main Demo โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
def main():
print("โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ")
print(" TASK SCHEDULER - LLD DEMO ")
print("โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ\n")
executor = ThreadPoolStrategy(pool_size=4)
scheduler = TaskScheduler(executor)
scheduler.add_observer(LoggingObserver())
scheduler.start()
# โโโ Schedule One-Time Delayed Tasks โโโโโโโโโโโโโโโโ
print("--- Scheduling Delayed Tasks ---")
t1 = Task("SendEmail", Priority.HIGH, lambda: print(" >> Sending email..."))
t2 = Task("GenerateReport", Priority.LOW, lambda: print(" >> Generating report..."))
t3 = Task("CriticalAlert", Priority.CRITICAL, lambda: print(" >> CRITICAL ALERT FIRED!"))
scheduler.schedule(t1, 0.5)
scheduler.schedule(t2, 1.0)
scheduler.schedule(t3, 0.5)
print(f"Pending tasks: {scheduler.pending_count}")
time.sleep(2)
# โโโ Schedule Recurring Task โโโโโโโโโโโโโโโโโโโโโโโโ
print("\n--- Scheduling Recurring Task ---")
heartbeat = Task("Heartbeat", Priority.MEDIUM,
lambda: print(f" >> heartbeat ping at {time.time():.3f}"))
heartbeat_id = scheduler.schedule_recurring(heartbeat, 0.6)
time.sleep(2)
# โโโ Cancel Recurring Task โโโโโโโโโโโโโโโโโโโโโโโโโโ
print("\n--- Cancelling Heartbeat ---")
cancelled = scheduler.cancel(heartbeat_id)
print(f"Cancelled: {cancelled}")
time.sleep(1)
# โโโ Task That Fails โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
print("\n--- Task That Throws Exception ---")
def failing_action():
raise RuntimeError("Something went wrong!")
failing_task = Task("FailingJob", Priority.HIGH, failing_action)
scheduler.schedule(failing_task, 0.1)
time.sleep(0.5)
# โโโ Shutdown โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
print("\n--- Shutting Down Scheduler ---")
scheduler.shutdown()
print("\nโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ")
print(" DEMO COMPLETE ")
print("โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ")
if __name__ == "__main__":
main()
#include <iostream>
#include <string>
#include <queue>
#include <unordered_map>
#include <vector>
#include <thread>
#include <mutex>
#include <condition_variable>
#include <functional>
#include <chrono>
#include <atomic>
#include <memory>
#include <sstream>
#include <random>
// โโโ Enums โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
enum class Priority { LOW = 0, MEDIUM = 1, HIGH = 2, CRITICAL = 3 };
enum class TaskType { ONE_TIME, RECURRING, DELAYED };
std::string priorityToString(Priority p) {
switch (p) {
case Priority::LOW: return "LOW";
case Priority::MEDIUM: return "MEDIUM";
case Priority::HIGH: return "HIGH";
case Priority::CRITICAL: return "CRITICAL";
default: return "UNKNOWN";
}
}
// โโโ Generate Short ID โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
std::string generateId() {
static std::mt19937 gen(std::random_device{}());
static std::uniform_int_distribution<int> dist(0, 15);
static const char* hex = "0123456789abcdef";
std::string id;
for (int i = 0; i < 8; ++i) id += hex[dist(gen)];
return id;
}
long long nowMs() {
return std::chrono::duration_cast<std::chrono::milliseconds>(
std::chrono::system_clock::now().time_since_epoch()).count();
}
// โโโ Task โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class Task {
private:
std::string id;
std::string name;
Priority priority;
std::function<void()> action;
public:
Task(std::string name, Priority priority, std::function<void()> action)
: id(generateId()), name(std::move(name)), priority(priority), action(std::move(action)) {}
void execute() { action(); }
const std::string& getId() const { return id; }
const std::string& getName() const { return name; }
Priority getPriority() const { return priority; }
std::string toString() const {
return "Task[" + name + " | " + priorityToString(priority) + " | id=" + id + "]";
}
};
// โโโ ScheduledTask โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class ScheduledTask {
private:
std::shared_ptr<Task> task;
long long nextRunTimeMs;
long long intervalMs;
TaskType type;
std::atomic<bool> cancelled{false};
public:
ScheduledTask(std::shared_ptr<Task> task, long long nextRunTimeMs,
long long intervalMs, TaskType type)
: task(std::move(task)), nextRunTimeMs(nextRunTimeMs),
intervalMs(intervalMs), type(type) {}
void cancel() { cancelled = true; }
bool isCancelled() const { return cancelled; }
std::shared_ptr<Task> getTask() const { return task; }
long long getNextRunTimeMs() const { return nextRunTimeMs; }
TaskType getType() const { return type; }
long long getIntervalMs() const { return intervalMs; }
void reschedule() {
if (type == TaskType::RECURRING) {
nextRunTimeMs = nowMs() + intervalMs;
}
}
// For priority queue (min-heap: smallest time first, then highest priority)
bool operator>(const ScheduledTask& other) const {
if (nextRunTimeMs != other.nextRunTimeMs)
return nextRunTimeMs > other.nextRunTimeMs;
return static_cast<int>(task->getPriority()) < static_cast<int>(other.task->getPriority());
}
};
// โโโ Observer Interface โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class TaskObserver {
public:
virtual ~TaskObserver() = default;
virtual void onTaskStarted(const Task& task) = 0;
virtual void onTaskCompleted(const Task& task) = 0;
virtual void onTaskFailed(const Task& task, const std::exception& e) = 0;
virtual void onTaskCancelled(const Task& task) = 0;
};
// โโโ Logging Observer โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class LoggingObserver : public TaskObserver {
public:
void onTaskStarted(const Task& task) override {
std::cout << " [STARTED] " << task.toString() << "\n";
}
void onTaskCompleted(const Task& task) override {
std::cout << " [COMPLETED] " << task.toString() << "\n";
}
void onTaskFailed(const Task& task, const std::exception& e) override {
std::cout << " [FAILED] " << task.toString() << " - " << e.what() << "\n";
}
void onTaskCancelled(const Task& task) override {
std::cout << " [CANCELLED] " << task.toString() << "\n";
}
};
// โโโ Thread Pool โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class ThreadPool {
private:
std::vector<std::thread> workers;
std::queue<std::function<void()>> jobs;
std::mutex mtx;
std::condition_variable cv;
std::atomic<bool> stopped{false};
public:
ThreadPool(int size) {
for (int i = 0; i < size; ++i) {
workers.emplace_back([this] {
while (true) {
std::function<void()> job;
{
std::unique_lock<std::mutex> lock(mtx);
cv.wait(lock, [this] { return stopped || !jobs.empty(); });
if (stopped && jobs.empty()) return;
job = std::move(jobs.front());
jobs.pop();
}
job();
}
});
}
}
void submit(std::function<void()> job) {
{
std::lock_guard<std::mutex> lock(mtx);
jobs.push(std::move(job));
}
cv.notify_one();
}
void shutdown() {
stopped = true;
cv.notify_all();
for (auto& w : workers) {
if (w.joinable()) w.join();
}
}
};
// โโโ Task Scheduler โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
class TaskScheduler {
private:
std::priority_queue<ScheduledTask, std::vector<ScheduledTask>,
std::greater<ScheduledTask>> taskQueue;
std::unordered_map<std::string, std::shared_ptr<ScheduledTask>> taskRegistry;
std::unique_ptr<ThreadPool> pool;
std::vector<std::shared_ptr<TaskObserver>> observers;
std::mutex mtx;
std::condition_variable cv;
std::atomic<bool> running{false};
std::thread dispatcherThread;
public:
TaskScheduler(int poolSize) : pool(std::make_unique<ThreadPool>(poolSize)) {}
void addObserver(std::shared_ptr<TaskObserver> observer) {
observers.push_back(std::move(observer));
}
std::string schedule(std::shared_ptr<Task> task, long long delayMs) {
long long runTime = nowMs() + delayMs;
auto st = std::make_shared<ScheduledTask>(task, runTime, 0, TaskType::DELAYED);
enqueue(st);
return task->getId();
}
std::string scheduleRecurring(std::shared_ptr<Task> task, long long intervalMs) {
long long runTime = nowMs() + intervalMs;
auto st = std::make_shared<ScheduledTask>(task, runTime, intervalMs, TaskType::RECURRING);
enqueue(st);
return task->getId();
}
bool cancel(const std::string& taskId) {
std::lock_guard<std::mutex> lock(mtx);
auto it = taskRegistry.find(taskId);
if (it == taskRegistry.end() || it->second->isCancelled()) return false;
it->second->cancel();
for (auto& obs : observers) obs->onTaskCancelled(*it->second->getTask());
return true;
}
void start() {
running = true;
dispatcherThread = std::thread([this] { dispatch(); });
}
void shutdown() {
running = false;
cv.notify_all();
if (dispatcherThread.joinable()) dispatcherThread.join();
pool->shutdown();
}
int getPendingCount() {
std::lock_guard<std::mutex> lock(mtx);
return static_cast<int>(taskQueue.size());
}
private:
void enqueue(std::shared_ptr<ScheduledTask> st) {
std::lock_guard<std::mutex> lock(mtx);
taskRegistry[st->getTask()->getId()] = st;
taskQueue.push(*st);
cv.notify_one();
}
void dispatch() {
while (running) {
std::unique_lock<std::mutex> lock(mtx);
if (taskQueue.empty()) {
cv.wait_for(lock, std::chrono::milliseconds(100));
continue;
}
auto top = taskQueue.top();
long long now = nowMs();
if (top.getNextRunTimeMs() <= now) {
taskQueue.pop();
if (top.isCancelled()) {
taskRegistry.erase(top.getTask()->getId());
continue;
}
auto task = top.getTask();
auto& obs = observers;
pool->submit([task, &obs] {
for (auto& o : obs) o->onTaskStarted(*task);
try {
task->execute();
for (auto& o : obs) o->onTaskCompleted(*task);
} catch (const std::exception& e) {
for (auto& o : obs) o->onTaskFailed(*task, e);
}
});
if (top.getType() == TaskType::RECURRING && !top.isCancelled()) {
auto regIt = taskRegistry.find(task->getId());
if (regIt != taskRegistry.end()) {
regIt->second->reschedule();
taskQueue.push(*regIt->second);
}
} else {
taskRegistry.erase(task->getId());
}
} else {
long long waitMs = top.getNextRunTimeMs() - now;
cv.wait_for(lock, std::chrono::milliseconds(waitMs));
}
}
}
};
// โโโ Main Demo โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
int main() {
std::cout << "โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ\n";
std::cout << " TASK SCHEDULER - LLD DEMO \n";
std::cout << "โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ\n\n";
TaskScheduler scheduler(4);
scheduler.addObserver(std::make_shared<LoggingObserver>());
scheduler.start();
// โโโ Schedule Delayed Tasks โโโโโโโโโโโโโโโโโโโโโโโโโ
std::cout << "--- Scheduling Delayed Tasks ---\n";
auto t1 = std::make_shared<Task>("SendEmail", Priority::HIGH,
[] { std::cout << " >> Sending email...\n"; });
auto t2 = std::make_shared<Task>("GenerateReport", Priority::LOW,
[] { std::cout << " >> Generating report...\n"; });
auto t3 = std::make_shared<Task>("CriticalAlert", Priority::CRITICAL,
[] { std::cout << " >> CRITICAL ALERT FIRED!\n"; });
scheduler.schedule(t1, 500);
scheduler.schedule(t2, 1000);
scheduler.schedule(t3, 500);
std::cout << "Pending tasks: " << scheduler.getPendingCount() << "\n";
std::this_thread::sleep_for(std::chrono::seconds(2));
// โโโ Schedule Recurring Task โโโโโโโโโโโโโโโโโโโโโโโโ
std::cout << "\n--- Scheduling Recurring Task ---\n";
auto heartbeat = std::make_shared<Task>("Heartbeat", Priority::MEDIUM,
[] { std::cout << " >> heartbeat ping at " << nowMs() << "\n"; });
std::string heartbeatId = scheduler.scheduleRecurring(heartbeat, 600);
std::this_thread::sleep_for(std::chrono::seconds(2));
// โโโ Cancel Recurring Task โโโโโโโโโโโโโโโโโโโโโโโโโโ
std::cout << "\n--- Cancelling Heartbeat ---\n";
bool cancelled = scheduler.cancel(heartbeatId);
std::cout << "Cancelled: " << (cancelled ? "true" : "false") << "\n";
std::this_thread::sleep_for(std::chrono::seconds(1));
// โโโ Task That Fails โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
std::cout << "\n--- Task That Throws Exception ---\n";
auto failingTask = std::make_shared<Task>("FailingJob", Priority::HIGH,
[] { throw std::runtime_error("Something went wrong!"); });
scheduler.schedule(failingTask, 100);
std::this_thread::sleep_for(std::chrono::milliseconds(500));
// โโโ Shutdown โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
std::cout << "\n--- Shutting Down Scheduler ---\n";
scheduler.shutdown();
std::cout << "\nโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ\n";
std::cout << " DEMO COMPLETE \n";
std::cout << "โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ\n";
return 0;
}
State Transitions
stateDiagram-v2
[*] --> SCHEDULED
SCHEDULED --> RUNNING : dispatcher picks up
SCHEDULED --> CANCELLED : cancel() called
RUNNING --> COMPLETED : execution succeeds
RUNNING --> FAILED : exception thrown
COMPLETED --> SCHEDULED : recurring reschedule
CANCELLED --> [*]
FAILED --> [*]
COMPLETED --> [*]
Sequence Diagram - Schedule and Execute
sequenceDiagram
participant Client
participant TS as TaskScheduler
participant Lock as ReentrantLock
participant PQ as PriorityQueue
participant Exec as ThreadPoolStrategy
participant Obs as TaskObserver
Client->>TS: schedule(task, delayMs)
TS->>Lock: lock()
TS->>PQ: offer(scheduledTask)
TS->>Lock: signal + unlock()
TS-->>Client: return taskId
Note over TS: Dispatcher thread wakes
TS->>PQ: peek() - check if due
TS->>PQ: poll() - extract task
TS->>Exec: execute(task)
Exec->>Obs: onTaskStarted(task)
Exec->>Exec: task.execute()
Exec->>Obs: onTaskCompleted(task)
How to Extend
| Extension | Implementation |
|---|---|
| Cron expressions | Parse cron string to compute nextRunTime, new CronScheduledTask subclass |
| Task dependencies | DAG of tasks; only enqueue a task when all predecessors complete |
| Persistence | Serialize task queue to DB; reload on restart for durability |
| Distributed scheduler | Use distributed lock (Redis/ZooKeeper) for leader election among scheduler instances |
| Rate limiting | Token bucket before executor.execute() to cap throughput |
| Task retries | On failure, re-enqueue with exponential backoff (max 3 retries) |
What Interviewers Look For
- โ Min-heap for efficient next-task retrieval
- โ Thread pool for bounded concurrent execution
- โ Strategy pattern for swappable execution models
- โ Observer pattern for lifecycle notifications
- โ Thread-safety with locks and condition variables
- โ Recurring task support with clean reschedule logic
- โ Cancellation that doesnโt corrupt the queue
- โ Graceful shutdown - no abrupt thread kills
Related Concepts
Scale this design past a single process and these are the concepts it runs into:
- Durable Execution โ โ Temporal, Cadence and Step Functions are what make a schedule survive a process restart
- Message Queues โ โ a broker replaces the in-process heap once workers live on separate machines
- Leader Election โ โ only one scheduler instance should be dispatching a given schedule
- Retry & Backoff โ โ a failed task needs bounded, backed-off retries rather than immediate re-execution
- Dead Letter Queue โ โ a task that keeps failing needs somewhere to land instead of retrying forever
- Idempotency โ โ at-least-once dispatch means a task body has to be safe to run twice
Discussion
Newest first