設計任務排程器 — Priority Queue、Time Wheel 與 Cron 排程實戰 | 資料結構與演算法

2026/07/26
設計任務排程器 — Priority Queue、Time Wheel 與 Cron 排程實戰 | 資料結構與演算法

任務排程器(Task Scheduler) 是後端系統中不可或缺的基礎設施——從簡單的 Cron Job 定時備份,到複雜的 DAG 依賴工作流,皆依賴排程器精確管理「什麼任務」在「何時」以「什麼優先序」被執行。本文將從需求分析出發,深入設計以 Priority Queue(優先佇列)Timer Wheel(時間輪) 為核心的排程系統,涵蓋 Cron 表達式解析DAG 拓撲排序依賴解析指數退避重試分散式任務調度,搭配 JavaScript/TypeScriptC++ 雙語言完整實作,帶你徹底掌握任務排程器的設計與實戰。

前言

在上一篇文章中,我們實作了

搜尋自動補全(Autocomplete)

,學會如何用 Trie 和 Top-K 排序實現即時的搜尋建議功能。今天我們要解決另一個在實際工作中無處不在的系統設計問題——任務排程器(Task Scheduler)

你每天都在使用排程器,只是可能沒有意識到。手機上的鬧鐘 App、電子郵件的定時發送、電商平台的「下單 30 分鐘未付款自動取消」、每天凌晨自動跑的資料庫備份——這些背後都有一個排程器在默默運作。更複雜的場景包括 CI/CD 流程中的 DAG 依賴工作流、機器學習訓練管線中的多步驟任務編排、以及微服務架構中的分散式任務調度。

這個問題的關鍵挑戰在於:如何同時支援 即時任務延遲任務週期任務,並且在百萬級任務規模下保持高吞吐量和低延遲。暴力的 setTimeout 顯然撐不住——我們需要精心設計的資料結構來管理時間軸上的任務調度。答案是我們在

第 006 篇

中學過的 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 平台排程器為參考:

指標數值說明
即時任務 QPS10,000Priority Queue 排程
延遲任務 QPS5,000Timer Wheel 管理
固定 Cron Job 數量1,000週期性任務
DAG 工作流100/小時多步驟依賴任務
單任務記錄大小~512 B任務 ID + 類型 + Payload + 狀態
Priority Queue 最大容量~100 萬任務記憶體約 512MB
Timer Wheel 記憶體~28 MB360 萬槽位 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 QueueO(log n)O(log n)O(n)O(n)任務量 < 10 萬,需精確計時
單層 Timer WheelO(1)O(1) 均攤O(1)O(slots)百萬級延遲任務
多層 Timer WheelO(1)O(1) 均攤O(1)O(slots x layers)超大範圍 + 高精度
Delay QueueO(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 QueueTimer 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 QueueO(n)n = 活躍任務數
單層 Timer WheelO(slots + n)slots 固定,n = 任務數
多層 Timer WheelO(slots x layers + n)768 槽 x 3 層 = 2304 個槽
DAG 鄰接表O(V + E)拓撲排序的圖結構

效能基準(參考值)

操作Timer WheelPriority Queue
100 萬任務插入~100ms~800ms
適用任務規模百萬級以上十萬級以下
記憶體布局固定陣列(cache-friendly)動態堆積
使用系統Netty, Kafka, LinuxJDK 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 關鍵機制:

  1. Leader Election:同一時刻只有一個排程器 Active,防止重複調度
  2. Task Heartbeat:Worker 持有任務期間定期更新心跳,排程器偵測超時後重新分發
  3. 分散式鎖:Redis SETNX 確保任務不被多個 Worker 同時搶佔
  4. 冪等性:任務攜帶唯一 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 生產環境工具比較

工具適用場景優點缺點
CeleryPython,簡單週期任務易上手,社群龐大不支援 DAG
Temporal複雜工作流,長時間運行可靠性極高,支援暫停/恢復學習曲線高
Apache Airflow資料管線,DAG 工作流視覺化強,Python DSL重量級,適合離線批次
BullMQNode.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. 總結

在這篇文章中,我們從需求分析出發,完整設計並實作了一套任務排程器系統:

  1. Priority Queue 排程:使用 Min-Heap 以 executeAt 時間戳排序,O(log n) 插入與取出。適合即時任務的精確排序和優先級管理,是排程器的核心引擎。

  2. Timer Wheel 延遲排程:將時間軸映射到環形陣列,O(1) 插入與觸發,記憶體固定。多層時間輪更進一步——用極少的槽位支援超大時間範圍,是 Netty、Kafka、Linux Kernel 的共同選擇。

  3. Cron 表達式解析:支援標準 5 欄 Cron 語法(分 時 日 月 週),透過逐欄位匹配算出下次執行時間,完美處理週期性任務排程。

  4. DAG 依賴解析:使用 Kahn’s Algorithm 拓撲排序,自動識別可並行執行的任務群組,並偵測循環依賴。

  5. 生產環境架構: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?

setTimeoutsetInterval 是基於相對時間的簡易定時器——「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——超過最大重試次數的任務進入死信佇列,由人工介入處理。

BenZ Software Developer

熱愛技術的軟體開發者,在這裡分享程式開發經驗與學習筆記。

本週主打

AI 自動化入門包

你每天手動在做的那些煩事,其實 AI 可以自己跑。這份給你 10 個照著做就會的自動化工作流 + 50 個複製即用的提示詞,不用會寫程式。

看看這個產品 →