Files
2026_DesignAI/officefile/md/supplements/01-foundations/01.2-state-and-state-machines.md
T
pengxiao a90f7adfa1 refactor(officefile): 按 md/latex/word 三层结构重组文档目录
将 Markdown 源文件移入 md/,LaTeX 工作目录保留在 latex/,
Word 导出移入 word/;删除临时脚本、调试截图和空 stub。

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
2026-05-29 14:25:21 +08:00

22 KiB
Raw Blame History

01.2 状态与状态机

核心问题

在AI系统中,"状态"到底是什么? 为什么状态管理是Agent系统的核心? LangGraph是如何用状态机设计工作流的?


概念讲解

什么是状态

状态是系统在某一时刻的快照,包含所有影响未来行为的信息:

系统状态 = 所有相关的变量值

例如:生态分析工作流的状态
{
    "input_data": {...},           # 输入数据
    "current_step": "buffer",      # 当前步骤
    "intermediate_results": {...}, # 中间结果
    "user_preferences": {...},     # 用户偏好
    "error_count": 0,              # 错误计数
    "checkpoint_reached": False    # 检查点状态
}

状态的类型

类型 说明 例子
静态状态 初始输入,不变化 输入文件路径、参数
动态状态 运行中变化 当前步骤、累积结果
控制状态 影响流程走向 分支条件、错误标志
会话状态 跨请求持久化 用户偏好、历史记录

为什么状态管理很重要

1. 断点续传

# 没有状态管理:出错后必须从头开始
def analysis_without_state():
    step1()
    step2()  # 如果这里出错
    step3()  # 这些都要重做

# 有状态管理:可以从断点继续
class AnalysisWithState:
    def __init__(self):
        self.state = {"current_step": 0}

    def run(self):
        if self.state["current_step"] < 1:
            step1()
            self.state["current_step"] = 1

        if self.state["current_step"] < 2:
            try:
                step2()
                self.state["current_step"] = 2
            except Exception:
                # 保存状态,下次可以从这里继续
                save_state(self.state)
                raise

        if self.state["current_step"] < 3:
            step3()

2. 人机协同

# HITL需要状态来知道在哪里需要人类介入
class HITLWorkflow:
    def __init__(self):
        self.state = {
            "step": "identify_sources",
            "pending_review": True,
            "sources": None,
            "human_feedback": None
        }

    def next_action(self):
        if self.state["pending_review"]:
            return "request_human_review"
        elif self.state["human_feedback"]:
            return "incorporate_feedback"
        else:
            return "proceed_to_next_step"

3. 调试和可解释性

# 状态历史记录了整个决策过程
class StatefulAgent:
    def __init__(self):
        self.state_history = []

    def decide(self, context):
        # 记录状态
        self.state_history.append({
            "timestamp": now(),
            "state": self.state.copy(),
            "context": context,
            "decision": None
        })

        # 做决策
        decision = self._make_decision(context)
        self.state_history[-1]["decision"] = decision

        return decision

    def explain(self):
        """回溯决策过程"""
        return self.state_history

状态机

状态机是描述系统状态转换的模型:

    ┌─────────┐
    │   初始   │
    │  state  │
    └────┬────┘
         │ event: start
         ↓
    ┌─────────┐
    │ 加载数据 │
    └────┬────┘
         │ success
         ↓
    ┌─────────┐     error      ┌─────────┐
    │ 分析处理 │ ─────────────→│  错误   │
    └────┬────┘                 └─────────┘
         │ success                      │ retry
         ↓                              │
    ┌─────────┐                        │
    │  人类   │                        │
    │  审查   │                        │
    └────┬────┘                        │
         │ approve                     │
         ↓                             │
    ┌─────────┐                        │
    │  完成   │←───────────────────────┘
    └─────────┘

状态机的要素

  1. 状态 (State):系统可能处于的情况
  2. 事件 (Event):触发状态转换的条件
  3. 转换 (Transition):从一个状态到另一个状态
  4. 动作 (Action):状态转换时执行的操作

设计原理

LangGraph的状态设计哲学

LangGraph是构建有状态Agent的框架,其核心思想:

from typing import TypedDict

# 定义状态类型
class AnalysisState(TypedDict):
    """生态网络分析状态"""

    # 输入数据
    input_path: str
    parameters: dict

    # 处理过程
    current_step: str
    intermediate_results: dict

    # 人机交互
    review_requested: bool
    human_feedback: str

    # 输出
    final_result: dict
    errors: list

# 状态图定义
workflow = StateGraph(AnalysisState)

# 添加节点(处理步骤)
workflow.add_node("load_data", load_data_node)
workflow.add_node("identify_sources", identify_sources_node)
workflow.add_node("human_review", human_review_node)
workflow.add_node("extract_corridors", extract_corridors_node)

# 添加边(状态转换)
workflow.add_edge("load_data", "identify_sources")
workflow.add_conditional_edge(
    "identify_sources",
    should_review,  # 条件函数
    {
        "review": "human_review",
        "continue": "extract_corridors"
    }
)

# 编译为可执行图
app = workflow.compile()

核心概念

  1. 状态即消息:状态在节点间传递
  2. 图即流程:有向图描述工作流
  3. 条件分支:基于状态的动态路由

工作流状态机实现

from enum import Enum
from typing import Dict, Any, Callable, Optional
from dataclasses import dataclass, field

class WorkflowState(Enum):
    """工作流状态枚举"""
    IDLE = "idle"
    LOADING = "loading"
    PROCESSING = "processing"
    REVIEWING = "reviewing"
    COMPLETED = "completed"
    ERROR = "error"

@dataclass
class WorkflowContext:
    """工作流上下文(状态数据)"""
    data: Dict[str, Any] = field(default_factory=dict)
    current_step: int = 0
    errors: list = field(default_factory=list)
    metadata: Dict[str, Any] = field(default_factory=dict)

class StateMachine:
    """通用状态机"""

    def __init__(self, initial_state: WorkflowState):
        self.state = initial_state
        self.context = WorkflowContext()
        self.transitions: Dict[WorkflowState, Dict[str, WorkflowState]] = {}
        self.actions: Dict[tuple[WorkflowState, WorkflowState], Callable] = {}

    def add_transition(self,
                       from_state: WorkflowState,
                       event: str,
                       to_state: WorkflowState,
                       action: Callable = None):
        """添加状态转换"""
        if from_state not in self.transitions:
            self.transitions[from_state] = {}
        self.transitions[from_state][event] = to_state

        if action:
            self.actions[(from_state, to_state)] = action

    def trigger(self, event: str, **kwargs) -> bool:
        """触发事件"""
        if self.state not in self.transitions:
            raise ValueError(f"没有从状态 {self.state} 的转换")

        if event not in self.transitions[self.state]:
            print(f"事件 {event} 在状态 {self.state} 下无效")
            return False

        # 获取目标状态
        new_state = self.transitions[self.state][event]
        old_state = self.state

        # 执行转换动作
        action = self.actions.get((old_state, new_state))
        if action:
            result = action(self.context, **kwargs)
            if result is False:  # 动作失败,不转换
                return False

        # 更新状态
        self.state = new_state
        print(f"状态转换: {old_state}{new_state}")
        return True

# 生态分析工作流状态机
class EcologicalAnalysisWorkflow:
    """生态网络分析工作流"""

    def __init__(self):
        # 创建状态机
        self.sm = StateMachine(WorkflowState.IDLE)

        # 定义转换
        self.sm.add_transition(WorkflowState.IDLE, "start", WorkflowState.LOADING)
        self.sm.add_transition(WorkflowState.LOADING, "loaded", WorkflowState.PROCESSING)
        self.sm.add_transition(WorkflowState.LOADING, "error", WorkflowState.ERROR)
        self.sm.add_transition(WorkflowState.PROCESSING, "complete", WorkflowState.REVIEWING)
        self.sm.add_transition(WorkflowState.PROCESSING, "error", WorkflowState.ERROR)
        self.sm.add_transition(WorkflowState.REVIEWING, "approved", WorkflowState.COMPLETED)
        self.sm.add_transition(WorkflowState.REVIEWING, "rejected", WorkflowState.PROCESSING)
        self.sm.add_transition(WorkflowState.ERROR, "retry", WorkflowState.LOADING)

    def run(self, data_path: str):
        """执行工作流"""

        # 启动
        self.sm.trigger("start", data_path=data_path)

        # 模拟加载
        print("加载数据...")
        self.sm.trigger("loaded")

        # 模拟处理
        print("处理数据...")
        self.sm.trigger("complete")

        # 审查
        print("等待审查...")
        # 这里会等待人类输入
        # 假设批准
        self.sm.trigger("approved")

        print(f"工作流完成,最终状态: {self.sm.state}")

代码示例

完整的状态机工作流

"""
完整的状态机工作流示例
"""
import json
from typing import Dict, Any, List, Optional
from dataclasses import dataclass, field, asdict
from enum import Enum
import time

class State(Enum):
    """状态枚举"""
    IDLE = "idle"
    LOAD_DATA = "load_data"
    IDENTIFY_SOURCES = "identify_sources"
    BUILD_RESISTANCE = "build_resistance"
    REVIEW_SOURCES = "review_sources"
    REVIEW_RESISTANCE = "review_resistance"
    EXTRACT_CORRIDORS = "extract_corridors"
    COMPLETED = "completed"
    ERROR = "error"

@dataclass
class WorkflowState:
    """工作流状态数据"""
    current: State = State.IDLE
    step_number: int = 0
    data_path: Optional[str] = None
    sources: Optional[List[Dict]] = None
    resistance_weights: Optional[Dict] = None
    corridors: Optional[List[Dict]] = None
    errors: List[str] = field(default_factory=list)
    history: List[Dict] = field(default_factory=list)

    def transition_to(self, new_state: State, action: str = ""):
        """状态转换"""
        old_state = self.current
        self.current = new_state
        self.step_number += 1

        # 记录历史
        self.history.append({
            "step": self.step_number,
            "from": old_state.value,
            "to": new_state.value,
            "action": action,
            "timestamp": time.time()
        })

    def to_dict(self) -> Dict:
        """序列化"""
        return {
            "current": self.current.value,
            "step_number": self.step_number,
            "data_path": self.data_path,
            "sources": self.sources,
            "resistance_weights": self.resistance_weights,
            "corridors": self.corridors,
            "errors": self.errors,
            "history": self.history
        }

    def save(self, path: str):
        """保存状态到文件"""
        with open(path, 'w') as f:
            json.dump(self.to_dict(), f, indent=2)

    @classmethod
    def load(cls, path: str) -> 'WorkflowState':
        """从文件加载状态"""
        with open(path, 'r') as f:
            data = json.load(f)

        # 转换State枚举
        data["current"] = State(data["current"])

        return cls(**{k: v for k, v in data.items() if k != "history"})

class EcologicalAnalysisAgent:
    """生态分析智能体(有状态)"""

    def __init__(self):
        self.state = WorkflowState()
        self.review_callbacks = {
            State.REVIEW_SOURCES: self._review_sources,
            State.REVIEW_RESISTANCE: self._review_resistance
        }

    def start(self, data_path: str):
        """启动分析"""
        self.state.data_path = data_path
        self.state.transition_to(State.LOAD_DATA, "开始加载数据")
        self._execute_current_step()

    def _execute_current_step(self):
        """执行当前状态对应的操作"""
        handlers = {
            State.LOAD_DATA: self._handle_load_data,
            State.IDENTIFY_SOURCES: self._handle_identify_sources,
            State.BUILD_RESISTANCE: self._handle_build_resistance,
            State.REVIEW_SOURCES: self._handle_review,
            State.REVIEW_RESISTANCE: self._handle_review,
            State.EXTRACT_CORRIDORS: self._handle_extract_corridors,
            State.COMPLETED: self._handle_completed,
            State.ERROR: self._handle_error
        }

        handler = handlers.get(self.state.current)
        if handler:
            handler()

    def _handle_load_data(self):
        """处理数据加载"""
        print(f"\n[状态: {self.state.current.value}] 加载数据: {self.state.data_path}")

        # 模拟加载
        try:
            # 这里实际会读取文件
            time.sleep(0.5)
            print("数据加载成功")
            self.state.transition_to(State.IDENTIFY_SOURCES, "数据加载完成")
            self._execute_current_step()
        except Exception as e:
            self.state.errors.append(str(e))
            self.state.transition_to(State.ERROR, f"加载失败: {e}")
            self._execute_current_step()

    def _handle_identify_sources(self):
        """处理源地识别"""
        print(f"\n[状态: {self.state.current.value}] 识别生态源地...")

        # 模拟识别
        self.state.sources = [
            {"id": 1, "area": 1500, "type": "forest"},
            {"id": 2, "area": 800, "type": "wetland"}
        ]
        print(f"识别到 {len(self.state.sources)} 个源地")

        self.state.transition_to(State.REVIEW_SOURCES, "源地识别完成,等待审查")
        self._execute_current_step()

    def _handle_build_resistance(self):
        """处理阻力面构建"""
        print(f"\n[状态: {self.state.current.value}] 构建阻力面...")

        # 模拟构建
        self.state.resistance_weights = {
            "forest": 1,
            "grassland": 10,
            "urban": 100,
            "water": 50
        }
        print("阻力面构建完成")

        self.state.transition_to(State.REVIEW_RESISTANCE, "阻力面构建完成,等待审查")
        self._execute_current_step()

    def _handle_review(self):
        """处理审查状态"""
        print(f"\n[状态: {self.state.current.value}] 等待人类审查...")

        callback = self.review_callbacks.get(self.state.current)
        if callback:
            result = callback()

            if result == "approve":
                if self.state.current == State.REVIEW_SOURCES:
                    self.state.transition_to(State.BUILD_RESISTANCE, "审查通过")
                elif self.state.current == State.REVIEW_RESISTANCE:
                    self.state.transition_to(State.EXTRACT_CORRIDORS, "审查通过")
                self._execute_current_step()
            else:
                # 拒绝,返回上一状态
                print("审查未通过,重新执行...")
                # 简化处理:直接继续

    def _review_sources(self) -> str:
        """审查源地"""
        print("\n=== 源地审查 ===")
        print(f"识别到 {len(self.state.sources)} 个源地:")
        for s in self.state.sources:
            print(f"  - ID {s['id']}: {s['type']}, 面积 {s['area']}")

        # 实际实现中这里会等待人类输入
        # 这里模拟自动批准
        print("\n[模拟] 审查: 批准")
        return "approve"

    def _review_resistance(self) -> str:
        """审查阻力面"""
        print("\n=== 阻力面审查 ===")
        print("阻力权重:")
        for land_type, weight in self.state.resistance_weights.items():
            print(f"  - {land_type}: {weight}")

        print("\n[模拟] 审查: 批准")
        return "approve"

    def _handle_extract_corridors(self):
        """处理廊道提取"""
        print(f"\n[状态: {self.state.current.value}] 提取生态廊道...")

        # 模拟提取
        self.state.corridors = [
            {"from": 1, "to": 2, "length": 3500}
        ]
        print(f"提取到 {len(self.state.corridors)} 条廊道")

        self.state.transition_to(State.COMPLETED, "分析完成")
        self._execute_current_step()

    def _handle_completed(self):
        """处理完成状态"""
        print(f"\n[状态: {self.state.current.value}] 工作流完成!")
        print(f"\n=== 结果摘要 ===")
        print(f"源地数量: {len(self.state.sources) if self.state.sources else 0}")
        print(f"廊道数量: {len(self.state.corridors) if self.state.corridors else 0}")
        print(f"执行步骤: {self.state.step_number}")

    def _handle_error(self):
        """处理错误状态"""
        print(f"\n[状态: {self.state.current.value}] 发生错误")
        for error in self.state.errors:
            print(f"  - {error}")

    def save_state(self, path: str):
        """保存当前状态"""
        self.state.save(path)
        print(f"状态已保存到: {path}")

    def resume_from(self, path: str):
        """从保存的状态恢复"""
        self.state = WorkflowState.load(path)
        print(f"从状态恢复: {self.state.current.value}")
        print(f"历史步骤: {self.state.step_number}")
        self._execute_current_step()

# 使用示例
if __name__ == "__main__":
    print("=== 生态分析状态机工作流 ===\n")

    agent = EcologicalAnalysisAgent()

    # 执行工作流
    agent.start("data.geojson")

    # 可以保存状态
    # agent.save_state("workflow_state.json")

    # 可以从状态恢复
    # new_agent = EcologicalAnalysisAgent()
    # new_agent.resume_from("workflow_state.json")

案例分析

LangGraph在空间分析中的应用

"""
LangGraph风格的生态网络分析工作流
"""
from typing import TypedDict, Annotated, Literal
from operator import add

class EcologicalState(TypedDict):
    """生态分析状态类型"""
    messages: Annotated[list, add]  # 消息历史
    input_data: dict
    sources: list
    resistance: dict
    corridors: list
    next_step: str
    human_feedback: str

# 节点函数
def load_data_node(state: EcologicalState) -> EcologicalState:
    """加载数据节点"""
    print("执行: load_data")
    state["sources"] = [{"id": 1, "area": 1000}]
    state["next_step"] = "identify"
    return state

def identify_sources_node(state: EcologicalState) -> EcologicalState:
    """识别源地节点"""
    print("执行: identify_sources")
    state["sources"] = [{"id": i, "area": i * 100} for i in range(1, 6)]
    state["next_step"] = "review"
    return state

def human_review_node(state: EcologicalState) -> EcologicalState:
    """人类审查节点"""
    print("执行: human_review")
    print(f"待审查: {state['sources']}")

    # 在实际实现中,这里会等待人类输入
    state["human_feedback"] = "approved"
    state["next_step"] = "build_resistance"
    return state

def build_resistance_node(state: EcologicalState) -> EcologicalState:
    """构建阻力面节点"""
    print("执行: build_resistance")
    state["resistance"] = {"forest": 1, "urban": 100}
    state["next_step"] = "complete"
    return state

# 路由函数
def should_review(state: EcologicalState) -> Literal["review", "skip"]:
    """决定是否需要审查"""
    if len(state.get("sources", [])) > 3:
        return "review"
    return "skip"

# 条件边
def route_after_identify(state: EcologicalState) -> str:
    """识别源地后的路由"""
    if state.get("human_feedback") == "approved":
        return "build_resistance"
    return "identify"  # 重新识别

print("""
┌──────────────┐
│ load_data    │
└──────┬───────┘
┌──────────────┐
│identify_sources│
└──────┬───────┘
       ├────→ [review?] ──→ human_review ──┐
       │ No                                │
       ↓                                   ↓
┌──────────────┐                   ┌──────────────┐
│build_resistance│◀──────────────────│   approved   │
└──────────────┘                   └──────────────┘
""")

反思与延伸

思考问题

  1. 状态粒度:状态应该有多细?太细会怎样,太粗会怎样?

  2. 持久化策略:哪些状态需要持久化?什么时候保存状态?

  3. 并发处理:如果多个Agent协同工作,如何管理共享状态?

  4. 调试:当状态机出错时,如何调试?

实践练习

  1. 状态审计:添加状态转换日志,分析工作流执行路径

  2. 状态压缩:实现状态序列化/反序列化,支持断点续传

  3. 条件路由:实现一个带多个分支的状态机

延伸阅读

  • "Designing Data-Intensive Applications" (Kleppmann) - 状态管理理论
  • LangGraph文档 - 实际框架使用
  • "State Machine Design Patterns" - 状态机设计模式

关键要点

  1. 状态是系统在某一时刻的完整快照
  2. 状态机描述系统如何随事件转换状态
  3. LangGraph用图结构表达有状态的工作流
  4. 良好的状态管理支持断点续传和HITL
  5. 状态历史是调试和可解释性的关键