昨天我們看了四種 Workflow 模式。今天要處理實務上真正會遇到的問題:
這些問題的答案,就是今天的主題:控制流與狀態管理。
最土法煉鋼的做法是用一堆變數:
def daily_brief():
now = get_time()
weather = get_weather(city)
todos = get_todos()
analysis = analyze(weather, todos)
summary = summarize(now, weather, todos, analysis)
return summary
步驟少的時候還好,但參數會越傳越多,而且中間出錯的話什麼都沒留下。
比較好的做法是用一個狀態物件貫穿整個流程:
"""workflows/state.py"""
import json
import time
from dataclasses import dataclass, field
from datetime import datetime
from pathlib import Path
from typing import Any, Optional
@dataclass
class WorkflowState:
"""一次 Workflow 執行的完整狀態。"""
# 輸入
user_input: str
started_at: float = field(default_factory=time.time)
# 各步驟的產出
data: dict = field(default_factory=dict)
# 執行紀錄
completed_steps: list = field(default_factory=list)
failed_steps: list = field(default_factory=list)
warnings: list = field(default_factory=list)
# 成本
input_tokens: int = 0
output_tokens: int = 0
# 結果
result: Optional[str] = None
error: Optional[str] = None
# --- 操作 ---
def set(self, key: str, value: Any):
self.data[key] = value
return self
def get(self, key: str, default=None):
return self.data.get(key, default)
def has(self, *keys) -> bool:
"""檢查是否所有指定的 key 都有值。"""
return all(self.data.get(k) is not None for k in keys)
def mark_done(self, step: str):
self.completed_steps.append(step)
def mark_failed(self, step: str, reason: str):
self.failed_steps.append({"step": step, "reason": reason})
def warn(self, message: str):
self.warnings.append(message)
def add_usage(self, response):
self.input_tokens += response.usage.input_tokens
self.output_tokens += response.usage.output_tokens
@property
def elapsed(self) -> float:
return time.time() - self.started_at
@property
def is_healthy(self) -> bool:
"""沒有任何步驟失敗。"""
return not self.failed_steps
def summary(self) -> str:
return (
f"完成 {len(self.completed_steps)} 步"
f"(失敗 {len(self.failed_steps)},警告 {len(self.warnings)}),"
f"耗時 {self.elapsed:.2f}s,"
f"用量 {self.input_tokens}+{self.output_tokens} tokens"
)
def to_dict(self) -> dict:
return {
"user_input": self.user_input,
"data": self.data,
"completed_steps": self.completed_steps,
"failed_steps": self.failed_steps,
"warnings": self.warnings,
"usage": {"input": self.input_tokens, "output": self.output_tokens},
"elapsed": round(self.elapsed, 2),
"result": self.result,
"error": self.error,
}
有了狀態物件,每個步驟的簽名就統一了:
def step_name(state: WorkflowState) -> WorkflowState:
...
return state
**每一步讀狀態、做事、寫狀態。**這讓步驟可以自由組合、重新排序、跳過。
"""workflows/step.py"""
import logging
from typing import Callable
logger = logging.getLogger(__name__)
class Step:
"""Workflow 中的一個步驟。"""
def __init__(
self,
name: str,
fn: Callable,
required: bool = True,
requires: tuple = (),
retries: int = 0,
):
"""
Args:
name: 步驟名稱
fn: 實際執行的函式,簽名為 fn(state) -> state
required: 失敗時是否中止整個流程
requires: 執行前必須已經存在的狀態 key
retries: 失敗時重試幾次
"""
self.name = name
self.fn = fn
self.required = required
self.requires = requires
self.retries = retries
def execute(self, state):
# 前置條件檢查
missing = [k for k in self.requires if not state.has(k)]
if missing:
reason = f"缺少必要資料:{', '.join(missing)}"
logger.warning("跳過步驟 %s:%s", self.name, reason)
state.mark_failed(self.name, reason)
if self.required:
raise WorkflowError(f"步驟「{self.name}」無法執行:{reason}")
return state
# 執行(含重試)
last_error = None
for attempt in range(self.retries + 1):
try:
logger.info("執行步驟:%s", self.name)
state = self.fn(state)
state.mark_done(self.name)
return state
except Exception as e:
last_error = e
if attempt < self.retries:
logger.warning(
"步驟 %s 失敗(第 %d 次):%s,重試中",
self.name, attempt + 1, e,
)
# 全部重試都失敗
reason = f"{type(last_error).__name__}: {last_error}"
state.mark_failed(self.name, reason)
if self.required:
logger.error("必要步驟 %s 失敗,中止流程", self.name)
raise WorkflowError(f"步驟「{self.name}」失敗:{last_error}") from last_error
logger.warning("選用步驟 %s 失敗,繼續執行", self.name)
state.warn(f"{self.name} 沒有完成:{last_error}")
return state
class WorkflowError(Exception):
"""Workflow 執行失敗。"""
兩個關鍵設計:
required 區分必要與選用
**這是 Agent 韌性的關鍵。**一個好的生活助理,在天氣 API 掛掉時應該說「今天天氣查不到,但你有三件事要處理」,而不是整個罷工。
requires 宣告前置條件
明確寫出「這一步需要哪些資料」,比在函式裡寫 if weather is None: return 清楚,而且可以在執行前就檢查。
"""workflows/engine.py"""
import logging
logger = logging.getLogger(__name__)
class Workflow:
"""一連串步驟組成的流程。"""
def __init__(self, name, steps=None):
self.name = name
self.steps = list(steps or [])
def add(self, step):
self.steps.append(step)
return self # 支援串接:wf.add(a).add(b).add(c)
def run(self, user_input, tracer=None):
state = WorkflowState(user_input=user_input)
logger.info("=== 開始執行 Workflow:%s ===", self.name)
try:
for step in self.steps:
state = step.execute(state)
except WorkflowError as e:
state.error = str(e)
logger.error("Workflow %s 中止:%s", self.name, e)
except Exception as e:
state.error = f"未預期的錯誤:{e}"
logger.exception("Workflow %s 發生未預期錯誤", self.name)
logger.info("=== %s 結束:%s ===", self.name, state.summary())
if tracer:
tracer._write("workflow_done", **state.to_dict())
return state
把前面所有東西組起來:
"""workflows/daily_brief.py"""
from datetime import datetime
import config
import prompts
from analysis.todo_insight import analyze_todos, make_todo_summary
from analysis.weather_insight import analyze_weather, make_advice
from tools.weather import get_weather
from todo_store import load_todos
from workflows.engine import Step, Workflow
from workflows.state import WorkflowState
# === 各步驟 ===
def step_get_time(state):
now = datetime.now()
weekdays = ["星期一", "星期二", "星期三", "星期四", "星期五", "星期六", "星期日"]
state.set("now", now)
state.set("time_text", f"{now.strftime('%Y年%m月%d日 %H:%M')}({weekdays[now.weekday()]})")
return state
def step_get_weather(state):
data = get_weather(config.DEFAULT_CITY, config.CWA_API_KEY)
insight = analyze_weather(data)
state.set("weather_raw", data)
state.set("weather_insight", insight)
state.set("weather_text", make_advice(insight))
return state
def step_get_todos(state):
todos = load_todos()
insight = analyze_todos(todos)
state.set("todos", todos)
state.set("todo_insight", insight)
state.set("todo_text", make_todo_summary(insight))
return state
def make_step_compose(client):
"""產生「合成簡報」這個步驟(需要注入 client)。"""
def step_compose(state):
prompt = prompts.DAILY_BRIEF.substitute(
current_time=state.get("time_text"),
weather=state.get("weather_text", "(今天沒有取得天氣資訊)"),
todos=state.get("todo_text", "(今天沒有取得待辦資訊)"),
)
response = client.messages.create(
model="claude-opus-5",
max_tokens=1024,
system=prompts.ASSISTANT_SYSTEM,
messages=[{"role": "user", "content": prompt}],
)
state.add_usage(response)
text = "\n".join(b.text for b in response.content if b.type == "text").strip()
state.result = text
# 如果有步驟失敗,附上說明
if state.warnings:
notes = "\n".join(f"(註:{w})" for w in state.warnings)
state.result = f"{text}\n\n{notes}"
return state
return step_compose
# === 組裝 ===
def build_daily_brief_workflow(client):
return Workflow("今日簡報", [
Step("取得時間", step_get_time, required=True),
Step("查詢天氣", step_get_weather, required=False, retries=1),
Step("讀取待辦", step_get_todos, required=False),
Step("合成簡報", make_step_compose(client),
required=True, requires=("time_text",)),
])
用起來:
import anthropic
import config
from logger_setup import setup_logging
from workflows.daily_brief import build_daily_brief_workflow
setup_logging()
client = anthropic.Anthropic(api_key=config.ANTHROPIC_API_KEY)
workflow = build_daily_brief_workflow(client)
state = workflow.run("早安")
print(state.result)
print()
print(state.summary())
正常情況的輸出:
早安!今天 10 月 8 日星期三,台北白天有短暫陣雨,降雨機率 60%,
出門記得帶傘,氣溫 25 到 30 度,短袖加件薄外套剛好。
今天最該處理的是「交鐵人賽文章」,今天就到期了,而且是高優先度。
繳電費還有三天,可以晚點再說。
下午雨勢可能比較明顯,如果要出門辦事建議挑上午。
完成 4 步(失敗 0,警告 0),耗時 3.21s,用量 892+186 tokens
天氣 API 掛掉的情況(我把 key 改錯試了一下):
早安!今天 10 月 8 日星期三上午 8 點 15 分。
今天有兩件事要處理,最急的是「交鐵人賽文章」,今天就到期。
繳電費還有三天。
天氣資訊目前拿不到,出門前建議自己確認一下。
(註:查詢天氣 沒有完成:氣象署授權碼無效,請檢查 CWA_API_KEY)
完成 3 步(失敗 1,警告 1),耗時 2.84s,用量 654+142 tokens
**流程沒有崩潰,而且誠實說明了缺什麼。**這就是 required=False 的價值。
有時候某些步驟只在特定條件下才該跑:
class ConditionalStep(Step):
"""只在條件成立時執行的步驟。"""
def __init__(self, name, fn, condition, **kwargs):
super().__init__(name, fn, **kwargs)
self.condition = condition
def execute(self, state):
if not self.condition(state):
logger.info("條件不成立,跳過步驟:%s", self.name)
return state
return super().execute(state)
用法:
def step_rain_reminder(state):
"""下雨天才產生的額外提醒。"""
insight = state.get("weather_insight")
state.set("extra_reminder",
f"降雨機率 {insight['max_rain']}%,記得帶傘,"
f"機車族建議提早出門。")
return state
workflow.add(ConditionalStep(
"雨天提醒",
step_rain_reminder,
condition=lambda s: s.get("weather_insight", {}).get("will_rain"),
required=False,
))
條件判斷用程式,不要用 AI。「降雨機率超過 40% 就提醒」這種規則是確定的,寫死比讓模型判斷可靠得多(Day 16 講過的原則)。
對清單裡的每個項目做一樣的事:
class ForEachStep(Step):
"""對某個 list 的每個元素執行操作。"""
def __init__(self, name, fn, items_key, result_key, **kwargs):
super().__init__(name, fn, **kwargs)
self.items_key = items_key
self.result_key = result_key
def execute(self, state):
items = state.get(self.items_key) or []
results = []
for i, item in enumerate(items):
try:
results.append(self.fn(item, state))
except Exception as e:
logger.warning("處理第 %d 項失敗:%s", i, e)
state.warn(f"{self.name} 第 {i + 1} 項失敗:{e}")
results.append(None)
state.set(self.result_key, results)
state.mark_done(self.name)
return state
每一項的失敗是獨立的——第三項壞掉不影響第四項。這在批次處理時很重要。
如果流程很長(或很貴),中途失敗要能從斷點接續:
import json
from pathlib import Path
CHECKPOINT_DIR = Path("data/checkpoints")
class ResumableWorkflow(Workflow):
"""支援斷點續跑的 Workflow。"""
def run(self, user_input, session_id, tracer=None):
CHECKPOINT_DIR.mkdir(parents=True, exist_ok=True)
path = CHECKPOINT_DIR / f"{session_id}.json"
# 嘗試載入之前的進度
state = WorkflowState(user_input=user_input)
done = set()
if path.exists():
try:
saved = json.loads(path.read_text(encoding="utf-8"))
state.data = saved["data"]
state.completed_steps = saved["completed_steps"]
done = set(saved["completed_steps"])
logger.info("從斷點繼續,已完成 %d 步", len(done))
except (json.JSONDecodeError, KeyError):
logger.warning("checkpoint 檔案損壞,從頭開始")
try:
for step in self.steps:
if step.name in done:
logger.info("跳過已完成的步驟:%s", step.name)
continue
state = step.execute(state)
# 每完成一步就存檔
path.write_text(
json.dumps(state.to_dict(), ensure_ascii=False,
indent=2, default=str),
encoding="utf-8",
)
except WorkflowError as e:
state.error = str(e)
logger.error("流程中斷,進度已儲存至 %s", path)
return state
# 成功完成,清掉 checkpoint
path.unlink(missing_ok=True)
return state
default=str 是為了處理 datetime 這類 JSON 不認識的物件(Day 10 提過的坑)。
日常對話的 Workflow 用不到這個,但如果是「每天凌晨批次處理一百個使用者的簡報」這種任務,斷點續跑能省下很多重跑的成本。
寫了這麼多,歸納幾個原則:
# ❌ 一步做太多
def step_get_everything(state):
state.set("weather", get_weather(...))
state.set("todos", load_todos())
state.set("calendar", get_events())
return state
# ✅ 拆開,各自可以設定 required / retries
Step("查天氣", step_weather, required=False)
Step("讀待辦", step_todos, required=True)
Step("讀行事曆", step_calendar, required=False)
問自己:這一步失敗,整件事還有意義嗎?
有意義 → required=False
沒意義 → required=True
不要默默跳過。state.warn() 記下來,最後告訴使用者。誠實的部分結果,比假裝完整的答案好。
能寫成 if x > 40: 的,就不要問模型。
每一次呼叫都是成本和延遲。能合併的合併,能用程式算的別問模型。上面的 daily_brief 四個步驟裡只有一步用 AI,其他都是純程式。
看完這兩天的內容,可以歸納出一個很清楚的對照:
| Workflow | Agent | |
|---|---|---|
| 控制流 | 你寫的 if / for / Step | 模型的 stop_reason 迴圈 |
| 狀態 | WorkflowState 物件 |
messages 陣列 |
| 錯誤處理 | required 標記、重試設定 |
錯誤訊息回傳給模型,它自己決定 |
| 成本 | 固定、可預估 | 浮動、要設上限 |
| 除錯 | 看 state、看 log | 看 trace |
| 擴充 | 加一個 Step | 加一個 Tool |
有趣的是兩者的結構其實同構:
Workflow:state → step → state → step → ... → result
Agent: messages → model → tool → messages → model → ... → answer
差別只在於「下一步是誰決定的」。
而這正是明天的主題——那條界線到底在哪裡?什麼時候該讓 AI 自己決定?
WorkflowState 物件貫穿流程,每一步讀狀態、做事、寫狀態Step 加上 required / requires / retries,把錯誤處理標準化ConditionalStep 做條件分支,ForEachStep 做迴圈,每項失敗互相獨立Checkpoint 讓長流程可以斷點續跑明天要談整個系列標題裡的那個詞:自主決策。什麼時候該放手,什麼時候該抓緊。