a90f7adfa1
将 Markdown 源文件移入 md/,LaTeX 工作目录保留在 latex/, Word 导出移入 word/;删除临时脚本、调试截图和空 stub。 Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
22 KiB
22 KiB
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 │
↓ │
┌─────────┐ │
│ 完成 │←───────────────────────┘
└─────────┘
状态机的要素:
- 状态 (State):系统可能处于的情况
- 事件 (Event):触发状态转换的条件
- 转换 (Transition):从一个状态到另一个状态
- 动作 (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()
核心概念:
- 状态即消息:状态在节点间传递
- 图即流程:有向图描述工作流
- 条件分支:基于状态的动态路由
工作流状态机实现
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 │
└──────────────┘ └──────────────┘
""")
反思与延伸
思考问题
-
状态粒度:状态应该有多细?太细会怎样,太粗会怎样?
-
持久化策略:哪些状态需要持久化?什么时候保存状态?
-
并发处理:如果多个Agent协同工作,如何管理共享状态?
-
调试:当状态机出错时,如何调试?
实践练习
-
状态审计:添加状态转换日志,分析工作流执行路径
-
状态压缩:实现状态序列化/反序列化,支持断点续传
-
条件路由:实现一个带多个分支的状态机
延伸阅读
- "Designing Data-Intensive Applications" (Kleppmann) - 状态管理理论
- LangGraph文档 - 实际框架使用
- "State Machine Design Patterns" - 状态机设计模式
关键要点
- 状态是系统在某一时刻的完整快照
- 状态机描述系统如何随事件转换状态
- LangGraph用图结构表达有状态的工作流
- 良好的状态管理支持断点续传和HITL
- 状态历史是调试和可解释性的关键