設計任務排程器 — Priority Queue、Time Wheel 與 Cron 排程實戰 | 資料結構與演算法
任務排程器(Task Scheduler) 是後端系統中不可或缺的基礎設施——從簡單的 Cron Job 定時備份,到複雜的 DAG 依賴工作流,皆依賴排程器精確管理「什麼任務」在「何時」以「什麼優先序」被執行。本文將從需求分析出發,深入設計以 Priority Queue(優先佇列) 與 Timer Wheel(時間輪) 為核心的排程系統,涵蓋 Cron 表達式解析、DAG 拓撲排序依賴解析、指數退避重試 與 分散式任務調度,搭配 JavaScript/TypeScript 與 C++ 雙語言完整實作,帶你徹底掌握任務排程器的設計與實戰。
前言
在上一篇文章中,我們實作了
搜尋自動補全(Autocomplete),學會如何用 Trie 和 Top-K 排序實現即時的搜尋建議功能。今天我們要解決另一個在實際工作中無處不在的系統設計問題——任務排程器(Task Scheduler)。
你每天都在使用排程器,只是可能沒有意識到。手機上的鬧鐘 App、電子郵件的定時發送、電商平台的「下單 30 分鐘未付款自動取消」、每天凌晨自動跑的資料庫備份——這些背後都有一個排程器在默默運作。更複雜的場景包括 CI/CD 流程中的 DAG 依賴工作流、機器學習訓練管線中的多步驟任務編排、以及微服務架構中的分散式任務調度。
這個問題的關鍵挑戰在於:如何同時支援 即時任務、延遲任務 和 週期任務,並且在百萬級任務規模下保持高吞吐量和低延遲。暴力的 setTimeout 顯然撐不住——我們需要精心設計的資料結構來管理時間軸上的任務調度。答案是我們在
中學過的 Priority Queue(優先佇列) 和高效的 Timer Wheel(時間輪) 演算法。
本文你將學到:
- 任務排程器的完整需求分析與規模估算
- 四種方案比較:Priority Queue、單層時間輪、多層時間輪、Delay Queue
- 完整可運行的 JavaScript/TypeScript 與 C++ 雙語言實作(含 Cron 解析器與 DAG 引擎)
- 分散式排程架構:Leader Election、冪等性保證、Dead Letter Queue
- 指數退避重試策略的設計原理
- 4 道 LeetCode 經典排程題目
1. 需求分析
功能性需求(Functional Requirements)
- 支援 即時任務(Immediate Task):按優先級排序,立即分發到 Worker 執行
- 支援 延遲任務(Delayed Task):指定 N 秒後執行,如「訂單 30 分鐘後未付款自動取消」
- 支援 週期任務(Recurring Task):使用 Cron 表達式定義,如「每週一至週五 9:00 跑報表」
- 支援 任務依賴(DAG Workflow):Task A 完成後才能啟動 Task B(有向無環圖關係)
- 任務具有 優先級(Priority):高優先級任務插隊執行
- 支援任務 取消(Cancel) 與狀態查詢
- 任務失敗時自動 重試(Retry),支援指數退避(Exponential Backoff)
非功能性需求(Non-Functional Requirements)
- 調度精度:延遲任務誤差 < 100ms
- 吞吐量:每秒可排程 100,000 個任務
- 可靠性:任務不重複執行、不丟失(At-Least-Once 語義 + 冪等性)
- 可用性:99.99%(排程器本身不成為單點故障)
排除範圍
- 不設計任務執行引擎(Worker 的實際執行邏輯)
- 不設計任務結果儲存(Artifact Storage)
規模估算
以中型 SaaS 平台排程器為參考:
| 指標 | 數值 | 說明 |
|---|---|---|
| 即時任務 QPS | 10,000 | Priority Queue 排程 |
| 延遲任務 QPS | 5,000 | Timer Wheel 管理 |
| 固定 Cron Job 數量 | 1,000 | 週期性任務 |
| DAG 工作流 | 100/小時 | 多步驟依賴任務 |
| 單任務記錄大小 | ~512 B | 任務 ID + 類型 + Payload + 狀態 |
| Priority Queue 最大容量 | ~100 萬任務 | 記憶體約 512MB |
| Timer Wheel 記憶體 | ~28 MB | 360 萬槽位 x 8B 指標 |
結論:記憶體瓶頸不在 Timer Wheel(固定且極小),而在 Priority Queue 中的任務數量。核心挑戰是排程器本身的高可用——避免單點故障。
2. 方案設計
2.1 方案一:Priority Queue(Min-Heap)
以 executeAt 時間戳作為 Min-Heap 的排序鍵,每次取出堆頂——即最早該執行的任務。排程迴圈 sleep 至堆頂任務的到期時間,避免忙碌等待(Busy Waiting)。
Min-Heap(以 executeAt 時間戳排序):
[Task#5, t=1000]
/ \
[Task#3, t=1200] [Task#8, t=1500]
/ \
[Task#1, t=1800] [Task#6, t=2000]
時間推進到 t=1000:
pop() → Task#5,送往 Worker 執行
heap 自動維持 Task#3 成為新堆頂
排程主迴圈:
sleep(heap.peek().executeAt - now())
→ 避免忙碌等待,精準喚醒
優點:實作簡單,精度不受限。缺點:插入和取出都是 O(log n),百萬級任務時效能下降。
2.2 方案二:單層 Timer Wheel(Hashed Timer Wheel)
Timer Wheel 是 Netty、Linux Kernel、Kafka 等系統處理大量定時任務的核心資料結構。將時間軸映射到固定大小的環形陣列,每個槽位掛載同一時刻到期的任務鏈結串列,指針每 tick 前進一格。
單層 Timer Wheel(精度=10ms,範圍=1秒,共 100 槽):
┌─────────────────────────────────────────────────────┐
│ 槽位: [0] [1] [2] [3] [4] ... [98] [99] │
│ 時間戳: 0 10 20 30 40 ... 980 990ms │
│ │
│ 指針 → [當前槽=42,對應時間 420ms] │
│ │
│ 槽位42: Task#3 → Task#7 → Task#15 → null │
│ (鏈結串列儲存同一時間觸發的多個任務) │
└─────────────────────────────────────────────────────┘
插入 Task(delay=580ms):
目標槽 = (42 + 580/10) % 100 = (42 + 58) % 100 = 0
掛載至槽位 [0],記錄 rounds = 0(未超過一圈)
優點:插入和觸發都是 O(1),記憶體固定。缺點:精度受槽位粒度限制,超過一圈需額外管理 rounds。
2.3 方案三:多層 Timer Wheel(Hierarchical Timer Wheel)
模仿時鐘的時 / 分 / 秒結構,用多個粒度不同的時間輪組成層級。粗粒度層管理遠期任務,到期時「降級」到精粒度層精確排程。
三層 Timer Wheel:
第三層(粗粒度):
[0][1]...[255] — 每 65536ms 一格,範圍 ~4.6 小時
第二層(中粒度):
[0][1]...[255] — 每 256ms 一格,範圍 ~65 秒
第一層(精粒度):
[0][1]...[255] — 每 1ms 一格,範圍 256ms
插入 delay=3661 秒的任務:
先掛到第三層 → 第三層指針到達時降級到第二層
→ 第二層到達時降級到第一層精確執行
總槽位數:256 x 3 = 768(極省記憶體)
2.4 方案四:Delay Queue(基於訊息佇列)
將任務發送到 Kafka / RabbitMQ 等訊息佇列,利用佇列本身的 delayed message 功能實現延遲排程。
方案比較
| 策略 | 插入 | 觸發/取出 | 取消 | 記憶體 | 適用場景 |
|---|---|---|---|---|---|
| Priority Queue | O(log n) | O(log n) | O(n) | O(n) | 任務量 < 10 萬,需精確計時 |
| 單層 Timer Wheel | O(1) | O(1) 均攤 | O(1) | O(slots) | 百萬級延遲任務 |
| 多層 Timer Wheel | O(1) | O(1) 均攤 | O(1) | O(slots x layers) | 超大範圍 + 高精度 |
| Delay Queue | O(1) | O(1) | O(1) | 外部 | 分散式環境,已有 MQ 基礎設施 |
結論:本文選擇 混合方案——Timer Wheel 管理大量延遲任務,Priority Queue 處理即時任務的精確排序,Cron Parser 計算週期任務的下次執行時間,DAG Engine 解析任務依賴。
3. 核心實作——JavaScript/TypeScript
3.1 型別定義與 Priority Queue
// ============================================================
// Task Scheduler — TypeScript 完整實作
// 包含:Priority Queue + Timer Wheel + DAG 依賴 + Cron 解析器
// ============================================================
// ─── 型別定義 ───────────────────────────────────────────────
type TaskStatus = "pending" | "ready" | "running" | "success" | "failed" | "cancelled";
type TaskType = "immediate" | "delayed" | "recurring";
type TaskPriority = 1 | 2 | 3 | 4 | 5; // 1=最高, 5=最低
interface Task {
id: string;
name: string;
type: TaskType;
priority: TaskPriority;
executeAt: number; // Unix timestamp(ms)
cronExpression?: string; // 週期任務的 Cron 表達式
dependencies: string[]; // 依賴的任務 ID 列表
payload: Record<string, unknown>;
status: TaskStatus;
retryCount: number;
maxRetries: number;
createdAt: number;
execute: () => Promise<void>;
}
// ─── Min-Heap Priority Queue ────────────────────────────────
class TaskPriorityQueue {
private heap: Task[] = [];
/**
* 比較函式:先按 executeAt 排序,再按 priority(數值小 = 優先級高)
*/
private compare(a: Task, b: Task): boolean {
if (a.executeAt !== b.executeAt) return a.executeAt < b.executeAt;
return a.priority < b.priority;
}
push(task: Task): void {
this.heap.push(task);
this.bubbleUp(this.heap.length - 1);
}
pop(): Task | undefined {
if (this.heap.length === 0) return undefined;
const top = this.heap[0];
const last = this.heap.pop()!;
if (this.heap.length > 0) {
this.heap[0] = last;
this.sinkDown(0);
}
return top;
}
peek(): Task | undefined {
return this.heap[0];
}
get size(): number {
return this.heap.length;
}
removeById(taskId: string): boolean {
const idx = this.heap.findIndex(t => t.id === taskId);
if (idx === -1) return false;
// 標記刪除(lazy deletion)— 取出時再跳過
this.heap[idx].status = "cancelled";
return true;
}
private bubbleUp(i: number): void {
while (i > 0) {
const parent = Math.floor((i - 1) / 2);
if (this.compare(this.heap[i], this.heap[parent])) {
[this.heap[i], this.heap[parent]] = [this.heap[parent], this.heap[i]];
i = parent;
} else break;
}
}
private sinkDown(i: number): void {
const n = this.heap.length;
while (true) {
let smallest = i;
const left = 2 * i + 1;
const right = 2 * i + 2;
if (left < n && this.compare(this.heap[left], this.heap[smallest])) smallest = left;
if (right < n && this.compare(this.heap[right], this.heap[smallest])) smallest = right;
if (smallest !== i) {
[this.heap[i], this.heap[smallest]] = [this.heap[smallest], this.heap[i]];
i = smallest;
} else break;
}
}
}
3.2 Cron 解析器
// ─── Cron 解析器 ─────────────────────────────────────────────
class CronParser {
/**
* 解析 Cron 表達式,計算下次執行時間
* 支援標準 5 欄 Cron:分 時 日 月 週
* 支援特殊值:* , - /
*/
static nextExecutionTime(expression: string, fromTime: Date = new Date()): Date {
const fields = expression.trim().split(/\s+/);
if (fields.length !== 5) {
throw new Error(`無效的 Cron 表達式:${expression}(需要 5 個欄位)`);
}
const [minExpr, hourExpr, domExpr, monthExpr, dowExpr] = fields;
const parsedMin = this.parseField(minExpr, 0, 59);
const parsedHour = this.parseField(hourExpr, 0, 23);
const parsedDom = this.parseField(domExpr, 1, 31);
const parsedMonth = this.parseField(monthExpr, 1, 12);
const parsedDow = this.parseField(dowExpr, 0, 6); // 0=週日, 6=週六
// 從 fromTime + 1 分鐘開始搜尋
const candidate = new Date(fromTime);
candidate.setSeconds(0, 0);
candidate.setMinutes(candidate.getMinutes() + 1);
// 最多搜尋 4 年(防止無限迴圈)
const maxDate = new Date(fromTime.getTime() + 4 * 365 * 24 * 60 * 60 * 1000);
while (candidate < maxDate) {
const month = candidate.getMonth() + 1; // 1-12
const dom = candidate.getDate();
const dow = candidate.getDay(); // 0-6
const hour = candidate.getHours();
const min = candidate.getMinutes();
if (!parsedMonth.has(month)) {
candidate.setMonth(candidate.getMonth() + 1, 1);
candidate.setHours(0, 0, 0, 0);
continue;
}
if (!parsedDom.has(dom) || !parsedDow.has(dow)) {
candidate.setDate(candidate.getDate() + 1);
candidate.setHours(0, 0, 0, 0);
continue;
}
if (!parsedHour.has(hour)) {
candidate.setHours(candidate.getHours() + 1, 0, 0, 0);
continue;
}
if (!parsedMin.has(min)) {
candidate.setMinutes(candidate.getMinutes() + 1, 0, 0);
continue;
}
return new Date(candidate);
}
throw new Error(`無法在 4 年內找到下次執行時間:${expression}`);
}
private static parseField(expr: string, min: number, max: number): Set<number> {
const result = new Set<number>();
// 星期別名
const dowAlias: Record<string, number> = {
SUN: 0, MON: 1, TUE: 2, WED: 3, THU: 4, FRI: 5, SAT: 6,
};
const resolveAlias = (s: string): number => {
const upper = s.toUpperCase();
return dowAlias[upper] ?? parseInt(s, 10);
};
for (const part of expr.split(",")) {
if (part === "*") {
for (let i = min; i <= max; i++) result.add(i);
} else if (part.includes("/")) {
const [range, step] = part.split("/");
const stepNum = parseInt(step, 10);
const start = range === "*" ? min : resolveAlias(range.split("-")[0]);
const end = range === "*"
? max
: range.includes("-") ? resolveAlias(range.split("-")[1]) : max;
for (let i = start; i <= end; i += stepNum) result.add(i);
} else if (part.includes("-")) {
const [start, end] = part.split("-").map(resolveAlias);
for (let i = start; i <= end; i++) result.add(i);
} else {
result.add(resolveAlias(part));
}
}
return result;
}
}
3.3 Timer Wheel 與 DAG Engine
// ─── Timer Wheel 實作 ──────────────────────────────────────
class TimerWheel {
private slots: Map<number, Task[]>;
private currentSlot: number = 0;
private readonly slotCount: number;
private readonly tickMs: number;
private intervalHandle?: ReturnType<typeof setInterval>;
constructor(slotCount: number = 3600, tickMs: number = 10) {
this.slotCount = slotCount;
this.tickMs = tickMs;
this.slots = new Map();
}
/** 插入延遲任務 — 時間複雜度:O(1) */
schedule(task: Task): void {
const delayMs = Math.max(0, task.executeAt - Date.now());
const ticks = Math.floor(delayMs / this.tickMs);
const rounds = Math.floor(ticks / this.slotCount);
const slot = (this.currentSlot + ticks) % this.slotCount;
// 記錄剩餘圈數
(task as any).__timerRounds = rounds;
if (!this.slots.has(slot)) this.slots.set(slot, []);
this.slots.get(slot)!.push(task);
}
/** 啟動 Timer Wheel */
start(onTaskReady: (task: Task) => void): void {
this.intervalHandle = setInterval(() => {
const slotTasks = this.slots.get(this.currentSlot) ?? [];
const remaining: Task[] = [];
for (const task of slotTasks) {
if (task.status === "cancelled") continue;
const rounds = (task as any).__timerRounds;
if (rounds === 0) {
onTaskReady(task);
} else {
(task as any).__timerRounds = rounds - 1;
remaining.push(task);
}
}
if (remaining.length > 0) {
this.slots.set(this.currentSlot, remaining);
} else {
this.slots.delete(this.currentSlot);
}
this.currentSlot = (this.currentSlot + 1) % this.slotCount;
}, this.tickMs);
}
stop(): void {
if (this.intervalHandle) clearInterval(this.intervalHandle);
}
}
// ─── DAG Engine — Kahn's 拓撲排序 ───────────────────────────
interface WorkflowResult {
order: string[];
parallelGroups: string[][];
hasCycle: boolean;
}
class DAGEngine {
/**
* Kahn's Algorithm — BFS 拓撲排序
* 同時計算可並行執行的任務群組
* 時間複雜度:O(V + E)
*/
static topoSort(tasks: Map<string, Task>): WorkflowResult {
const inDegree = new Map<string, number>();
const adjList = new Map<string, string[]>();
for (const [id] of tasks) {
inDegree.set(id, 0);
adjList.set(id, []);
}
for (const [id, task] of tasks) {
for (const dep of task.dependencies) {
if (!tasks.has(dep)) {
throw new Error(`任務 ${id} 依賴不存在的任務 ${dep}`);
}
inDegree.set(id, (inDegree.get(id) ?? 0) + 1);
adjList.get(dep)!.push(id);
}
}
const order: string[] = [];
const parallelGroups: string[][] = [];
let queue = [...inDegree.entries()]
.filter(([, deg]) => deg === 0)
.map(([id]) => id);
while (queue.length > 0) {
parallelGroups.push([...queue]);
const nextQueue: string[] = [];
for (const id of queue) {
order.push(id);
for (const neighbor of adjList.get(id) ?? []) {
const newDeg = (inDegree.get(neighbor) ?? 1) - 1;
inDegree.set(neighbor, newDeg);
if (newDeg === 0) nextQueue.push(neighbor);
}
}
queue = nextQueue;
}
const hasCycle = order.length !== tasks.size;
return { order, parallelGroups, hasCycle };
}
}
3.4 TaskScheduler 主排程器
// ─── TaskScheduler(主排程器)──────────────────────────────────
class TaskScheduler {
private priorityQueue: TaskPriorityQueue;
private timerWheel: TimerWheel;
private allTasks: Map<string, Task> = new Map();
private completedTasks: Set<string> = new Set();
private runningTasks: Set<string> = new Set();
private readonly maxConcurrency: number;
private isRunning: boolean = false;
private schedulerInterval?: ReturnType<typeof setInterval>;
constructor(maxConcurrency: number = 10) {
this.maxConcurrency = maxConcurrency;
this.priorityQueue = new TaskPriorityQueue();
this.timerWheel = new TimerWheel(360000, 1); // 100 小時,1ms 精度
}
/** 提交即時/延遲任務 */
submit(task: Omit<Task, "status" | "retryCount" | "createdAt">): string {
const fullTask: Task = {
...task,
status: "pending",
retryCount: 0,
createdAt: Date.now(),
};
this.allTasks.set(task.id, fullTask);
if (task.dependencies.length === 0) {
this.enqueueTask(fullTask);
}
console.log(`[Scheduler] 任務已提交:${task.id} (${task.type}, priority=${task.priority})`);
return task.id;
}
/** 提交 DAG 工作流 */
submitWorkflow(tasks: Omit<Task, "status" | "retryCount" | "createdAt">[]): WorkflowResult {
const taskMap = new Map<string, Task>();
for (const t of tasks) {
const fullTask: Task = {
...t,
status: "pending",
retryCount: 0,
createdAt: Date.now(),
};
this.allTasks.set(t.id, fullTask);
taskMap.set(t.id, fullTask);
}
const result = DAGEngine.topoSort(taskMap);
if (result.hasCycle) {
throw new Error("工作流包含循環依賴,無法執行");
}
console.log("[Scheduler] 工作流拓撲排序結果:");
result.parallelGroups.forEach((group, i) => {
console.log(` 第 ${i + 1} 波(可並行):${group.join(", ")}`);
});
// 入度為 0 的任務加入排程
for (const id of (result.parallelGroups[0] ?? [])) {
this.enqueueTask(taskMap.get(id)!);
}
return result;
}
/** 提交週期任務(Cron) */
submitCron(
taskTemplate: Omit<Task, "id" | "status" | "retryCount" | "createdAt" | "executeAt" | "type">,
cronExpression: string
): void {
const scheduleNext = () => {
const nextTime = CronParser.nextExecutionTime(cronExpression);
const taskId = `${taskTemplate.name}_${nextTime.getTime()}`;
const cronTask: Task = {
...taskTemplate,
id: taskId,
type: "recurring",
status: "pending",
retryCount: 0,
createdAt: Date.now(),
executeAt: nextTime.getTime(),
cronExpression,
execute: async () => {
await taskTemplate.execute();
scheduleNext(); // 執行完畢後安排下次
},
};
this.allTasks.set(taskId, cronTask);
this.timerWheel.schedule(cronTask);
console.log(`[Scheduler] Cron 任務排程:${taskTemplate.name},下次:${nextTime.toISOString()}`);
};
scheduleNext();
}
/** 取消任務 */
cancel(taskId: string): boolean {
const task = this.allTasks.get(taskId);
if (!task || task.status === "running") return false;
task.status = "cancelled";
this.priorityQueue.removeById(taskId);
console.log(`[Scheduler] 任務已取消:${taskId}`);
return true;
}
/** 查詢任務狀態 */
getStatus(taskId: string): TaskStatus | null {
return this.allTasks.get(taskId)?.status ?? null;
}
/** 啟動排程器主迴圈 */
start(): void {
if (this.isRunning) return;
this.isRunning = true;
this.timerWheel.start((task) => {
if (task.status === "pending") this.enqueueTask(task);
});
this.schedulerInterval = setInterval(() => this.processQueue(), 10);
console.log("[Scheduler] 排程器啟動");
}
stop(): void {
this.isRunning = false;
this.timerWheel.stop();
if (this.schedulerInterval) clearInterval(this.schedulerInterval);
console.log("[Scheduler] 排程器停止");
}
private enqueueTask(task: Task): void {
if (task.type === "delayed" && task.executeAt > Date.now() + 10) {
this.timerWheel.schedule(task);
} else {
task.status = "ready";
this.priorityQueue.push(task);
}
}
private processQueue(): void {
while (
this.runningTasks.size < this.maxConcurrency &&
this.priorityQueue.size > 0
) {
const task = this.priorityQueue.peek();
if (!task) break;
if (task.status === "cancelled") { this.priorityQueue.pop(); continue; }
if (task.executeAt > Date.now()) break;
this.priorityQueue.pop();
this.executeTask(task);
}
}
private async executeTask(task: Task): Promise<void> {
task.status = "running";
this.runningTasks.add(task.id);
try {
await task.execute();
task.status = "success";
this.completedTasks.add(task.id);
console.log(`[Scheduler] 任務完成:${task.id}`);
this.onTaskComplete(task.id);
} catch (error) {
console.error(`[Scheduler] 任務失敗:${task.id}`, error);
if (task.retryCount < task.maxRetries) {
task.retryCount++;
task.status = "pending";
// 指數退避:2^retryCount 秒
task.executeAt = Date.now() + 2 ** task.retryCount * 1000;
this.enqueueTask(task);
console.log(`[Scheduler] 重試(${task.retryCount}/${task.maxRetries}):${task.id}`);
} else {
task.status = "failed";
console.log(`[Scheduler] 任務永久失敗(進入 Dead Letter Queue):${task.id}`);
this.onTaskComplete(task.id);
}
} finally {
this.runningTasks.delete(task.id);
}
}
/** 任務完成後,檢查依賴此任務的任務是否可就緒 */
private onTaskComplete(completedId: string): void {
for (const [, task] of this.allTasks) {
if (task.status !== "pending") continue;
if (!task.dependencies.includes(completedId)) continue;
const allDepsComplete = task.dependencies.every(
depId => this.completedTasks.has(depId)
);
if (allDepsComplete) {
console.log(`[Scheduler] 依賴解除,任務就緒:${task.id}`);
task.executeAt = Date.now();
this.enqueueTask(task);
}
}
}
}
3.5 完整使用示範
// ─── 使用示範 ──────────────────────────────────────────────
async function demo(): Promise<void> {
console.log("=== 任務排程器 Demo ===\n");
const scheduler = new TaskScheduler(5);
scheduler.start();
// 建立任務的輔助函式
const makeTask = (
id: string, name: string, priority: TaskPriority,
deps: string[] = [], delay: number = 0
): Omit<Task, "status" | "retryCount" | "createdAt"> => ({
id, name,
type: delay > 0 ? "delayed" : "immediate",
priority,
executeAt: Date.now() + delay,
dependencies: deps,
payload: {},
maxRetries: 2,
execute: async () => {
await new Promise(r => setTimeout(r, 100));
console.log(` → 執行任務:${name} (ID: ${id})`);
},
});
// 1. 即時優先級任務
scheduler.submit(makeTask("t1", "高優先級緊急任務", 1));
scheduler.submit(makeTask("t2", "普通任務A", 3));
scheduler.submit(makeTask("t3", "低優先級任務", 5));
// 2. DAG 工作流(ML 訓練管線)
console.log("\n--- 提交 ML 訓練工作流 ---");
scheduler.submitWorkflow([
makeTask("fetch", "抓取資料", 3),
makeTask("augment", "擴增資料", 3),
makeTask("clean", "清洗資料", 2, ["fetch"]),
makeTask("transform", "格式轉換", 2, ["augment"]),
makeTask("feature", "特徵工程", 2, ["clean", "transform"]),
makeTask("train", "訓練模型", 1, ["feature"]),
]);
// 3. Cron 排程(示範計算下次時間)
console.log("\n--- Cron 排程 ---");
const nextTime = CronParser.nextExecutionTime("0 9 * * MON-FRI");
console.log(` 每週一至週五 9:00,下次:${nextTime.toLocaleString("zh-TW")}`);
const nextEveryFive = CronParser.nextExecutionTime("*/5 * * * *");
console.log(` 每 5 分鐘,下次:${nextEveryFive.toLocaleString("zh-TW")}`);
// 4. 延遲任務
scheduler.submit(makeTask("delayed1", "3 秒後執行", 2, [], 3000));
await new Promise(r => setTimeout(r, 2000));
scheduler.stop();
console.log("\n=== Demo 結束 ===");
}
demo().catch(console.error);
// 輸出:
// [Scheduler] 排程器啟動
// [Scheduler] 任務已提交:t1 (immediate, priority=1)
// [Scheduler] 任務已提交:t2 (immediate, priority=3)
// [Scheduler] 任務已提交:t3 (immediate, priority=5)
// → 執行任務:高優先級緊急任務 (ID: t1)
// [Scheduler] 任務完成:t1
// → 執行任務:普通任務A (ID: t2)
// [Scheduler] 任務完成:t2
// → 執行任務:低優先級任務 (ID: t3)
// [Scheduler] 任務完成:t3
// --- 提交 ML 訓練工作流 ---
// [Scheduler] 工作流拓撲排序結果:
// 第 1 波(可並行):fetch, augment
// 第 2 波(可並行):clean, transform
// 第 3 波(可並行):feature
// 第 4 波(可並行):train
// --- Cron 排程 ---
// 每週一至週五 9:00,下次:2026/7/27 09:00:00
// [Scheduler] 排程器停止
// === Demo 結束 ===
4. 核心實作——C++
C++ 版本提供高效能的 Timer Wheel 實作,利用多執行緒和精確計時實現百萬級延遲任務排程。
4.1 單層與多層 Timer Wheel
// ============================================================
// Timer Wheel — C++ 高效實作
// 多層時間輪,支援百萬級延遲任務,O(1) 插入與觸發
// ============================================================
#include <vector>
#include <list>
#include <functional>
#include <string>
#include <memory>
#include <iostream>
#include <chrono>
#include <thread>
#include <mutex>
#include <atomic>
#include <unordered_map>
#include <queue>
// ─── 任務定義 ──────────────────────────────────────────────
struct TimerTask {
std::string id;
std::function<void()> callback;
int64_t executeAtMs; // 絕對執行時間(ms)
int32_t rounds; // 剩餘圈數
bool cancelled = false;
};
using TaskPtr = std::shared_ptr<TimerTask>;
// ─── 單層 Timer Wheel ──────────────────────────────────────
class SingleTimerWheel {
public:
explicit SingleTimerWheel(int slotCount = 512, int tickMs = 10)
: slotCount_(slotCount), tickMs_(tickMs),
slots_(slotCount), currentSlot_(0) {}
/** 插入任務 — O(1) */
void insert(TaskPtr task) {
std::lock_guard<std::mutex> lock(mutex_);
int64_t delayMs = std::max(0LL, task->executeAtMs - nowMs());
int64_t ticks = delayMs / tickMs_;
task->rounds = static_cast<int32_t>(ticks / slotCount_);
int slot = (currentSlot_ + static_cast<int>(ticks % slotCount_)) % slotCount_;
slots_[slot].push_back(task);
taskIndex_[task->id] = task;
}
/** 取消任務 — O(1) 標記刪除 */
bool cancel(const std::string& id) {
std::lock_guard<std::mutex> lock(mutex_);
auto it = taskIndex_.find(id);
if (it == taskIndex_.end()) return false;
it->second->cancelled = true;
taskIndex_.erase(it);
return true;
}
/** 推進一個 tick,回傳到期任務 */
std::vector<TaskPtr> tick() {
std::lock_guard<std::mutex> lock(mutex_);
std::vector<TaskPtr> ready;
auto& slot = slots_[currentSlot_];
for (auto it = slot.begin(); it != slot.end(); ) {
auto& task = *it;
if (task->cancelled) {
it = slot.erase(it);
continue;
}
if (task->rounds == 0) {
ready.push_back(task);
taskIndex_.erase(task->id);
it = slot.erase(it);
} else {
task->rounds--;
++it;
}
}
currentSlot_ = (currentSlot_ + 1) % slotCount_;
return ready;
}
private:
static int64_t nowMs() {
return std::chrono::duration_cast<std::chrono::milliseconds>(
std::chrono::steady_clock::now().time_since_epoch()
).count();
}
int slotCount_;
int64_t tickMs_;
std::vector<std::list<TaskPtr>> slots_;
int currentSlot_;
std::mutex mutex_;
std::unordered_map<std::string, TaskPtr> taskIndex_;
};
4.2 多層 Timer Wheel + Priority Queue 排程器
// ─── 多層 Timer Wheel(Hierarchical)────────────────────────
class HierarchicalTimerWheel {
public:
HierarchicalTimerWheel()
: wheel1_(256, 1), // 256 槽,1ms/tick
wheel2_(256, 256), // 256 槽,256ms/tick
wheel3_(256, 65536), // 256 槽,65536ms/tick
running_(false) {}
void insert(TaskPtr task) {
int64_t delayMs = task->executeAtMs - nowMs();
if (delayMs < 256) {
wheel1_.insert(task);
} else if (delayMs < 65536) {
wheel2_.insert(task);
} else {
wheel3_.insert(task);
}
}
bool cancel(const std::string& id) {
return wheel1_.cancel(id) || wheel2_.cancel(id) || wheel3_.cancel(id);
}
/** 啟動後台 Tick 執行緒 */
void start(std::function<void(TaskPtr)> onTaskReady) {
running_ = true;
tickThread_ = std::thread([this, onTaskReady]() {
while (running_) {
// 推進第一層
auto tasks = wheel1_.tick();
for (auto& task : tasks) {
if (!task->cancelled) onTaskReady(task);
}
// 每 256 次推進第二層(降級到第一層)
if (tick1Count_++ % 256 == 0) {
auto tasks2 = wheel2_.tick();
for (auto& task : tasks2) wheel1_.insert(task);
}
// 每 65536 次推進第三層
if (tick2Count_++ % 256 == 0) {
auto tasks3 = wheel3_.tick();
for (auto& task : tasks3) wheel2_.insert(task);
}
std::this_thread::sleep_for(std::chrono::milliseconds(1));
}
});
}
void stop() {
running_ = false;
if (tickThread_.joinable()) tickThread_.join();
}
private:
static int64_t nowMs() {
return std::chrono::duration_cast<std::chrono::milliseconds>(
std::chrono::steady_clock::now().time_since_epoch()
).count();
}
SingleTimerWheel wheel1_, wheel2_, wheel3_;
std::atomic<bool> running_;
std::thread tickThread_;
uint64_t tick1Count_ = 0;
uint64_t tick2Count_ = 0;
};
// ─── Priority Queue 排程器(對比實作)─────────────────────────
class PriorityQueueScheduler {
public:
struct TaskComparator {
bool operator()(const TaskPtr& a, const TaskPtr& b) const {
return a->executeAtMs > b->executeAtMs; // Min-Heap
}
};
void insert(TaskPtr task) {
std::lock_guard<std::mutex> lock(mutex_);
pq_.push(task);
cv_.notify_one();
}
void start(std::function<void(TaskPtr)> onTaskReady) {
running_ = true;
thread_ = std::thread([this, onTaskReady]() {
while (running_) {
std::unique_lock<std::mutex> lock(mutex_);
if (pq_.empty()) {
cv_.wait_for(lock, std::chrono::milliseconds(100));
continue;
}
auto now = nowMs();
if (pq_.top()->executeAtMs <= now) {
auto task = pq_.top();
pq_.pop();
lock.unlock();
if (!task->cancelled) onTaskReady(task);
} else {
auto sleepMs = pq_.top()->executeAtMs - now;
cv_.wait_for(lock, std::chrono::milliseconds(sleepMs));
}
}
});
}
void stop() {
running_ = false;
cv_.notify_all();
if (thread_.joinable()) thread_.join();
}
private:
static int64_t nowMs() {
return std::chrono::duration_cast<std::chrono::milliseconds>(
std::chrono::steady_clock::now().time_since_epoch()
).count();
}
std::priority_queue<TaskPtr, std::vector<TaskPtr>, TaskComparator> pq_;
std::mutex mutex_;
std::condition_variable cv_;
std::atomic<bool> running_{false};
std::thread thread_;
};
// ─── main 功能示範 ─────────────────────────────────────────
int main() {
std::cout << "=== Timer Wheel 排程器 Demo ===\n\n";
HierarchicalTimerWheel scheduler;
std::atomic<int> executedCount{0};
std::mutex printMutex;
scheduler.start([&](TaskPtr task) {
std::lock_guard<std::mutex> lock(printMutex);
std::cout << " [執行] " << task->id << "\n";
task->callback();
executedCount++;
});
auto nowMs = []() {
return std::chrono::duration_cast<std::chrono::milliseconds>(
std::chrono::steady_clock::now().time_since_epoch()
).count();
};
// 插入不同延遲的任務
std::vector<std::pair<std::string, int>> taskDefs = {
{"task-immediate", 0},
{"task-50ms", 50},
{"task-200ms", 200},
{"task-500ms", 500},
{"task-to-cancel", 300},
};
for (const auto& [id, delay] : taskDefs) {
auto task = std::make_shared<TimerTask>();
task->id = id;
task->executeAtMs = nowMs() + delay;
task->rounds = 0;
task->callback = [id]() {
std::cout << " → 回呼:" << id << "\n";
};
scheduler.insert(task);
std::cout << "已排程:" << id << "(delay=" << delay << "ms)\n";
}
// 取消一個任務
std::this_thread::sleep_for(std::chrono::milliseconds(100));
bool cancelled = scheduler.cancel("task-to-cancel");
std::cout << "\n取消 task-to-cancel:" << (cancelled ? "成功" : "失敗") << "\n";
std::this_thread::sleep_for(std::chrono::milliseconds(700));
scheduler.stop();
std::cout << "\n共執行任務:" << executedCount.load() << " 個\n";
// 輸出:共執行任務:4 個(task-to-cancel 已被取消)
std::cout << "=== Demo 結束 ===\n";
return 0;
}
5. 效能分析
時間複雜度
| 操作 | Priority Queue | Timer Wheel | 說明 |
|---|---|---|---|
| 插入任務 | O(log n) | O(1) | n = 佇列中任務數 |
| 取出/觸發 | O(log n) | O(1) 均攤 | Timer Wheel 每 tick O(k),k 為到期任務數 |
| 取消任務 | O(n) | O(1) | Timer Wheel 用標記刪除 |
| DAG 排程 | O(V + E) | — | V = 任務數,E = 依賴邊數 |
| Cron 解析 | O(1) 均攤 | — | 最壞情況遍歷所有候選時間 |
空間複雜度
| 結構 | 空間 | 說明 |
|---|---|---|
| Priority Queue | O(n) | n = 活躍任務數 |
| 單層 Timer Wheel | O(slots + n) | slots 固定,n = 任務數 |
| 多層 Timer Wheel | O(slots x layers + n) | 768 槽 x 3 層 = 2304 個槽 |
| DAG 鄰接表 | O(V + E) | 拓撲排序的圖結構 |
效能基準(參考值)
| 操作 | Timer Wheel | Priority Queue |
|---|---|---|
| 100 萬任務插入 | ~100ms | ~800ms |
| 適用任務規模 | 百萬級以上 | 十萬級以下 |
| 記憶體布局 | 固定陣列(cache-friendly) | 動態堆積 |
| 使用系統 | Netty, Kafka, Linux | JDK DelayQueue |
6. 生產環境考量
6.1 分散式排程架構
┌───────────────────────────────────────────────────────────────┐
│ Scheduler Cluster(排程器叢集) │
│ │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ Scheduler 1 │ │ Scheduler 2 │ │ Scheduler 3 │ │
│ │ (Active) │ │ (Standby) │ │ (Standby) │ │
│ └──────────────┘ └──────────────┘ └──────────────┘ │
│ │ │
│ └───── Leader Election (etcd / ZooKeeper) ──────── │
└───────────────────────────┬───────────────────────────────────┘
│ 分發任務
▼
┌───────────────────────────────────────────────────────────────┐
│ Message Broker(Kafka / RabbitMQ) │
│ Topic: tasks.high │ Topic: tasks.normal │ Topic: retry │
└───────────────────────────┬───────────────────────────────────┘
┌─────────────┼─────────────┐
▼ ▼ ▼
┌──────────┐ ┌──────────┐ ┌──────────┐
│ Worker 1 │ │ Worker 2 │ │ Worker N │ ← 水平擴展
└──────────┘ └──────────┘ └──────────┘
HA 關鍵機制:
- Leader Election:同一時刻只有一個排程器 Active,防止重複調度
- Task Heartbeat:Worker 持有任務期間定期更新心跳,排程器偵測超時後重新分發
- 分散式鎖:Redis
SETNX確保任務不被多個 Worker 同時搶佔 - 冪等性:任務攜帶唯一 ID,Worker 執行前判重
6.2 冪等性保證
冪等性設計(防止任務重複執行):
方案一:唯一執行 ID + 去重資料表
Worker 收到 task_id="order-pay-001"
→ INSERT INTO executions(task_id) ON CONFLICT DO NOTHING
→ 若 INSERT 成功:執行業務邏輯
→ 若 INSERT 失敗:已有記錄,跳過(冪等)
方案二:業務邏輯本身冪等
UPDATE orders SET status='paid' WHERE id=? AND status='pending'
→ 只有 status='pending' 的訂單會被更新,多次執行結果相同
方案三:分散式鎖
SET task:{id}:lock {worker_id} NX EX 30
→ 30 秒內只有一個 Worker 可持有鎖
→ Worker 崩潰後鎖自動釋放
6.3 指數退避重試策略
Exponential Backoff with Jitter:
失敗次數 基礎延遲 Jitter 範圍 實際延遲
──────── ──────── ────────── ────────
第 1 次 1 秒 +0~1 秒 1~2 秒
第 2 次 2 秒 +0~2 秒 2~4 秒
第 3 次 4 秒 +0~4 秒 4~8 秒
第 4 次 8 秒 +0~8 秒 8~16 秒
第 5 次 16 秒 +0~16 秒 16~32 秒
放棄 ─ ─ → Dead Letter Queue
計算公式:
delay = min(maxDelay, baseDelay x 2^retryCount)
delay = delay x (0.5 + random() x 0.5) // 加入隨機性
加入 Jitter 的原因:
避免所有失敗任務在同一時刻重試,造成下游瞬間過載
6.4 Dead Letter Queue
超過最大重試次數的任務進入 死信佇列(Dead Letter Queue, DLQ),由人工或自動化流程介入處理。DLQ 的關鍵設計:
- 保留完整的任務資訊與錯誤日誌
- 提供管理介面供人工審閱與重新提交
- 設定告警閾值——當 DLQ 堆積超過一定數量時自動通知 On-Call 工程師
6.5 生產環境工具比較
| 工具 | 適用場景 | 優點 | 缺點 |
|---|---|---|---|
| Celery | Python,簡單週期任務 | 易上手,社群龐大 | 不支援 DAG |
| Temporal | 複雜工作流,長時間運行 | 可靠性極高,支援暫停/恢復 | 學習曲線高 |
| Apache Airflow | 資料管線,DAG 工作流 | 視覺化強,Python DSL | 重量級,適合離線批次 |
| BullMQ | Node.js,即時任務 | Redis 後端,輕量高效 | 單 Redis 可能成為瓶頸 |
| 自建 Timer Wheel | 超低延遲,百萬級任務 | 最高效能 | 開發成本高 |
7. LeetCode 練習
以下是與任務排程相關的經典題目,建議按順序練習:
題目一:LeetCode 621 — Task Scheduler
| 項目 | 內容 |
|---|---|
| 難度 | Medium |
| 連結 | LeetCode 621 |
| 核心 | 貪心 + 冷卻時間,計算最短完成時間 |
| 提示 | 關鍵在於出現頻率最高的任務決定了最小間隔,其他任務可以填入空閒槽位 |
題目二:LeetCode 1834 — Single-Threaded CPU
| 項目 | 內容 |
|---|---|
| 難度 | Medium |
| 連結 | LeetCode 1834 |
| 核心 | 排序 + Priority Queue 模擬 CPU 排程 |
| 提示 | 按到達時間排序後,用 Min-Heap 按處理時間排序待執行任務 |
題目三:LeetCode 207 — Course Schedule
| 項目 | 內容 |
|---|---|
| 難度 | Medium |
| 連結 | LeetCode 207 |
| 核心 | 拓撲排序偵測 DAG 是否有環(循環依賴) |
| 提示 | 這就是 DAG 工作流排程的環偵測問題——Kahn’s Algorithm 的直接應用 |
題目四:LeetCode 1882 — Process Tasks Using Servers
| 項目 | 內容 |
|---|---|
| 難度 | Medium |
| 連結 | LeetCode 1882 |
| 核心 | 雙 Priority Queue 模擬任務分配到 Worker |
| 提示 | 一個 Heap 管理空閒伺服器,另一個管理忙碌伺服器的完成時間 |
練習建議:先完成 621(經典排程問題)和 207(DAG 環偵測),再挑戰 1834 和 1882。621 的核心在於理解冷卻時間的貪心策略——出現頻率最高的任務決定了最短完成時間,其他任務填入空閒槽位。207 則是本文 DAG Engine 的理論基礎——Kahn’s Algorithm 拓撲排序。建議在面試時先畫出排程流程圖,再開始寫程式碼。
8. 總結
在這篇文章中,我們從需求分析出發,完整設計並實作了一套任務排程器系統:
Priority Queue 排程:使用 Min-Heap 以
executeAt時間戳排序,O(log n) 插入與取出。適合即時任務的精確排序和優先級管理,是排程器的核心引擎。Timer Wheel 延遲排程:將時間軸映射到環形陣列,O(1) 插入與觸發,記憶體固定。多層時間輪更進一步——用極少的槽位支援超大時間範圍,是 Netty、Kafka、Linux Kernel 的共同選擇。
Cron 表達式解析:支援標準 5 欄 Cron 語法(分 時 日 月 週),透過逐欄位匹配算出下次執行時間,完美處理週期性任務排程。
DAG 依賴解析:使用 Kahn’s Algorithm 拓撲排序,自動識別可並行執行的任務群組,並偵測循環依賴。
生產環境架構:Leader Election 確保排程器高可用;冪等性設計防止任務重複執行;指數退避 + Jitter 避免重試雪崩;Dead Letter Queue 兜底處理永久失敗的任務。
任務排程器是 Priority Queue、拓撲排序、分散式協調的綜合應用。理解了這個系統的設計原理後,你不僅能在面試中從容應對 LeetCode 621 等排程題目,還能在實際工作中為後端系統設計可靠的任務調度基礎設施。
在下一篇文章中,我們將探討另一個經典的系統設計問題——
設計短網址系統(URL Shortener),學習如何用雜湊函數和分散式 ID 產生器實現高效的短網址服務,敬請期待!
FAQ
Q: Priority Queue 排程和 Timer Wheel 排程該怎麼選?
兩者的選擇取決於任務規模和精度需求。Priority Queue(Min-Heap) 插入和取出都是 O(log n),沒有精度限制,適合任務數量在十萬級以下且需要精確計時的場景,例如 JDK 的 DelayQueue 就是基於 Priority Queue 實作。Timer Wheel 插入和觸發都是 O(1),記憶體用量固定,適合百萬級以上的大量延遲任務場景,Netty、Linux Kernel 和 Kafka 都採用 Timer Wheel。實務上,兩者經常混合使用:Timer Wheel 負責管理大量延遲任務,當任務即將到期時再丟進 Priority Queue 做最終精確排程。如果你的系統任務量不大(< 10 萬),直接用 Priority Queue 即可。
Q: Cron 表達式和 setTimeout/setInterval 有什麼區別?為什麼需要 Cron?
setTimeout 和 setInterval 是基於相對時間的簡易定時器——「3 秒後執行」或「每 5 秒執行一次」,適合短暫的程式內定時。Cron 則是基於絕對時間的排程語法——「每週一至週五早上 9:00」或「每月 1 號凌晨 2:30」,適合長期運行的週期性任務。核心差異在於:Cron 能表達複雜的時間規則、需要持久化(重啟後不丟失)、每次執行完成後重新計算下次時間而非簡單累加間隔。setInterval 還有累積誤差問題——如果回調耗時超過間隔會導致任務堆積,Cron 排程器則不存在此問題。
Q: 如何確保分散式任務排程器中任務不被重複執行?
分散式環境下防止重複執行需要多層保障。第一層:Leader Election——用 etcd 或 ZooKeeper 選舉唯一的 Active Scheduler,避免多排程器重複分發。第二層:分散式鎖——Worker 接收任務前用 Redis SETNX 搶鎖,搶到鎖才能執行,鎖設合理 TTL 防死鎖。第三層:冪等性——即使任務被執行兩次,結果也一致。常見做法是唯一執行 ID + 去重資料表(INSERT ON CONFLICT DO NOTHING),或讓業務操作本身冪等(UPDATE WHERE status='pending')。第四層:Dead Letter Queue——超過最大重試次數的任務進入死信佇列,由人工介入處理。