一. 前言:任務不是排成清單就會自動變合理 #
假設拍拍君要做一份每日報表:
- 下載訂單與商品資料。
- 清理兩份資料。
- 合併資料。
- 計算指標。
- 產生 HTML。
- 發布報表。
把它們依序寫成六行程式,當然可以跑。
但「下載訂單」與「下載商品」其實互不相依,應該能同時進行; 「合併資料」卻必須等兩邊都清理完成,不能只是碰巧排在後面。
真正的規則不是一張線性清單,而是一張依賴圖。
Python 標準庫的 graphlib 提供 TopologicalSorter,
專門處理這類有方向、無循環的圖,也就是 DAG(Directed Acyclic Graph)。
這篇會從最小拓撲排序開始,一路做到:
- 用資料結構描述任務依賴。
- 找出合法執行順序。
- 在啟動工作前偵測循環依賴。
- 分批取得「目前已就緒」的任務。
- 實作可並行、失敗就停止的簡易 workflow runner。
先說清楚:graphlib 不是 cron,也不是背景任務平台。
它負責的是依賴順序,不是「幾點執行」。
二. 安裝:不用裝,標準庫已經準備好了 #
graphlib 從 Python 3.9 起就是標準庫的一部分。
先確認版本:
python --version
只要是 Python 3.9 以上,就可以直接匯入:
from graphlib import TopologicalSorter
如果你用 uv 建一個練習專案:
uv init graphlib-demo
cd graphlib-demo
uv run python -c "from graphlib import TopologicalSorter; print('ready')"
沒有額外 dependency,也不用偷偷把另一個同名套件裝進來。
三. 第一個 TopologicalSorter
#
TopologicalSorter 最方便的輸入格式,是:
{
"任務": {"它依賴的任務", "另一個前置任務"},
}
例如:
from graphlib import TopologicalSorter
dependencies = {
"download": set(),
"clean": {"download"},
"report": {"clean"},
}
sorter = TopologicalSorter(dependencies)
order = list(sorter.static_order())
print(order)
輸出會符合:
["download", "clean", "report"]
這裡最容易寫反。
字典的 value 是「這個 key 的前置任務」,不是「接下來要做什麼」。
所以要寫:
"report": {"clean"}
而不是:
"clean": {"report"} # 方向反了
拍拍君建議變數直接叫 dependencies 或 predecessors,
不要只叫 graph,可以少掉不少腦內翻譯。
四. 實戰資料管線:同一層可以有多個任務 #
來描述前言的報表流程:
from graphlib import TopologicalSorter
dependencies = {
"download_users": set(),
"download_orders": set(),
"clean_users": {"download_users"},
"clean_orders": {"download_orders"},
"join_data": {"clean_users", "clean_orders"},
"calculate_metrics": {"join_data"},
"render_html": {"calculate_metrics"},
"publish": {"render_html"},
}
order = list(TopologicalSorter(dependencies).static_order())
for index, task in enumerate(order, start=1):
print(f"{index}. {task}")
你可能看到:
1. download_users
2. download_orders
3. clean_users
4. clean_orders
5. join_data
6. calculate_metrics
7. render_html
8. publish
但前四項的相對順序不應被當成 API 契約。
只要依賴關係仍然成立,另一個合法順序也完全正確。
因此測試時不要直接斷言完整 list,除非順序真的只有一種。
比較穩定的測試方式,是檢查每個前置任務都在後續任務之前:
def assert_dependencies_respected(
order: list[str],
dependencies: dict[str, set[str]],
) -> None:
positions = {task: index for index, task in enumerate(order)}
for task, predecessors in dependencies.items():
for predecessor in predecessors:
assert positions[predecessor] < positions[task]
這是在測真正的規則,不是在測剛好出現的排列。
五. 循環依賴:不是排序困難,是根本無解 #
如果 A 等 B、B 等 C、C 又等 A:
A -> B -> C -> A
沒有任何任務可以先開始。
TopologicalSorter 會丟出 CycleError:
from graphlib import CycleError, TopologicalSorter
dependencies = {
"extract": {"load"},
"transform": {"extract"},
"load": {"transform"},
}
try:
order = list(TopologicalSorter(dependencies).static_order())
except CycleError as error:
print("依賴圖有循環:", error)
在 workflow 系統裡,這種錯誤應該在啟動時就擋住, 不要等 worker 全部閒著時才猜它們為什麼沒有工作。
可以包成一個設定驗證函式:
from collections.abc import Iterable
from graphlib import CycleError, TopologicalSorter
def validate_dag(
dependencies: dict[str, set[str]],
) -> tuple[str, ...]:
try:
return tuple(TopologicalSorter(dependencies).static_order())
except CycleError as error:
cycle: Iterable[str] = error.args[1]
path = " -> ".join(map(str, cycle))
raise ValueError(f"workflow 存在循環依賴:{path}") from error
對外轉成自己的 domain error,呼叫端就不用理解 graphlib 的例外細節。
六. prepare()、get_ready()、done() 的生命週期
#
互動式排序器的流程固定是:
- 建立 sorter。
- 呼叫
prepare()鎖定圖並檢查循環。 - 用
get_ready()取出目前沒有未完成前置任務的節點。 - 執行這批節點。
- 對完成的節點呼叫
done()。 - 重複直到 sorter 不再 active。
先看不做真正並行的版本:
from graphlib import TopologicalSorter
dependencies = {
"download_users": set(),
"download_orders": set(),
"clean_users": {"download_users"},
"clean_orders": {"download_orders"},
"join_data": {"clean_users", "clean_orders"},
}
sorter = TopologicalSorter(dependencies)
sorter.prepare()
while sorter.is_active():
ready = sorter.get_ready()
print("這一批可以執行:", ready)
for task in ready:
print("執行", task)
sorter.done(task)
概念輸出會像:
這一批可以執行: ('download_users', 'download_orders')
這一批可以執行: ('clean_users', 'clean_orders')
這一批可以執行: ('join_data',)
done() 非常重要。
排序器不知道你的函式何時結束,必須由 runner 明確回報。
如果忘記呼叫,後續節點永遠不會變成 ready。
七. 用 ThreadPoolExecutor 執行可並行任務
#
對網路請求、檔案等待這類 I/O 工作,可以用 thread pool:
from collections.abc import Callable
from concurrent.futures import (
FIRST_COMPLETED,
Future,
ThreadPoolExecutor,
wait,
)
from graphlib import TopologicalSorter
Task = Callable[[], None]
def run_workflow(
tasks: dict[str, Task],
dependencies: dict[str, set[str]],
*,
max_workers: int = 4,
) -> None:
sorter = TopologicalSorter(dependencies)
sorter.prepare()
running: dict[Future[None], str] = {}
with ThreadPoolExecutor(max_workers=max_workers) as pool:
while sorter.is_active() or running:
for name in sorter.get_ready():
future = pool.submit(tasks[name])
running[future] = name
if not running:
raise RuntimeError("workflow 尚未完成,卻沒有可執行任務")
finished, _ = wait(running, return_when=FIRST_COMPLETED)
for future in finished:
name = running.pop(future)
future.result()
sorter.done(name)
這個 runner 有幾個重要行為:
- 每次把所有 ready task 送進 pool。
- 至少等一個 task 完成,再解鎖後續依賴。
future.result()會重新拋出 worker 裡的例外。- 只有成功完成的 task 才會交給
sorter.done()。
它已經能處理小型內部工具,但還不是完整工作流平台。
例如:失敗重試、持久化、跨機器 queue、取消傳播與 observability, 都需要另外設計。
八. graphlib、APScheduler、asyncio 怎麼選?
#
它們解的問題不同:
| 工具 | 核心問題 | 典型用途 |
|---|---|---|
graphlib |
哪些任務必須先完成? | 建置、資料管線、dependency plan |
| APScheduler | 任務應該在什麼時間觸發? | cron、interval、指定時間工作 |
asyncio / AnyIO |
如何管理非同步 I/O 與任務生命週期? | API client、socket、並發 I/O |
| Celery / workflow platform | 如何跨 process 或機器排隊、重試與保存狀態? | production background jobs |
它們也可以組合。
例如 APScheduler 每天早上觸發一次報表;
報表內部用 graphlib 決定步驟依賴;
每個下載 task 再用 AnyIO 並發呼叫 API。
工具分層之後,每一層都只處理自己擅長的規則。
結語 #
graphlib.TopologicalSorter 的 API 很小,卻能把一類常見問題說得很清楚:
- 用
task -> predecessors建立依賴圖。 - 用
static_order()取得合法線性順序。 - 用
CycleError在開工前抓出無解設定。 - 用
prepare()、get_ready()、done()管理可並行的就緒批次。 - 只有成功完成的任務,才應該標記為 done。
當你的腳本開始充滿「先跑這個、如果完成再跑那個」,
先別急著堆更多 if。
把依賴畫成 DAG,通常會立刻看出哪些規則是真的, 哪些順序只是歷史包袱。
拍拍君的建議是:先用 static_order() 驗證模型,
確定依賴方向正確,再升級到 parallel runner。
小工具不用一開始就變成工作流平台, 但至少可以先學會不讓任務互相等到天荒地老。🔗