AI Agent 設計:Workflow 工作流編排——Cyclic Graph 循環圖與 Checkpoint 斷點續跑實作
這個章節我們要來討論如何用「有狀態的工作流圖」把多步 AI 任務組成穩定、可控、可恢復的流程,原因在於自由 Agent 靈活但難預測;所以固定流程用 Workflow,分支/循環/終止都顯式可控
Cyclic Graph(循環圖)
代表的是圖中至少存在一條路徑,沿著邊前進後,可以回到曾經經過的節點。
例如:
A → B → C
↑ │
└───┘
循環如下:
answer → quality → answer
與 DAG 的主要差異是:
| 圖類型 | 可回到先前節點 | 常見用途 |
|---|---|---|
| DAG | 不可以 | 固定順序、一次性資料管線 |
| Cyclic Graph | 可以 | 重試、迭代、反思、審核修正 |
這邊是簡單的workflow架構程式
# ─────────────────────────────────────────────────────────────────────────────
# 迷你工作流引擎:節點是 (state)->state 的函式;邊可帶條件函式做路由
# ─────────────────────────────────────────────────────────────────────────────
class Workflow:
def __init__(self):
self.nodes = {} # name -> fn(state)->state
self.edges = {} # name -> [(condition_fn, next_name)]
self.entry = None
def node(self, name, fn):
self.nodes[name] = fn
return self
def edge(self, src, dst, cond=None):
self.edges.setdefault(src, []).append((cond or (lambda s: True), dst))
return self
def set_entry(self, name):
self.entry = name
return self
def _next(self, name, state):
for cond, dst in self.edges.get(name, []):
if cond(state):
return dst
return None
def run(self, state, start=None, on_step=None, max_steps=20):
cur = start or self.entry
steps = 0
while cur and cur != "END" and steps < max_steps:
state = self.nodes[cur](dict(state))
state["_last_node"] = cur
if on_step:
on_step(cur, state)
cur = self._next(cur, state)
steps += 1
state["_done"] = (cur == "END" or cur is None)
return state
def cmd_graph(_args):
wf = build_workflow()
print(f"入口節點:{wf.entry}\n\n節點與出邊(條件分支 / 循環):")
for src, outs in wf.edges.items():
for cond, dst in outs:
tag = " [有條件]" if cond.__code__.co_argcount and "lambda" in cond.__qualname__ else ""
print(f" {src:10s} ──→ {dst}{tag}")
print("\n要點:工作流把『流程結構』顯式畫出來,分支/循環/終止都可控可測——"
"這是它相對自由 Agent 的優勢。")
然後展開 workflow 流程結構
python workflow_lab.py graph
結果如下:
入口節點:classify
節點與出邊(條件分支 / 循環):
classify ──→ escalate [有條件]
classify ──→ answer [有條件]
escalate ──→ END [有條件]
answer ──→ quality [有條件]
quality ──→ END [有條件]
quality ──→ answer [有條件]

在 answer 寫死是回答營業時間,然後 quality 做了簡單的回復内文長度檢查,如下:
# ─────────────────────────────────────────────────────────────────────────────
# 範例業務:客服工單分流工作流
# classify → (抱怨? → escalate) / (問題? → answer→quality_check→(過關?END / 否則重answer))
# ─────────────────────────────────────────────────────────────────────────────
def n_classify(s):
text = s["input"]
s["intent"] = "complaint" if any(w in text for w in ("退費", "壞", "爛", "投訴")) else "question"
return s
def n_escalate(s):
s["reply"] = "已將您的問題升級給專人處理,將於 24 小時內回覆。"
return s
def n_answer(s):
s["attempt"] = s.get("attempt", 0) + 1
# 模擬第一次答得太短、第二次補足(觸發 quality 循環)
s["reply"] = ("營業時間 9-18 點。" if s["attempt"] >= 2 else "9-18")
return s
def n_quality(s):
s["quality_ok"] = len(s.get("reply", "")) >= 8 # 太短視為不合格 → 回去重答
return s
執行測試:
python workflow_lab.py run --input "我要退費,你們東西是壞的"
結果如下:
輸入:我要退費,你們東西是壞的
▶ 節點 classify → {'intent': 'complaint'}
▶ 節點 escalate → {'intent': 'complaint', 'reply': '已將您的問題升級給專人處理,將於 24 小時內回覆。'}
最終回覆:已將您的問題升級給專人處理,將於 24 小時內回覆。
完成:True(共經過分流,answer 嘗試 0 次)
執行測試2:
python workflow_lab.py run --input "請問營業時間?"
結果如下:
輸入:請問營業時間?
▶ 節點 classify → {'intent': 'question'}
▶ 節點 answer → {'intent': 'question', 'attempt': 1, 'reply': '9-18'}
▶ 節點 quality → {'intent': 'question', 'attempt': 1, 'reply': '9-18', 'quality_ok': False}
▶ 節點 answer → {'intent': 'question', 'attempt': 2, 'reply': '營業時間 9-18 點。', 'quality_ok': False}
▶ 節點 quality → {'intent': 'question', 'attempt': 2, 'reply': '營業時間 9-18 點。', 'quality_ok': True}
最終回覆:營業時間 9-18 點。
完成:True(共經過分流,answer 嘗試 2 次)
在下面嘗試解釋如何在 workflow 流程執行中出現 crash 後再次執行流程,需要將執行記錄寫入到 checkpoint,這樣在下次執行時可以知道中斷在哪從而繼續
CKPT = Path(__file__).resolve().parent / "checkpoint.json"
def cmd_run(args):
wf = build_workflow()
print(f"輸入:{args.input}\n")
def show(node, state):
extra = {k: v for k, v in state.items() if not k.startswith("_") and k != "input"}
print(f" ▶ 節點 {node:10s} → {extra}")
final = wf.run({"input": args.input}, on_step=show)
print(f"\n最終回覆:{final.get('reply')}")
print(f"完成:{final['_done']}(共經過分流,answer 嘗試 {final.get('attempt', 0)} 次)")
def cmd_resume(_args):
wf = build_workflow()
print("情境:工作流跑到 classify 後『崩潰』,狀態已存檔,之後從 checkpoint 續跑\n")
# 第一階段:只跑 entry 節點就「中斷」,把狀態存檔
state = wf.nodes["classify"](dict({"input": "請問營業時間?"}))
state["_last_node"] = "classify"
CKPT.write_text(json.dumps(state, ensure_ascii=False), encoding="utf-8")
print(f" [存檔] 已完成 classify,state={ {k:v for k,v in state.items() if not k.startswith('_')} }")
print(f" [崩潰] 程式中斷…\n")
# 第二階段:重新啟動,讀回 checkpoint,從中斷處的下一個節點續跑
restored = json.loads(CKPT.read_text(encoding="utf-8"))
resume_from = wf._next(restored["_last_node"], restored)
print(f" [恢復] 讀回 checkpoint,從節點 `{resume_from}` 續跑")
final = wf.run(restored, start=resume_from,
on_step=lambda n, s: print(f" ▶ {n} → reply={s.get('reply')!r}"))
print(f"\n最終回覆:{final.get('reply')}")
CKPT.unlink(missing_ok=True)
print("要點:持久化 checkpoint 讓長流程可中斷、可恢復、可人工介入後再續——生產長任務必備。")
測試:
python workflow_lab.py resume
執行結果:
情境:工作流跑到 classify 後『崩潰』,狀態已存檔,之後從 checkpoint 續跑
[存檔] 已完成 classify,state={'input': '請問營業時間?', 'intent': 'question'}
[崩潰] 程式中斷…
[恢復] 讀回 checkpoint,從節點 `answer` 續跑
▶ answer → reply='9-18'
▶ quality → reply='9-18'
▶ answer → reply='營業時間 9-18 點。'
▶ quality → reply='營業時間 9-18 點。'
最終回覆:營業時間 9-18 點。
要點:持久化 checkpoint 讓長流程可中斷、可恢復、可人工介入後再續——生產長任務必備。
可以看到 checkpoint 内容如下:
