這個章節我們要來討論如何用「有狀態的工作流圖」把多步 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 [有條件]

17089-b1stc15jvtr.png

在 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 内容如下:
53853-suo50u4tbie.png

無標籤

關注作者:

新增評論