Files
2026_DesignAI/officefile/supplements/03-autonomous-design/03.1-workflow-orchestration.md
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

22 KiB
Raw Blame History

03.1 工作流编排原理

核心问题

如何设计复杂的多步骤分析流程? DAG(有向无环图)如何表达工作流? 如何处理工作流中的条件分支和错误?


概念讲解

工作流编排的核心

工作流编排是指协调多个处理步骤按顺序执行的能力:

简单顺序                     复杂编排
    │                            │
    ↓                            ↓
┌─────┐    ┌─────┐    ┌─────┐   ╔═════════════════╗
│Step1│───→│Step2│───→│Step3│   ║   条件分支       ║
└─────┘    └─────┘    └─────┘   ║                  ║
                                  ║  ┌─────┐       ║
    一维流程                       ║  │     │       ║
                                  ║  ↓     ↓       ║
                                  ║  Yes   No      ║
                                  ║  │     │       ║
                                  ║  ↓     ↓       ║
                                  ║┌─────┐ ┌─────┐ ║
                                  ║│StepA│ │StepB│ ║
                                  ║└─────┘ └─────┘ ║
                                  ║       │        ║
                                  ║       └────┬─── ║
                                  ║            │     ║
                                  ╔═════════════════╝

DAG:有向无环图

DAG (Directed Acyclic Graph) 是工作流编排的基础数据结构:

"""
DAG的数学表示

DAG = (V, E)
其中:
- V: 节点集合(处理步骤)
- E: 边集合(依赖关系)
- 条件:无环(没有节点能通过边回到自己)

性质:
1. 有方向:边从上游指向下游
2. 无环:没有循环依赖
3. 可拓扑排序:可以找到线性执行顺序
"""

为什么DAG适合工作流?

特性 说明
明确依赖 边定义了步骤间的依赖关系
可并行化 无依赖的步骤可并行执行
可验证 可以检测循环依赖
可可视化 容易理解和调试

工作流的组成要素

工作流 = 节点 + 边 + 条件 + 错误处理

┌─────────────────────────────────────────────────────────────┐
│                                                             │
│  ┌──────────────┐                                           │
│  │    节点      │  ←─────────────────────────────────────   │
│  │  ──────────  │                                           │
│  │  - 执行函数  │  输入 → 处理 → 输出                       │
│  │  - 输入/输出 │                                           │
│  │  - 副后置    │  before(), execute(), after()             │
│  └──────────────┘                                           │
│         │                                                    │
│         ↓                                                    │
│  ┌──────────────┐                                           │
│  │    边        │  ←─────────────────────────────────────   │
│  │  ──────────  │                                           │
│  │  - 数据流    │  上游输出 → 下游输入                       │
│  │  - 依赖关系  │  顺序执行                                  │
│  │  - 条件路由  │  基于状态选择路径                          │
│  └──────────────┘                                           │
│         │                                                    │
│         ↓                                                    │
│  ┌──────────────┐                                           │
│  │   条件分支   │  ←─────────────────────────────────────   │
│  │  ──────────  │                                           │
│  │  - 分支条件  │  if state.value > threshold: ...          │
│  │  - 路由选择  │  switch-case模式                          │
│  │  - 合并点    │  多路径汇聚                                │
│  └──────────────┘                                           │
│         │                                                    │
│         ↓                                                    │
│  ┌──────────────┐                                           │
│  │  错误处理    │  ←─────────────────────────────────────   │
│  │  ──────────  │                                           │
│  │  - 重试      │  失败后重新执行                            │
│  │  - 回滚      │  恢复到之前状态                            │
│  │  - 降级      │  使用备选方案                              │
│  │  - 告警      │  通知相关人员                              │
│  └──────────────┘                                           │
│                                                             │
└─────────────────────────────────────────────────────────────┘

设计原理

工作流设计模式

1. 线性流水线 (Linear Pipeline)

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)

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)

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)

错误处理模式

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)

代码示例

生态网络六阶段工作流

"""
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:
        """阶段4MCR分析"""
        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的工作流实现

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. 并行执行可以显著提升效率