iT邦幫忙

2026 iThome 鐵人賽

DAY 24
0
AI Engineering

從 Prompt 到自主決策:用 Python × Agentic Workflow 實作生活助理系列 第 24 篇

讓流程能轉彎:Workflow 的控制流與狀態管理

  • 分享至 

  • xImage
  •  

昨天我們看了四種 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 清楚,而且可以在執行前就檢查。


三、組裝成 Workflow

"""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

四、實作:今日簡報 Workflow

把前面所有東西組起來:

"""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

每一項的失敗是獨立的——第三項壞掉不影響第四項。這在批次處理時很重要。


七、Checkpoint:讓流程可以續跑

如果流程很長(或很貴),中途失敗要能從斷點接續:

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 用不到這個,但如果是「每天凌晨批次處理一百個使用者的簡報」這種任務,斷點續跑能省下很多重跑的成本。


八、控制流的設計原則

寫了這麼多,歸納幾個原則:

1. 每一步只做一件事

# ❌ 一步做太多
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)

2. 明確區分必要與選用

問自己:這一步失敗,整件事還有意義嗎?

有意義 → required=False
沒意義 → required=True

3. 失敗要留下痕跡

不要默默跳過。state.warn() 記下來,最後告訴使用者。誠實的部分結果,比假裝完整的答案好。

4. 條件判斷用程式,不用 AI

能寫成 if x > 40: 的,就不要問模型。

5. AI 呼叫要集中

每一次呼叫都是成本和延遲。能合併的合併,能用程式算的別問模型。上面的 daily_brief 四個步驟裡只有一步用 AI,其他都是純程式。


九、Workflow 與 Agent 的界線

看完這兩天的內容,可以歸納出一個很清楚的對照:

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 讓長流程可以斷點續跑
  • 每一步只做一件事;條件判斷用程式;AI 呼叫盡量集中
  • Workflow 和 Agent 結構同構,差別只在「誰決定下一步」

明天要談整個系列標題裡的那個詞:自主決策。什麼時候該放手,什麼時候該抓緊。


上一篇
從單次呼叫到流程:什麼是 Workflow
系列文
從 Prompt 到自主決策:用 Python × Agentic Workflow 實作生活助理 共 24 篇
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言