219232de74
以讲义内容为骨架迁移到标准目录格式: - 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>
635 lines
22 KiB
Markdown
635 lines
22 KiB
Markdown
# 03.1 工作流编排原理
|
||
|
||
## 核心问题
|
||
|
||
> 如何设计复杂的多步骤分析流程?
|
||
> DAG(有向无环图)如何表达工作流?
|
||
> 如何处理工作流中的条件分支和错误?
|
||
|
||
---
|
||
|
||
## 概念讲解
|
||
|
||
### 工作流编排的核心
|
||
|
||
工作流编排是指**协调多个处理步骤按顺序执行**的能力:
|
||
|
||
```
|
||
简单顺序 复杂编排
|
||
│ │
|
||
↓ ↓
|
||
┌─────┐ ┌─────┐ ┌─────┐ ╔═════════════════╗
|
||
│Step1│───→│Step2│───→│Step3│ ║ 条件分支 ║
|
||
└─────┘ └─────┘ └─────┘ ║ ║
|
||
║ ┌─────┐ ║
|
||
一维流程 ║ │ │ ║
|
||
║ ↓ ↓ ║
|
||
║ Yes No ║
|
||
║ │ │ ║
|
||
║ ↓ ↓ ║
|
||
║┌─────┐ ┌─────┐ ║
|
||
║│StepA│ │StepB│ ║
|
||
║└─────┘ └─────┘ ║
|
||
║ │ ║
|
||
║ └────┬─── ║
|
||
║ │ ║
|
||
╔═════════════════╝
|
||
```
|
||
|
||
### DAG:有向无环图
|
||
|
||
**DAG** (Directed Acyclic Graph) 是工作流编排的基础数据结构:
|
||
|
||
```python
|
||
"""
|
||
DAG的数学表示
|
||
|
||
DAG = (V, E)
|
||
其中:
|
||
- V: 节点集合(处理步骤)
|
||
- E: 边集合(依赖关系)
|
||
- 条件:无环(没有节点能通过边回到自己)
|
||
|
||
性质:
|
||
1. 有方向:边从上游指向下游
|
||
2. 无环:没有循环依赖
|
||
3. 可拓扑排序:可以找到线性执行顺序
|
||
"""
|
||
```
|
||
|
||
**为什么DAG适合工作流?**
|
||
|
||
| 特性 | 说明 |
|
||
|-----|------|
|
||
| **明确依赖** | 边定义了步骤间的依赖关系 |
|
||
| **可并行化** | 无依赖的步骤可并行执行 |
|
||
| **可验证** | 可以检测循环依赖 |
|
||
| **可可视化** | 容易理解和调试 |
|
||
|
||
### 工作流的组成要素
|
||
|
||
```
|
||
工作流 = 节点 + 边 + 条件 + 错误处理
|
||
|
||
┌─────────────────────────────────────────────────────────────┐
|
||
│ │
|
||
│ ┌──────────────┐ │
|
||
│ │ 节点 │ ←───────────────────────────────────── │
|
||
│ │ ────────── │ │
|
||
│ │ - 执行函数 │ 输入 → 处理 → 输出 │
|
||
│ │ - 输入/输出 │ │
|
||
│ │ - 副后置 │ before(), execute(), after() │
|
||
│ └──────────────┘ │
|
||
│ │ │
|
||
│ ↓ │
|
||
│ ┌──────────────┐ │
|
||
│ │ 边 │ ←───────────────────────────────────── │
|
||
│ │ ────────── │ │
|
||
│ │ - 数据流 │ 上游输出 → 下游输入 │
|
||
│ │ - 依赖关系 │ 顺序执行 │
|
||
│ │ - 条件路由 │ 基于状态选择路径 │
|
||
│ └──────────────┘ │
|
||
│ │ │
|
||
│ ↓ │
|
||
│ ┌──────────────┐ │
|
||
│ │ 条件分支 │ ←───────────────────────────────────── │
|
||
│ │ ────────── │ │
|
||
│ │ - 分支条件 │ if state.value > threshold: ... │
|
||
│ │ - 路由选择 │ switch-case模式 │
|
||
│ │ - 合并点 │ 多路径汇聚 │
|
||
│ └──────────────┘ │
|
||
│ │ │
|
||
│ ↓ │
|
||
│ ┌──────────────┐ │
|
||
│ │ 错误处理 │ ←───────────────────────────────────── │
|
||
│ │ ────────── │ │
|
||
│ │ - 重试 │ 失败后重新执行 │
|
||
│ │ - 回滚 │ 恢复到之前状态 │
|
||
│ │ - 降级 │ 使用备选方案 │
|
||
│ │ - 告警 │ 通知相关人员 │
|
||
│ └──────────────┘ │
|
||
│ │
|
||
└─────────────────────────────────────────────────────────────┘
|
||
```
|
||
|
||
---
|
||
|
||
## 设计原理
|
||
|
||
### 工作流设计模式
|
||
|
||
**1. 线性流水线 (Linear Pipeline)**
|
||
|
||
```python
|
||
class LinearWorkflow:
|
||
"""最简单的工作流:顺序执行"""
|
||
|
||
def __init__(self):
|
||
self.steps = []
|
||
|
||
def add_step(self, func, name=None):
|
||
"""添加步骤"""
|
||
self.steps.append({
|
||
'function': func,
|
||
'name': name or func.__name__
|
||
})
|
||
return self
|
||
|
||
def execute(self, initial_data):
|
||
"""执行工作流"""
|
||
result = initial_data
|
||
|
||
for step in self.steps:
|
||
print(f"执行: {step['name']}")
|
||
result = step['function'](result)
|
||
|
||
return result
|
||
|
||
# 使用示例
|
||
workflow = LinearWorkflow()
|
||
workflow.add_step(load_data, "加载数据")
|
||
workflow.add_step(clean_data, "清理数据")
|
||
workflow.add_step(analyze_data, "分析数据")
|
||
|
||
result = workflow.execute("data.geojson")
|
||
```
|
||
|
||
**2. 条件分支工作流 (Conditional Workflow)**
|
||
|
||
```python
|
||
class ConditionalWorkflow:
|
||
"""带条件分支的工作流"""
|
||
|
||
def __init__(self):
|
||
self.steps = {}
|
||
self.conditions = {}
|
||
self.transitions = {}
|
||
|
||
def add_step(self, name, func):
|
||
"""添加步骤"""
|
||
self.steps[name] = func
|
||
return self
|
||
|
||
def add_condition(self, name, condition_func):
|
||
"""添加条件判断"""
|
||
self.conditions[name] = condition_func
|
||
return self
|
||
|
||
def add_transition(self, from_step, condition, to_step):
|
||
"""添加状态转换"""
|
||
if from_step not in self.transitions:
|
||
self.transitions[from_step] = {}
|
||
self.transitions[from_step][condition] = to_step
|
||
return self
|
||
|
||
def execute(self, initial_data, start_step='start'):
|
||
"""执行工作流"""
|
||
current_step = start_step
|
||
state = {'data': initial_data}
|
||
|
||
while current_step != 'end':
|
||
print(f"当前步骤: {current_step}")
|
||
|
||
# 执行步骤
|
||
if current_step in self.steps:
|
||
result = self.steps[current_step](state['data'])
|
||
state['data'] = result
|
||
|
||
# 检查条件
|
||
if current_step in self.transitions:
|
||
transition = self.transitions[current_step]
|
||
matched = False
|
||
|
||
for condition_name, next_step in transition.items():
|
||
if condition_name in self.conditions:
|
||
if self.conditions[condition_name](state):
|
||
current_step = next_step
|
||
matched = True
|
||
break
|
||
|
||
if not matched and 'default' in transition:
|
||
current_step = transition['default']
|
||
elif not matched:
|
||
current_step = 'end'
|
||
else:
|
||
current_step = 'end'
|
||
|
||
return state['data']
|
||
|
||
# 使用示例:生态网络工作流
|
||
workflow = ConditionalWorkflow()
|
||
|
||
# 添加步骤
|
||
workflow.add_step('load_data', load_data)
|
||
workflow.add_step('identify_sources', identify_sources)
|
||
workflow.add_step('human_review', human_review)
|
||
workflow.add_step('build_resistance', build_resistance)
|
||
|
||
# 添加条件
|
||
workflow.add_condition('high_uncertainty',
|
||
lambda s: s.get('uncertainty', 0) > 0.3)
|
||
workflow.add_condition('approved',
|
||
lambda s: s.get('decision') == 'approve')
|
||
|
||
# 添加转换
|
||
workflow.add_transition('identify_sources', 'high_uncertainty', 'human_review')
|
||
workflow.add_transition('identify_sources', 'default', 'build_resistance')
|
||
workflow.add_transition('human_review', 'approved', 'build_resistance')
|
||
workflow.add_transition('human_review', 'default', 'identify_sources') # 重做
|
||
workflow.add_transition('build_resistance', 'default', 'end')
|
||
```
|
||
|
||
**3. 并行工作流 (Parallel Workflow)**
|
||
|
||
```python
|
||
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||
|
||
class ParallelWorkflow:
|
||
"""支持并行执行的工作流"""
|
||
|
||
def __init__(self, max_workers=4):
|
||
self.max_workers = max_workers
|
||
self.parallel_groups = {}
|
||
|
||
def add_parallel_group(self, group_name, tasks):
|
||
"""添加可并行执行的任务组"""
|
||
self.parallel_groups[group_name] = tasks
|
||
return self
|
||
|
||
def execute_group(self, group_name, shared_data):
|
||
"""执行一个并行任务组"""
|
||
if group_name not in self.parallel_groups:
|
||
raise ValueError(f"未知的任务组: {group_name}")
|
||
|
||
tasks = self.parallel_groups[group_name]
|
||
results = {}
|
||
|
||
with ThreadPoolExecutor(max_workers=self.max_workers) as executor:
|
||
# 提交所有任务
|
||
future_to_task = {
|
||
executor.submit(task['func'], shared_data): task['name']
|
||
for task in tasks
|
||
}
|
||
|
||
# 收集结果
|
||
for future in as_completed(future_to_task):
|
||
task_name = future_to_task[future]
|
||
try:
|
||
results[task_name] = future.result()
|
||
except Exception as e:
|
||
results[task_name] = {'error': str(e)}
|
||
|
||
return results
|
||
|
||
# 使用示例:同时处理多个区域
|
||
workflow = ParallelWorkflow(max_workers=4)
|
||
|
||
workflow.add_parallel_group('process_regions', [
|
||
{'name': 'region_north', 'func': process_north_region},
|
||
{'name': 'region_south', 'func': process_south_region},
|
||
{'name': 'region_east', 'func': process_east_region},
|
||
{'name': 'region_west', 'func': process_west_region},
|
||
])
|
||
|
||
results = workflow.execute_group('process_regions', shared_data)
|
||
```
|
||
|
||
### 错误处理模式
|
||
|
||
```python
|
||
class WorkflowErrorHandling:
|
||
"""工作流错误处理"""
|
||
|
||
class RetryPolicy:
|
||
"""重试策略"""
|
||
def __init__(self, max_retries=3, backoff=2.0):
|
||
self.max_retries = max_retries
|
||
self.backoff = backoff # 指数退避因子
|
||
|
||
def should_retry(self, attempt, error):
|
||
return attempt < self.max_retries
|
||
|
||
def get_delay(self, attempt):
|
||
return self.backoff ** attempt
|
||
|
||
def execute_with_retry(self, func, *args, retry_policy=None, **kwargs):
|
||
"""带重试的执行"""
|
||
if retry_policy is None:
|
||
retry_policy = self.RetryPolicy()
|
||
|
||
last_error = None
|
||
for attempt in range(retry_policy.max_retries + 1):
|
||
try:
|
||
return func(*args, **kwargs)
|
||
except Exception as e:
|
||
last_error = e
|
||
if retry_policy.should_retry(attempt, e):
|
||
delay = retry_policy.get_delay(attempt)
|
||
print(f"尝试 {attempt + 1} 失败,{delay}秒后重试...")
|
||
time.sleep(delay)
|
||
else:
|
||
break
|
||
|
||
raise last_error
|
||
|
||
def execute_with_fallback(self, primary_func, fallback_func, *args, **kwargs):
|
||
"""带降级的执行"""
|
||
try:
|
||
return primary_func(*args, **kwargs)
|
||
except Exception as e:
|
||
print(f"主函数失败: {e},使用降级方案")
|
||
return fallback_func(*args, **kwargs)
|
||
```
|
||
|
||
---
|
||
|
||
## 代码示例
|
||
|
||
### 生态网络六阶段工作流
|
||
|
||
```python
|
||
"""
|
||
ENAgent的完整工作流编排实现
|
||
|
||
六阶段:
|
||
1. 数据准备
|
||
2. 源地识别
|
||
3. 阻力面构建
|
||
4. MCR分析
|
||
5. 廊道提取
|
||
6. 结果评估
|
||
"""
|
||
import time
|
||
from typing import Dict, List, Optional, Callable
|
||
from enum import Enum
|
||
from dataclasses import dataclass
|
||
|
||
class Stage(Enum):
|
||
"""工作流阶段"""
|
||
DATA_PREPARATION = "data_preparation"
|
||
SOURCE_IDENTIFICATION = "source_identification"
|
||
RESISTANCE_SURFACE = "resistance_surface"
|
||
MCR_ANALYSIS = "mcr_analysis"
|
||
CORRIDOR_EXTRACTION = "corridor_extraction"
|
||
RESULT_EVALUATION = "result_evaluation"
|
||
COMPLETED = "completed"
|
||
ERROR = "error"
|
||
|
||
@dataclass
|
||
class WorkflowState:
|
||
"""工作流状态"""
|
||
current_stage: Stage
|
||
data: Dict
|
||
results: Dict
|
||
errors: List[str]
|
||
stage_history: List[Stage]
|
||
checkpoint_data: Optional[Dict] = None
|
||
|
||
class EcologicalNetworkWorkflow:
|
||
"""生态网络分析工作流"""
|
||
|
||
def __init__(self):
|
||
self.stages = {
|
||
Stage.DATA_PREPARATION: self._data_preparation,
|
||
Stage.SOURCE_IDENTIFICATION: self._source_identification,
|
||
Stage.RESISTANCE_SURFACE: self._resistance_surface,
|
||
Stage.MCR_ANALYSIS: self._mcr_analysis,
|
||
Stage.CORRIDOR_EXTRACTION: self._corridor_extraction,
|
||
Stage.RESULT_EVALUATION: self._result_evaluation,
|
||
}
|
||
|
||
self.transitions = {
|
||
Stage.DATA_PREPARATION: Stage.SOURCE_IDENTIFICATION,
|
||
Stage.SOURCE_IDENTIFICATION: Stage.RESISTANCE_SURFACE,
|
||
Stage.RESISTANCE_SURFACE: Stage.MCR_ANALYSIS,
|
||
Stage.MCR_ANALYSIS: Stage.CORRIDOR_EXTRACTION,
|
||
Stage.CORRIDOR_EXTRACTION: Stage.RESULT_EVALUATION,
|
||
Stage.RESULT_EVALUATION: Stage.COMPLETED,
|
||
}
|
||
|
||
# HITL审查点
|
||
self.checkpoints = {
|
||
Stage.SOURCE_IDENTIFICATION: True,
|
||
Stage.RESISTANCE_SURFACE: True,
|
||
Stage.CORRIDOR_EXTRACTION: False,
|
||
}
|
||
|
||
def execute(self, initial_data: Dict) -> WorkflowState:
|
||
"""执行完整工作流"""
|
||
state = WorkflowState(
|
||
current_stage=Stage.DATA_PREPARATION,
|
||
data=initial_data,
|
||
results={},
|
||
errors=[],
|
||
stage_history=[Stage.DATA_PREPARATION]
|
||
)
|
||
|
||
while state.current_stage != Stage.COMPLETED:
|
||
if state.current_stage == Stage.ERROR:
|
||
print("工作流因错误终止")
|
||
break
|
||
|
||
# 执行当前阶段
|
||
state = self._execute_stage(state)
|
||
|
||
# 检查是否需要审查
|
||
if self.checkpoints.get(state.current_stage, False):
|
||
state = self._handle_checkpoint(state)
|
||
|
||
# 转换到下一阶段
|
||
if state.current_stage != Stage.ERROR:
|
||
next_stage = self.transitions.get(
|
||
state.current_stage,
|
||
Stage.COMPLETED
|
||
)
|
||
state.current_stage = next_stage
|
||
state.stage_history.append(next_stage)
|
||
|
||
return state
|
||
|
||
def _execute_stage(self, state: WorkflowState) -> WorkflowState:
|
||
"""执行单个阶段"""
|
||
stage = state.current_stage
|
||
print(f"\n{'='*50}")
|
||
print(f"执行阶段: {stage.value}")
|
||
print('='*50)
|
||
|
||
try:
|
||
# 执行阶段函数
|
||
result = self.stages[stage](state.data, state.results)
|
||
|
||
# 保存结果
|
||
state.results[stage.value] = result
|
||
|
||
except Exception as e:
|
||
print(f"阶段 {stage.value} 执行失败: {e}")
|
||
state.errors.append(str(e))
|
||
state.current_stage = Stage.ERROR
|
||
|
||
return state
|
||
|
||
def _handle_checkpoint(self, state: WorkflowState) -> WorkflowState:
|
||
"""处理HITL审查点"""
|
||
stage = state.current_stage
|
||
print(f"\n[审查点: {stage.value}]")
|
||
|
||
# 实际实现中,这里会等待人类输入
|
||
# 模拟审查通过
|
||
approval = self._get_human_approval(state)
|
||
|
||
if not approval:
|
||
print("审查未通过,调整参数后重新执行...")
|
||
# 可以在这里修改state.data后返回同一阶段
|
||
|
||
return state
|
||
|
||
def _get_human_approval(self, state: WorkflowState) -> bool:
|
||
"""获取人类批准(模拟)"""
|
||
print(f"待审查结果: {list(state.results.keys())}")
|
||
# 实际实现中等待输入
|
||
return True
|
||
|
||
# === 阶段实现 ===
|
||
|
||
def _data_preparation(self, data: Dict, results: Dict) -> Dict:
|
||
"""阶段1:数据准备"""
|
||
print("加载和处理原始数据...")
|
||
time.sleep(0.5)
|
||
return {
|
||
'landcover': 'loaded',
|
||
'elevation': 'loaded',
|
||
'boundary': 'loaded'
|
||
}
|
||
|
||
def _source_identification(self, data: Dict, results: Dict) -> Dict:
|
||
"""阶段2:源地识别"""
|
||
print("识别生态源地...")
|
||
sources = [
|
||
{'id': 1, 'area': 1500, 'type': 'forest'},
|
||
{'id': 2, 'area': 800, 'type': 'wetland'}
|
||
]
|
||
return {'sources': sources, 'n_sources': len(sources)}
|
||
|
||
def _resistance_surface(self, data: Dict, results: Dict) -> Dict:
|
||
"""阶段3:阻力面构建"""
|
||
print("构建生态阻力面...")
|
||
return {
|
||
'resistance_surface': 'computed',
|
||
'weights': {'forest': 1, 'urban': 100}
|
||
}
|
||
|
||
def _mcr_analysis(self, data: Dict, results: Dict) -> Dict:
|
||
"""阶段4:MCR分析"""
|
||
print("执行最小累积阻力分析...")
|
||
return {'mcr_surface': 'computed'}
|
||
|
||
def _corridor_extraction(self, data: Dict, results: Dict) -> Dict:
|
||
"""阶段5:廊道提取"""
|
||
print("提取生态廊道...")
|
||
corridors = [
|
||
{'from': 1, 'to': 2, 'length': 3500}
|
||
]
|
||
return {'corridors': corridors, 'n_corridors': len(corridors)}
|
||
|
||
def _result_evaluation(self, data: Dict, results: Dict) -> Dict:
|
||
"""阶段6:结果评估"""
|
||
print("评估分析结果...")
|
||
return {
|
||
'connectivity_index': 0.75,
|
||
'network_efficiency': 0.82
|
||
}
|
||
|
||
# 使用示例
|
||
if __name__ == "__main__":
|
||
print("=== 生态网络分析工作流 ===\n")
|
||
|
||
workflow = EcologicalNetworkWorkflow()
|
||
|
||
initial_data = {
|
||
'landcover_path': 'data/landcover.tif',
|
||
'species': 'target_species',
|
||
'parameters': {}
|
||
}
|
||
|
||
final_state = workflow.execute(initial_data)
|
||
|
||
print("\n=== 工作流完成 ===")
|
||
print(f"执行的阶段数: {len(final_state.stage_history)}")
|
||
print(f"产生的错误: {final_state.errors}")
|
||
print(f"最终结果: {list(final_state.results.keys())}")
|
||
```
|
||
|
||
---
|
||
|
||
## 案例分析
|
||
|
||
### LangGraph的工作流实现
|
||
|
||
```python
|
||
from langgraph.graph import StateGraph, END
|
||
from typing import TypedDict
|
||
|
||
class ENAgentState(TypedDict):
|
||
"""ENAgent工作流状态"""
|
||
stage: str
|
||
data: dict
|
||
results: dict
|
||
requires_review: bool
|
||
|
||
def build_enagent_workflow():
|
||
"""构建ENAgent工作流"""
|
||
|
||
# 创建图
|
||
workflow = StateGraph(ENAgentState)
|
||
|
||
# 添加节点
|
||
workflow.add_node("prepare_data", prepare_data_node)
|
||
workflow.add_node("identify_sources", identify_sources_node)
|
||
workflow.add_node("build_resistance", build_resistance_node)
|
||
workflow.add_node("mcr_analysis", mcr_analysis_node)
|
||
workflow.add_node("extract_corridors", extract_corridors_node)
|
||
workflow.add_node("evaluate_results", evaluate_results_node)
|
||
|
||
# 添加边(线性流程)
|
||
workflow.set_entry_point("prepare_data")
|
||
workflow.add_edge("prepare_data", "identify_sources")
|
||
workflow.add_edge("identify_sources", "build_resistance")
|
||
workflow.add_edge("build_resistance", "mcr_analysis")
|
||
workflow.add_edge("mcr_analysis", "extract_corridors")
|
||
workflow.add_edge("extract_corridors", "evaluate_results")
|
||
workflow.add_edge("evaluate_results", END)
|
||
|
||
# 编译
|
||
return workflow.compile()
|
||
```
|
||
|
||
---
|
||
|
||
## 反思与延伸
|
||
|
||
### 思考问题
|
||
|
||
1. **工作流设计**:你的项目中有哪些可以自动化的步骤?
|
||
|
||
2. **错误处理**:当某个步骤失败时,应该重试、跳过还是终止?
|
||
|
||
3. **审查点**:在你的工作流中,哪些地方需要人类介入?
|
||
|
||
4. **并行化**:哪些步骤可以并行执行以提升效率?
|
||
|
||
### 延伸阅读
|
||
|
||
- **"Dataflow Programming"** - 数据流编程范式
|
||
- **"Workflow Patterns"** (Van der Aalst) - 工作流模式
|
||
- **Apache Airflow 文档** - 工作流调度系统
|
||
|
||
---
|
||
|
||
## 关键要点
|
||
|
||
1. **DAG是工作流的核心数据结构**,表达依赖关系
|
||
2. **节点**是处理步骤,**边**定义执行顺序
|
||
3. **条件分支**根据状态动态选择执行路径
|
||
4. **错误处理**是生产工作流的关键
|
||
5. **并行执行**可以显著提升效率
|