快轉到主要內容
  1. 教學文章/

Python graphlib 實戰:拓撲排序、依賴圖與工作流程排程

·5 分鐘· loading · loading · ·
Python Graphlib TopologicalSorter DAG Topological-Sort Workflow
每日拍拍
作者
每日拍拍
科學家 X 科技宅宅
目錄
Python 學習 - 本文屬於一個選集。
§ 97: 本文

featured

一. 前言:任務不是排成清單就會自動變合理
#

假設拍拍君要做一份每日報表:

  1. 下載訂單與商品資料。
  2. 清理兩份資料。
  3. 合併資料。
  4. 計算指標。
  5. 產生 HTML。
  6. 發布報表。

把它們依序寫成六行程式,當然可以跑。

但「下載訂單」與「下載商品」其實互不相依,應該能同時進行; 「合併資料」卻必須等兩邊都清理完成,不能只是碰巧排在後面。

真正的規則不是一張線性清單,而是一張依賴圖

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"}  # 方向反了

拍拍君建議變數直接叫 dependenciespredecessors, 不要只叫 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() 的生命週期
#

互動式排序器的流程固定是:

  1. 建立 sorter。
  2. 呼叫 prepare() 鎖定圖並檢查循環。
  3. get_ready() 取出目前沒有未完成前置任務的節點。
  4. 執行這批節點。
  5. 對完成的節點呼叫 done()
  6. 重複直到 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。

小工具不用一開始就變成工作流平台, 但至少可以先學會不讓任務互相等到天荒地老。🔗

延伸閱讀
#

Python 學習 - 本文屬於一個選集。
§ 97: 本文

相關文章

Python fsspec 實戰:統一讀寫本機、S3、HTTP 與資料管線路徑
·7 分鐘· loading · loading
Python Fsspec Filesystem S3 Data-Engineering ETL
Python OpenTelemetry 實戰:Trace、Span 與 FastAPI 觀測流程
·6 分鐘· loading · loading
Python OpenTelemetry Observability FastAPI Tracing Developer-Tools
Python uv build/publish 實戰:從 wheel 到 private package workflow
·11 分鐘· loading · loading
Python Uv Packaging Wheel PyPI Private-Package Developer-Tools
Streamlit + SQLModel 實戰:做一個本機 CRUD 小後台
·9 分鐘· loading · loading
Python Streamlit SQLModel SQLite CRUD Developer-Tools
Python socket 實戰:TCP client/server、timeout 與簡易通訊協定
·8 分鐘· loading · loading
Python Socket TCP Networking Standard-Library Developer-Tools
Python shutil 實戰:檔案複製、搬移、壓縮與安全清理
·7 分鐘· loading · loading
Python Shutil Filesystem Automation Standard-Library Developer-Tools