Files
2026_DesignAI/dofile/examples/01-foundations/state_machine.py
T
pengxiao 219232de74 refactor: 重组项目目录结构
以讲义内容为骨架迁移到标准目录格式:
- officefile/ 主内容(12章 + 附录 + CC4SI补充)
- dofile/ 代码示例(11个Python脚本)
- data/ 图片资源
- output/ 生成输出(忽略)
- Archive/ 归档旧目录(忽略)
- .claude/skills/ 保留markdown-to-docx工具链
- .pandoc/ 保留CSL和本地化配置

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

749 lines
24 KiB
Python

"""
状态机工作流示例 (State Machine Workflow Example)
==================================================
本示例展示如何使用状态机模式管理复杂的空间分析工作流。
状态机是一种行为设计模式,允许对象在其内部状态改变时改变其行为。
核心概念:
1. 状态 (State) - 系统在特定时刻的模式
2. 转换 (Transition) - 从一个状态到另一个状态的变化
3. 事件 (Event) - 触发状态转换的外部或内部条件
4. 动作 (Action) - 状态转换时执行的操作
应用场景:
- 空间数据处理流水线
- 多阶段决策流程
- 任务调度与监控
- 用户交互流程控制
作者: CC4SI 项目组
"""
from abc import ABC, abstractmethod
from typing import Dict, List, Optional, Callable, Any
from dataclasses import dataclass, field
from enum import Enum
import json
from datetime import datetime
# ============================================================================
# 状态定义
# ============================================================================
class WorkflowState(Enum):
"""工作流状态枚举"""
# 初始状态
IDLE = "idle"
INITIALIZED = "initialized"
# 数据处理状态
LOADING_DATA = "loading_data"
DATA_LOADED = "data_loaded"
VALIDATING_DATA = "validating_data"
DATA_VALIDATED = "data_validated"
# 分析状态
ANALYZING = "analyzing"
ANALYSIS_COMPLETE = "analysis_complete"
# 决策状态
DECIDING = "deciding"
DECISION_MADE = "decision_made"
# 输出状态
GENERATING_OUTPUT = "generating_output"
OUTPUT_COMPLETE = "output_complete"
# 异常状态
ERROR = "error"
PAUSED = "paused"
CANCELLED = "cancelled"
# 最终状态
COMPLETED = "completed"
class EventType(Enum):
"""事件类型枚举"""
# 控制事件
START = "start"
PAUSE = "pause"
RESUME = "resume"
CANCEL = "cancel"
RESET = "reset"
# 数据事件
DATA_LOAD_REQUEST = "data_load_request"
DATA_LOAD_SUCCESS = "data_load_success"
DATA_LOAD_FAILURE = "data_load_failure"
DATA_VALIDATE_REQUEST = "data_validate_request"
DATA_VALIDATE_SUCCESS = "data_validate_success"
DATA_VALIDATE_FAILURE = "data_validate_failure"
# 分析事件
ANALYZE_REQUEST = "analyze_request"
ANALYSIS_SUCCESS = "analysis_success"
ANALYSIS_FAILURE = "analysis_failure"
# 决策事件
DECIDE_REQUEST = "decide_request"
DECISION_SUCCESS = "decision_success"
DECISION_FAILURE = "decision_failure"
# 输出事件
OUTPUT_REQUEST = "output_request"
OUTPUT_SUCCESS = "output_success"
OUTPUT_FAILURE = "output_failure"
# 错误事件
ERROR_OCCURRED = "error_occurred"
RETRY = "retry"
# ============================================================================
# 状态机数据结构
# ============================================================================
@dataclass
class StateTransition:
"""状态转换定义"""
from_state: WorkflowState
event: EventType
to_state: WorkflowState
action: Optional[Callable] = None
guard: Optional[Callable[[], bool]] = None # 守卫条件
description: str = ""
def can_execute(self) -> bool:
"""检查转换是否可执行"""
if self.guard is None:
return True
return self.guard()
@dataclass
class StateContext:
"""状态上下文 - 存储工作流数据"""
data: Dict[str, Any] = field(default_factory=dict)
errors: List[str] = field(default_factory=list)
warnings: List[str] = field(default_factory=list)
history: List[Dict[str, Any]] = field(default_factory=list)
start_time: Optional[datetime] = None
end_time: Optional[datetime] = None
def add_history(self, from_state: WorkflowState, event: EventType,
to_state: WorkflowState, timestamp: datetime = None):
"""添加历史记录"""
self.history.append({
"from_state": from_state.value,
"event": event.value,
"to_state": to_state.value,
"timestamp": timestamp or datetime.now()
})
def get_data(self, key: str, default: Any = None) -> Any:
"""获取数据"""
return self.data.get(key, default)
def set_data(self, key: str, value: Any) -> None:
"""设置数据"""
self.data[key] = value
def add_error(self, error: str) -> None:
"""添加错误"""
self.errors.append(error)
def add_warning(self, warning: str) -> None:
"""添加警告"""
self.warnings.append(warning)
# ============================================================================
# 状态机实现
# ============================================================================
class StateMachine:
"""
状态机实现
管理状态转换和状态相关的行为。
"""
def __init__(self, initial_state: WorkflowState = WorkflowState.IDLE):
"""
初始化状态机
Args:
initial_state: 初始状态
"""
self._current_state = initial_state
self._transitions: Dict[WorkflowState, Dict[EventType, StateTransition]] = {}
self._context = StateContext()
self._state_listeners: Dict[WorkflowState, List[Callable]] = {}
print(f"[状态机] 初始化,初始状态: {initial_state.value}")
@property
def current_state(self) -> WorkflowState:
"""获取当前状态"""
return self._current_state
@property
def context(self) -> StateContext:
"""获取状态上下文"""
return self._context
def add_transition(self, transition: StateTransition) -> None:
"""
添加状态转换
Args:
transition: 状态转换定义
"""
if transition.from_state not in self._transitions:
self._transitions[transition.from_state] = {}
self._transitions[transition.from_state][transition.event] = transition
print(f"[状态机] 添加转换: {transition.from_state.value} + {transition.event.value} -> {transition.to_state.value}")
def add_state_listener(self, state: WorkflowState, listener: Callable) -> None:
"""
添加状态监听器
Args:
state: 要监听的状态
listener: 状态进入时调用的函数
"""
if state not in self._state_listeners:
self._state_listeners[state] = []
self._state_listeners[state].append(listener)
def trigger(self, event: EventType, payload: Any = None) -> bool:
"""
触发事件
Args:
event: 事件类型
payload: 事件负载
Returns:
是否成功触发状态转换
"""
# 检查当前状态是否有对应转换
if self._current_state not in self._transitions:
print(f"[状态机] 当前状态 {self._current_state.value} 没有定义任何转换")
return False
if event not in self._transitions[self._current_state]:
print(f"[状态机] 状态 {self._current_state.value} 不处理事件 {event.value}")
return False
transition = self._transitions[self._current_state][event]
# 检查守卫条件
if not transition.can_execute():
print(f"[状态机] 守卫条件不满足,转换被阻止")
return False
# 执行状态转换
old_state = self._current_state
self._current_state = transition.to_state
# 记录历史
if payload:
self._context.set_data("last_payload", payload)
self._context.add_history(old_state, event, self._current_state)
print(f"[状态机] 状态转换: {old_state.value} -> {self._current_state.value} (事件: {event.value})")
# 执行转换动作
if transition.action:
try:
transition.action(self._context, payload)
except Exception as e:
print(f"[状态机] 执行动作时出错: {e}")
self._context.add_error(f"转换动作执行失败: {e}")
# 触发状态监听器
if self._current_state in self._state_listeners:
for listener in self._state_listeners[self._current_state]:
try:
listener(self._current_state, self._context)
except Exception as e:
print(f"[状态机] 监听器执行出错: {e}")
return True
def can_trigger(self, event: EventType) -> bool:
"""
检查是否可以触发指定事件
Args:
event: 事件类型
Returns:
是否可以触发
"""
if self._current_state not in self._transitions:
return False
if event not in self._transitions[self._current_state]:
return False
transition = self._transitions[self._current_state][event]
return transition.can_execute()
def get_available_events(self) -> List[EventType]:
"""获取当前状态下可用的事件列表"""
if self._current_state not in self._transitions:
return []
available = []
for event, transition in self._transitions[self._current_state].items():
if transition.can_execute():
available.append(event)
return available
def reset(self) -> None:
"""重置状态机"""
self._current_state = WorkflowState.IDLE
self._context = StateContext()
print(f"[状态机] 状态机已重置")
def print_state(self) -> None:
"""打印当前状态"""
print(f"\n当前状态: {self._current_state.value}")
available = self.get_available_events()
if available:
print(f"可用事件: {', '.join(e.value for e in available)}")
else:
print("可用事件: 无")
# ============================================================================
# 空间分析工作流状态机
# ============================================================================
class SpatialAnalysisWorkflow:
"""
空间分析工作流
使用状态机实现的空间数据处理和分析工作流。
"""
def __init__(self):
"""初始化工作流"""
self.state_machine = StateMachine()
self._setup_transitions()
self._setup_listeners()
def _setup_transitions(self):
"""设置状态转换"""
sm = self.state_machine
# 启动流程
sm.add_transition(StateTransition(
from_state=WorkflowState.IDLE,
event=EventType.START,
to_state=WorkflowState.INITIALIZED,
action=self._action_initialize,
description="初始化工作流"
))
# 数据加载
sm.add_transition(StateTransition(
from_state=WorkflowState.INITIALIZED,
event=EventType.DATA_LOAD_REQUEST,
to_state=WorkflowState.LOADING_DATA,
action=self._action_load_data,
description="开始加载数据"
))
sm.add_transition(StateTransition(
from_state=WorkflowState.LOADING_DATA,
event=EventType.DATA_LOAD_SUCCESS,
to_state=WorkflowState.DATA_LOADED,
action=self._action_on_data_loaded,
description="数据加载成功"
))
sm.add_transition(StateTransition(
from_state=WorkflowState.LOADING_DATA,
event=EventType.DATA_LOAD_FAILURE,
to_state=WorkflowState.ERROR,
action=self._action_on_error,
description="数据加载失败"
))
# 数据验证
sm.add_transition(StateTransition(
from_state=WorkflowState.DATA_LOADED,
event=EventType.DATA_VALIDATE_REQUEST,
to_state=WorkflowState.VALIDATING_DATA,
action=self._action_validate_data,
description="开始验证数据"
))
sm.add_transition(StateTransition(
from_state=WorkflowState.VALIDATING_DATA,
event=EventType.DATA_VALIDATE_SUCCESS,
to_state=WorkflowState.DATA_VALIDATED,
description="数据验证成功"
))
sm.add_transition(StateTransition(
from_state=WorkflowState.VALIDATING_DATA,
event=EventType.DATA_VALIDATE_FAILURE,
to_state=WorkflowState.ERROR,
action=self._action_on_error,
description="数据验证失败"
))
# 分析
sm.add_transition(StateTransition(
from_state=WorkflowState.DATA_VALIDATED,
event=EventType.ANALYZE_REQUEST,
to_state=WorkflowState.ANALYZING,
action=self._action_analyze,
description="开始分析"
))
sm.add_transition(StateTransition(
from_state=WorkflowState.ANALYZING,
event=EventType.ANALYSIS_SUCCESS,
to_state=WorkflowState.ANALYSIS_COMPLETE,
description="分析完成"
))
# 决策
sm.add_transition(StateTransition(
from_state=WorkflowState.ANALYSIS_COMPLETE,
event=EventType.DECIDE_REQUEST,
to_state=WorkflowState.DECIDING,
action=self._action_decide,
description="开始决策"
))
sm.add_transition(StateTransition(
from_state=WorkflowState.DECIDING,
event=EventType.DECISION_SUCCESS,
to_state=WorkflowState.DECISION_MADE,
description="决策完成"
))
# 输出
sm.add_transition(StateTransition(
from_state=WorkflowState.DECISION_MADE,
event=EventType.OUTPUT_REQUEST,
to_state=WorkflowState.GENERATING_OUTPUT,
action=self._action_generate_output,
description="生成输出"
))
sm.add_transition(StateTransition(
from_state=WorkflowState.GENERATING_OUTPUT,
event=EventType.OUTPUT_SUCCESS,
to_state=WorkflowState.OUTPUT_COMPLETE,
description="输出完成"
))
# 完成
sm.add_transition(StateTransition(
from_state=WorkflowState.OUTPUT_COMPLETE,
event=EventType.START,
to_state=WorkflowState.COMPLETED,
action=self._action_complete,
description="工作流完成"
))
# 错误恢复
sm.add_transition(StateTransition(
from_state=WorkflowState.ERROR,
event=EventType.RETRY,
to_state=WorkflowState.INITIALIZED,
guard=lambda: len(self.state_machine.context.errors) < 3,
description="重试"
))
sm.add_transition(StateTransition(
from_state=WorkflowState.ERROR,
event=EventType.RESET,
to_state=WorkflowState.IDLE,
action=self._action_reset,
description="重置"
))
def _setup_listeners(self):
"""设置状态监听器"""
sm = self.state_machine
# 错误状态监听器
sm.add_state_listener(WorkflowState.ERROR, self._on_error_state)
# 完成状态监听器
sm.add_state_listener(WorkflowState.COMPLETED, self._on_complete_state)
# ------------------------------------------------------------------------
# 状态动作
# ------------------------------------------------------------------------
def _action_initialize(self, ctx: StateContext, payload: Any):
"""初始化动作"""
ctx.start_time = datetime.now()
ctx.set_data("workflow_id", f"WF-{datetime.now().strftime('%Y%m%d%H%M%S')}")
print(" [动作] 工作流初始化完成")
def _action_load_data(self, ctx: StateContext, payload: Any):
"""加载数据动作"""
source = payload or "默认数据源"
print(f" [动作] 从 '{source}' 加载数据...")
# 模拟数据加载
ctx.set_data("raw_data", [
{"id": 1, "x": 10, "y": 20, "value": 100},
{"id": 2, "x": 30, "y": 40, "value": 200},
{"id": 3, "x": 50, "y": 60, "value": 150},
])
# 模拟成功
self.state_machine.trigger(EventType.DATA_LOAD_SUCCESS)
def _action_on_data_loaded(self, ctx: StateContext, payload: Any):
"""数据加载完成动作"""
data_count = len(ctx.get_data("raw_data", []))
print(f" [动作] 数据加载完成,共 {data_count} 条记录")
def _action_validate_data(self, ctx: StateContext, payload: Any):
"""验证数据动作"""
print(" [动作] 验证数据...")
data = ctx.get_data("raw_data", [])
valid = all("id" in item and "x" in item and "y" in item for item in data)
if valid:
self.state_machine.trigger(EventType.DATA_VALIDATE_SUCCESS)
else:
ctx.add_error("数据验证失败: 缺少必需字段")
self.state_machine.trigger(EventType.DATA_VALIDATE_FAILURE)
def _action_analyze(self, ctx: StateContext, payload: Any):
"""分析动作"""
print(" [动作] 执行空间分析...")
data = ctx.get_data("raw_data", [])
values = [item.get("value", 0) for item in data]
avg = sum(values) / len(values) if values else 0
ctx.set_data("analysis_result", {
"average": avg,
"count": len(data),
"min": min(values) if values else 0,
"max": max(values) if values else 0
})
print(f" [动作] 分析完成,平均值: {avg:.2f}")
self.state_machine.trigger(EventType.ANALYSIS_SUCCESS)
def _action_decide(self, ctx: StateContext, payload: Any):
"""决策动作"""
print(" [动作] 执行决策...")
analysis = ctx.get_data("analysis_result", {})
avg = analysis.get("average", 0)
if avg > 150:
decision = "高价值区域"
elif avg > 100:
decision = "中等价值区域"
else:
decision = "低价值区域"
ctx.set_data("decision", decision)
print(f" [动作] 决策完成: {decision}")
self.state_machine.trigger(EventType.DECISION_SUCCESS)
def _action_generate_output(self, ctx: StateContext, payload: Any):
"""生成输出动作"""
print(" [动作] 生成输出报告...")
report = {
"workflow_id": ctx.get_data("workflow_id"),
"data_count": len(ctx.get_data("raw_data", [])),
"analysis": ctx.get_data("analysis_result"),
"decision": ctx.get_data("decision")
}
ctx.set_data("output", report)
print(" [动作] 输出生成完成")
self.state_machine.trigger(EventType.OUTPUT_SUCCESS)
def _action_complete(self, ctx: StateContext, payload: Any):
"""完成动作"""
ctx.end_time = datetime.now()
duration = (ctx.end_time - ctx.start_time).total_seconds() if ctx.start_time else 0
ctx.set_data("duration", duration)
print(f" [动作] 工作流完成,耗时: {duration:.2f}")
def _action_on_error(self, ctx: StateContext, payload: Any):
"""错误处理动作"""
print(f" [动作] 发生错误")
def _action_reset(self, ctx: StateContext, payload: Any):
"""重置动作"""
print(" [动作] 重置工作流")
# ------------------------------------------------------------------------
# 状态监听器
# ------------------------------------------------------------------------
def _on_error_state(self, state: WorkflowState, ctx: StateContext):
"""错误状态处理"""
print(f" [监听器] 进入错误状态")
print(f" [监听器] 错误列表: {ctx.errors}")
def _on_complete_state(self, state: WorkflowState, ctx: StateContext):
"""完成状态处理"""
print(f" [监听器] 工作流已完成")
output = ctx.get_data("output")
if output:
print(f" [监听器] 最终输出: {json.dumps(output, ensure_ascii=False, indent=2)}")
# ------------------------------------------------------------------------
# 公共接口
# ------------------------------------------------------------------------
def start(self, data_source: str = None) -> bool:
"""启动工作流"""
return self.state_machine.trigger(EventType.START, data_source)
def execute_full_workflow(self, data_source: str = None) -> Dict[str, Any]:
"""
执行完整工作流
Args:
data_source: 数据源
Returns:
执行结果
"""
print("\n" + "="*60)
print("执行完整空间分析工作流")
print("="*60)
# 启动
if not self.start(data_source):
return {"success": False, "error": "启动失败"}
# 等待异步操作完成 (简化版: 手动触发)
# 在实际应用中,这些事件会由异步操作触发
return {
"success": True,
"final_state": self.state_machine.current_state.value,
"context": self.state_machine.context.data
}
def print_history(self):
"""打印状态转换历史"""
history = self.state_machine.context.history
print(f"\n状态转换历史 (共 {len(history)} 次):")
print("-" * 70)
for i, h in enumerate(history, 1):
ts = h.get("timestamp", datetime.now()).strftime("%H:%M:%S")
print(f"{i:2d}. [{ts}] {h['from_state']:20s} -> {h['to_state']:20s} ({h['event']})")
print("-" * 70)
# ============================================================================
# 主程序
# ============================================================================
def main():
"""主程序 - 演示状态机工作流的使用"""
print("="*70)
print("状态机工作流示例演示")
print("="*70)
# 1. 创建工作流
print("\n[步骤 1] 创建空间分析工作流")
workflow = SpatialAnalysisWorkflow()
# 2. 显示初始状态
print("\n[步骤 2] 初始状态")
workflow.state_machine.print_state()
# 3. 手动执行状态转换
print("\n[步骤 3] 手动执行状态转换")
# 启动
print("\n3.1 启动工作流:")
workflow.state_machine.trigger(EventType.START)
workflow.state_machine.print_state()
# 请求数据加载 (这将触发加载动作,然后自动触发成功事件)
print("\n3.2 请求数据加载:")
workflow.state_machine.trigger(EventType.DATA_LOAD_REQUEST, "sample.csv")
workflow.state_machine.print_state()
# 请求数据验证
print("\n3.3 请求数据验证:")
workflow.state_machine.trigger(EventType.DATA_VALIDATE_REQUEST)
workflow.state_machine.print_state()
# 请求分析
print("\n3.4 请求分析:")
workflow.state_machine.trigger(EventType.ANALYZE_REQUEST)
workflow.state_machine.print_state()
# 请求决策
print("\n3.5 请求决策:")
workflow.state_machine.trigger(EventType.DECIDE_REQUEST)
workflow.state_machine.print_state()
# 请求输出
print("\n3.6 请求输出:")
workflow.state_machine.trigger(EventType.OUTPUT_REQUEST)
workflow.state_machine.print_state()
# 完成
print("\n3.7 完成工作流:")
workflow.state_machine.trigger(EventType.START)
workflow.state_machine.print_state()
# 4. 显示转换历史
print("\n[步骤 4] 状态转换历史")
workflow.print_history()
# 5. 演示错误处理
print("\n[步骤 5] 演示错误处理和恢复")
print("\n5.1 重置状态机:")
workflow.state_machine.reset()
workflow.state_machine.print_state()
print("\n5.2 启动后触发错误:")
workflow.state_machine.trigger(EventType.START)
workflow.state_machine.trigger(EventType.DATA_LOAD_FAILURE)
workflow.state_machine.print_state()
print("\n5.3 尝试重试:")
if workflow.state_machine.can_trigger(EventType.RETRY):
workflow.state_machine.trigger(EventType.RETRY)
workflow.state_machine.print_state()
else:
print(" 无法重试 (已达到最大重试次数)")
print("\n5.4 重置工作流:")
workflow.state_machine.trigger(EventType.RESET)
workflow.state_machine.print_state()
print("\n" + "="*70)
print("演示完成!")
print("="*70)
if __name__ == "__main__":
main()