# 01.2 状态与状态机 ## 核心问题 > 在AI系统中,"状态"到底是什么? > 为什么状态管理是Agent系统的核心? > LangGraph是如何用状态机设计工作流的? --- ## 概念讲解 ### 什么是状态 **状态**是系统在某一时刻的快照,包含所有影响未来行为的信息: ``` 系统状态 = 所有相关的变量值 例如:生态分析工作流的状态 { "input_data": {...}, # 输入数据 "current_step": "buffer", # 当前步骤 "intermediate_results": {...}, # 中间结果 "user_preferences": {...}, # 用户偏好 "error_count": 0, # 错误计数 "checkpoint_reached": False # 检查点状态 } ``` **状态的类型**: | 类型 | 说明 | 例子 | |-----|------|------| | 静态状态 | 初始输入,不变化 | 输入文件路径、参数 | | 动态状态 | 运行中变化 | 当前步骤、累积结果 | | 控制状态 | 影响流程走向 | 分支条件、错误标志 | | 会话状态 | 跨请求持久化 | 用户偏好、历史记录 | ### 为什么状态管理很重要 **1. 断点续传** ```python # 没有状态管理:出错后必须从头开始 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. 人机协同** ```python # 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. 调试和可解释性** ```python # 状态历史记录了整个决策过程 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的框架,其核心思想: ```python 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. **条件分支**:基于状态的动态路由 ### 工作流状态机实现 ```python 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}") ``` --- ## 代码示例 ### 完整的状态机工作流 ```python """ 完整的状态机工作流示例 """ 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在空间分析中的应用 ```python """ 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. **状态历史是调试和可解释性的关键**