Agent工作流编排引擎设计与实现

举报
柠檬🍋 发表于 2026/09/17 20:20:35 2026/09/17
【摘要】 Agent工作流编排引擎设计与实现在现代Agent系统中,工作流编排引擎是整个架构的核心枢纽。一个设计良好的工作流引擎能够将复杂的业务逻辑分解为可复用、可测试、可监控的任务单元,并通过DAG(有向无环图)来描述任务之间的依赖关系,实现自动化执行与故障恢复。本文将从工作流DSL设计出发,逐步深入到DAG任务图构建、条件分支与循环、并行执行控制、失败补偿机制,最终给出一个完整的工作流引擎代码实...

Agent工作流编排引擎设计与实现

在现代Agent系统中,工作流编排引擎是整个架构的核心枢纽。一个设计良好的工作流引擎能够将复杂的业务逻辑分解为可复用、可测试、可监控的任务单元,并通过DAG(有向无环图)来描述任务之间的依赖关系,实现自动化执行与故障恢复。本文将从工作流DSL设计出发,逐步深入到DAG任务图构建、条件分支与循环、并行执行控制、失败补偿机制,最终给出一个完整的工作流引擎代码实现。

一、工作流编排的核心问题

Agent系统面临的典型场景包括:多步骤推理链、工具调用编排、多Agent协作、数据处理流水线等。这些场景的共同特征是:任务之间存在复杂的依赖关系,部分任务可以并行执行,某些分支需要根据前置结果动态决定走向,执行过程中可能出现失败需要重试或补偿。

传统的方式是硬编码这些逻辑,用if-else和函数调用堆叠出执行流程。这种方式的问题在于:流程不可视、难以修改、无法复用、缺乏统一的监控和恢复机制。工作流编排引擎的思路是将流程定义与执行分离——用DSL描述"做什么"和"怎么做",由引擎负责"何时做"和"做失败了怎么办"。

工作流引擎需要解决以下几个核心问题:

第一,流程描述能力。引擎需要提供足够丰富的DSL来描述任务节点、依赖关系、条件分支、循环结构等。DSL的设计要在表达力和简洁性之间取得平衡,既不能过于简单导致无法描述复杂流程,也不能过于复杂导致学习成本过高。

第二,执行调度能力。引擎需要根据DAG拓扑排序确定执行顺序,识别可并行任务并并发调度,管理任务的生命周期(等待、执行、完成、失败),处理超时和资源限制。

第三,容错恢复能力。引擎需要提供重试、补偿、跳过等失败处理策略,支持检查点机制实现断点续跑,保证流程的最终一致性。

第四,可观测性。引擎需要记录每个任务的执行状态、耗时、输入输出,提供流程可视化能力,支持审计和调试。

二、工作流DSL设计

DSL是工作流引擎与使用者之间的接口。一个好的DSL设计应该具备声明式、可组合、可序列化的特征。我们采用YAML格式作为DSL的载体,因为它兼具人类可读性和机器可解析性。

name: "data-analysis-pipeline"
version: "1.0"
description: "数据分析流水线工作流"
inputs:
  - name: data_source
    type: string
    required: true
  - name: analysis_mode
    type: string
    default: "full"
nodes:
  - id: fetch_data
    type: task
    handler: "DataFetcher"
    params:
      source: "${inputs.data_source}"
    retry:
      max_attempts: 3
      backoff: exponential
      base_delay: 1000
  - id: validate_data
    type: task
    handler: "DataValidator"
    depends_on: [fetch_data]
    params:
      data: "${fetch_data.output}"
  - id: check_quality
    type: condition
    depends_on: [validate_data]
    branches:
      - when: "${validate_data.output.quality_score > 0.8}"
        next: process_data
      - when: "${validate_data.output.quality_score > 0.5}"
        next: clean_data
      - default: reject_data
  - id: clean_data
    type: task
    handler: "DataCleaner"
    depends_on: [check_quality]
    params:
      data: "${validate_data.output}"
  - id: process_data
    type: task
    handler: "DataProcessor"
    depends_on: [check_quality, clean_data]
    params:
      data: "${clean_data.output}"
      mode: "${inputs.analysis_mode}"
  - id: generate_report
    type: task
    handler: "ReportGenerator"
    depends_on: [process_data]
    params:
      result: "${process_data.output}"
  - id: reject_data
    type: task
    handler: "RejectionHandler"
    depends_on: [check_quality]
    params:
      reason: "${validate_data.output.issues}"
outputs:
  report: "${generate_report.output}"

这个DSL定义了一个数据分析流水线,包含数据获取、校验、质量检查(条件分支)、清洗、处理、报告生成等节点。每个节点有唯一的id、类型(task/condition)、处理器、参数、依赖关系和重试策略。

DSL中的表达式语法${node_id.output.field}用于引用其他节点的输出,引擎在执行时负责解析这些引用并注入实际值。条件节点通过branches定义分支逻辑,支持多条件匹配和默认分支。

三、DAG任务图构建

将DSL解析为内存中的DAG图是引擎的第一步工作。DAG的节点对应DSL中的task节点,边对应depends_on声明的依赖关系。构建过程需要做以下几件事:

解析DSL文件,为每个节点创建Node对象。Node对象包含id、type、handler、params、depends_on、retry策略等字段。然后根据depends_on建立边关系,构建邻接表。接着进行拓扑排序,检测是否存在环——如果存在环则报错,因为DAG不允许有环。最后计算每个节点的入度(前置依赖数量),入度为0的节点可以作为起始节点立即执行。

from dataclasses import dataclass, field
from typing import Dict, List, Any, Optional, Set
from collections import defaultdict, deque
import yaml
import re
import time
import threading
import traceback
from enum import Enum

class NodeType(Enum):
    TASK = "task"
    CONDITION = "condition"
    LOOP = "loop"

class NodeStatus(Enum):
    PENDING = "pending"
    RUNNING = "running"
    SUCCESS = "success"
    FAILED = "failed"
    SKIPPED = "skipped"
    COMPENSATING = "compensating"

@dataclass
class RetryPolicy:
    max_attempts: int = 1
    backoff: str = "fixed"
    base_delay: float = 1000

@dataclass
class Node:
    id: str
    node_type: NodeType
    handler: str = ""
    params: Dict[str, Any] = field(default_factory=dict)
    depends_on: List[str] = field(default_factory=list)
    retry: RetryPolicy = field(default_factory=RetryPolicy)
    branches: List[Dict[str, Any]] = field(default_factory=list)
    status: NodeStatus = NodeStatus.PENDING
    output: Any = None
    error: Optional[str] = None
    attempts: int = 0
    started_at: Optional[float] = None
    finished_at: Optional[float] = None

@dataclass
class WorkflowContext:
    workflow_id: str
    inputs: Dict[str, Any]
    outputs: Dict[str, Any] = field(default_factory=dict)
    node_outputs: Dict[str, Any] = field(default_factory=dict)
    checkpoint: Dict[str, Any] = field(default_factory=dict)

class DAGBuilder:
    """DAG构建器:解析DSL并构建任务图"""
    
    def __init__(self):
        self.nodes: Dict[str, Node] = {}
        self.edges: Dict[str, List[str]] = defaultdict(list)
        self.reverse_edges: Dict[str, List[str]] = defaultdict(list)
    
    def from_yaml(self, yaml_content: str) -> 'DAGBuilder':
        spec = yaml.safe_load(yaml_content)
        for node_spec in spec.get('nodes', []):
            node = Node(
                id=node_spec['id'],
                node_type=NodeType(node_spec.get('type', 'task')),
                handler=node_spec.get('handler', ''),
                params=node_spec.get('params', {}),
                depends_on=node_spec.get('depends_on', []),
                branches=node_spec.get('branches', []),
            )
            if 'retry' in node_spec:
                r = node_spec['retry']
                node.retry = RetryPolicy(
                    max_attempts=r.get('max_attempts', 1),
                    backoff=r.get('backoff', 'fixed'),
                    base_delay=r.get('base_delay', 1000),
                )
            self.nodes[node.id] = node
            for dep in node.depends_on:
                self.edges[dep].append(node.id)
                self.reverse_edges[node.id].append(dep)
        return self
    
    def topological_sort(self) -> List[str]:
        """拓扑排序,检测环"""
        in_degree = {nid: len(n.depends_on) for nid, n in self.nodes.items()}
        queue = deque([nid for nid, d in in_degree.items() if d == 0])
        result = []
        while queue:
            nid = queue.popleft()
            result.append(nid)
            for child in self.edges.get(nid, []):
                in_degree[child] -= 1
                if in_degree[child] == 0:
                    queue.append(child)
        if len(result) != len(self.nodes):
            raise ValueError("DAG中检测到环,无法进行拓扑排序")
        return result
    
    def get_ready_nodes(self) -> List[str]:
        """获取所有依赖已满足且未执行的节点"""
        ready = []
        for nid, node in self.nodes.items():
            if node.status != NodeStatus.PENDING:
                continue
            deps_satisfied = all(
                self.nodes[dep].status == NodeStatus.SUCCESS
                for dep in node.depends_on
            )
            if deps_satisfied:
                ready.append(nid)
        return ready
    
    def get_dependents(self, node_id: str) -> List[str]:
        """获取直接依赖该节点的后续节点"""
        return self.edges.get(node_id, [])

上面的代码实现了DAG构建器,包含从YAML解析、拓扑排序(含环检测)、就绪节点查询等功能。get_ready_nodes方法在每次有节点完成时被调用,返回所有前置依赖已满足的待执行节点,这是并行调度的基础。

四、条件分支与循环

条件分支是工作流引擎的高级特性。当某个节点执行完成后,可能需要根据其输出结果决定后续执行路径。条件节点的branches列表定义了多个when表达式,引擎按顺序匹配,第一个匹配成功的分支决定下一个执行的节点。如果没有匹配项,则走default分支。

循环结构允许对集合数据进行迭代处理。例如,需要对一批文档逐一分析时,可以用loop节点遍历文档列表,每次迭代执行一组子任务。循环的实现方式是将loop节点展开为多个等价的任务节点,每个对应一次迭代。

class ExpressionResolver:
    """表达式解析器:解析${node.output.field}形式的引用"""
    
    PATTERN = re.compile(r'\$\{([^}]+)\}')
    
    def resolve(self, expr: str, context: WorkflowContext) -> Any:
        if not isinstance(expr, str):
            return expr
        matches = self.PATTERN.findall(expr)
        if not matches:
            return expr
        result = expr
        for match in matches:
            value = self._lookup(match.strip(), context)
            if len(matches) == 1 and expr.strip() == f"${{{match}}}":
                return value
            result = result.replace(f"${{{match}}}", str(value))
        return result
    
    def _lookup(self, path: str, context: WorkflowContext) -> Any:
        parts = path.split('.')
        if parts[0] == 'inputs':
            obj = context.inputs
        elif parts[0] == 'outputs':
            obj = context.outputs
        else:
            obj = context.node_outputs.get(parts[0])
        for part in parts[1:]:
            if obj is None:
                return None
            if isinstance(obj, dict):
                obj = obj.get(part)
            elif isinstance(obj, list):
                obj = obj[int(part)] if part.isdigit() and int(part) < len(obj) else None
            else:
                obj = getattr(obj, part, None)
        return obj
    
    def resolve_params(self, params: Dict[str, Any], context: WorkflowContext) -> Dict[str, Any]:
        resolved = {}
        for key, value in params.items():
            if isinstance(value, dict):
                resolved[key] = self.resolve_params(value, context)
            elif isinstance(value, list):
                resolved[key] = [self.resolve(v, context) if isinstance(v, str) else v for v in value]
            elif isinstance(value, str):
                resolved[key] = self.resolve(value, context)
            else:
                resolved[key] = value
        return resolved
    
    def evaluate_condition(self, expr: str, context: WorkflowContext) -> bool:
        resolved = self.resolve(expr, context)
        if isinstance(resolved, bool):
            return resolved
        if isinstance(resolved, str):
            try:
                return bool(eval(resolved, {"__builtins__": {}}, {}))
            except Exception:
                return False
        return bool(resolved)


class BranchResolver:
    """分支解析器:处理条件节点的分支选择"""
    
    def __init__(self, resolver: ExpressionResolver):
        self.resolver = resolver
    
    def resolve_branch(self, node: Node, context: WorkflowContext) -> str:
        """返回下一个应该执行的节点id"""
        for branch in node.branches:
            if 'when' in branch:
                if self.resolver.evaluate_condition(branch['when'], context):
                    return branch['next']
        for branch in node.branches:
            if 'default' in branch:
                return branch['default']
        raise ValueError(f"条件节点 {node.id} 没有匹配的分支且无默认分支")

表达式解析器ExpressionResolver负责将${node.output.field}形式的引用替换为实际值。它支持嵌套路径访问(如fetch_data.output.quality_score),支持字典和列表的索引访问。BranchResolver利用表达式解析器来评估条件分支的when表达式,返回匹配的下一个节点id。

五、并行执行控制

工作流引擎的一个关键能力是并行执行。当多个节点的依赖都已满足时,这些节点可以同时执行,从而缩短整体流程时间。并行执行需要考虑几个问题:并发度控制(避免资源耗尽)、线程安全(节点状态更新需要加锁)、错误传播(某个节点失败后如何处理其后续节点)。

class WorkflowEngine:
    """工作流引擎:调度执行DAG中的任务"""
    
    def __init__(self, dag: DAGBuilder, handlers: Dict[str, Any],
                 max_workers: int = 4):
        self.dag = dag
        self.handlers = handlers
        self.max_workers = max_workers
        self.resolver = ExpressionResolver()
        self.branch_resolver = BranchResolver(self.resolver)
        self.lock = threading.Lock()
        self.executor = threading.ThreadPoolExecutor(max_workers=max_workers)
        self.event_listeners: List[Any] = []
    
    def add_listener(self, listener):
        self.event_listeners.append(listener)
    
    def _emit_event(self, event_type: str, node: Node):
        for listener in self.event_listeners:
            try:
                listener.on_event(event_type, node)
            except Exception:
                pass
    
    def run(self, inputs: Dict[str, Any]) -> WorkflowContext:
        context = WorkflowContext(
            workflow_id=f"wf-{int(time.time() * 1000)}",
            inputs=inputs,
        )
        self.dag.topological_sort()
        futures = {}
        while True:
            with self.lock:
                ready = self.dag.get_ready_nodes()
                pending_count = sum(
                    1 for n in self.dag.nodes.values()
                    if n.status in (NodeStatus.PENDING, NodeStatus.RUNNING)
                )
            if not ready and pending_count == 0:
                break
            if not ready and pending_count > 0:
                time.sleep(0.05)
                continue
            for node_id in ready:
                with self.lock:
                    node = self.dag.nodes[node_id]
                    if node.status != NodeStatus.PENDING:
                        continue
                    node.status = NodeStatus.RUNNING
                    node.started_at = time.time()
                self._emit_event("node_started", node)
                future = self.executor.submit(self._execute_node, node, context)
                futures[future] = node_id
            done_futures = [f for f in futures if f.done()]
            for f in done_futures:
                nid = futures.pop(f)
                node = self.dag.nodes[nid]
                try:
                    f.result()
                except Exception as e:
                    pass
                self._handle_node_completion(node, context)
        return context
    
    def _execute_node(self, node: Node, context: WorkflowContext):
        """执行单个节点,包含重试逻辑"""
        last_error = None
        for attempt in range(node.retry.max_attempts):
            node.attempts = attempt + 1
            try:
                if node.node_type == NodeType.CONDITION:
                    next_id = self.branch_resolver.resolve_branch(node, context)
                    with self.lock:
                        for nid, n in self.dag.nodes.items():
                            if nid != next_id and nid in self.dag.get_dependents(node.id):
                                n.status = NodeStatus.SKIPPED
                    node.output = {"selected_branch": next_id}
                    node.status = NodeStatus.SUCCESS
                    return
                handler = self.handlers.get(node.handler)
                if handler is None:
                    raise ValueError(f"未找到处理器: {node.handler}")
                params = self.resolver.resolve_params(node.params, context)
                node.output = handler(**params)
                node.status = NodeStatus.SUCCESS
                return
            except Exception as e:
                last_error = e
                if attempt < node.retry.max_attempts - 1:
                    delay = self._calc_backoff(node.retry, attempt)
                    time.sleep(delay / 1000.0)
        node.status = NodeStatus.FAILED
        node.error = str(last_error)
    
    def _handle_node_completion(self, node: Node, context: WorkflowContext):
        node.finished_at = time.time()
        if node.status == NodeStatus.SUCCESS:
            with self.lock:
                context.node_outputs[node.id] = node.output
            self._emit_event("node_success", node)
        elif node.status == NodeStatus.FAILED:
            self._emit_event("node_failed", node)
            self._handle_failure(node, context)
    
    def _handle_failure(self, node: Node, context: WorkflowContext):
        """失败处理:标记后续节点为SKIPPED,触发补偿"""
        with self.lock:
            to_skip = set()
            queue = deque(self.dag.get_dependents(node.id))
            while queue:
                nid = queue.popleft()
                if nid not in to_skip:
                    to_skip.add(nid)
                    queue.extend(self.dag.get_dependents(nid))
            for nid in to_skip:
                n = self.dag.nodes[nid]
                if n.status == NodeStatus.PENDING:
                    n.status = NodeStatus.SKIPPED
        self._run_compensation(node, context)
    
    def _run_compensation(self, failed_node: Node, context: WorkflowContext):
        """补偿机制:逆序执行已完成节点的补偿操作"""
        completed = [
            n for n in self.dag.nodes.values()
            if n.status == NodeStatus.SUCCESS and n.id != failed_node.id
        ]
        completed.sort(key=lambda n: n.started_at or 0, reverse=True)
        for node in completed:
            handler = self.handlers.get(node.handler)
            if handler and hasattr(handler, 'compensate'):
                try:
                    node.status = NodeStatus.COMPENSATING
                    self._emit_event("compensation_started", node)
                    handler.compensate(node.output, context.node_outputs)
                    self._emit_event("compensation_done", node)
                except Exception as e:
                    self._emit_event("compensation_failed", node)
    
    def _calc_backoff(self, retry: RetryPolicy, attempt: int) -> float:
        if retry.backoff == "exponential":
            return retry.base_delay * (2 ** attempt)
        elif retry.backoff == "linear":
            return retry.base_delay * (attempt + 1)
        return retry.base_delay

上面的WorkflowEngine类是整个引擎的核心。它使用线程池实现并行执行,通过锁保证线程安全。执行循环的逻辑是:不断查询就绪节点,提交到线程池执行,等待完成的future,处理节点完成事件。当节点失败时,标记所有后续节点为SKIPPED,并逆序执行已完成节点的补偿操作。

六、失败补偿机制

失败补偿是工作流引擎保证最终一致性的关键机制。当流程中某个节点失败且无法重试成功时,之前已经成功执行的节点可能产生了副作用(如创建了文件、发送了消息、修改了数据库),这些副作用需要被撤销。

补偿机制的设计思路是:每个任务处理器除了正常的执行方法外,还可以提供一个compensate方法。当流程失败时,引擎按照执行顺序的逆序调用已完成节点的compensate方法,逐步撤销副作用。

补偿操作本身也可能失败,因此需要记录补偿状态,支持后续手动干预。引擎将补偿结果写入检查点,如果补偿过程中引擎崩溃,重启后可以从检查点恢复继续补偿。

class CheckpointManager:
    """检查点管理器:持久化工作流状态以支持断点续跑"""
    
    def __init__(self, storage_path: str):
        self.storage_path = storage_path
        self._ensure_dir()
    
    def _ensure_dir(self):
        import os
        os.makedirs(self.storage_path, exist_ok=True)
    
    def save(self, context: WorkflowContext, dag: DAGBuilder):
        import json
        import os
        checkpoint = {
            "workflow_id": context.workflow_id,
            "inputs": context.inputs,
            "node_outputs": self._serialize(context.node_outputs),
            "nodes": {
                nid: {
                    "status": n.status.value,
                    "output": self._serialize(n.output),
                    "error": n.error,
                    "attempts": n.attempts,
                    "started_at": n.started_at,
                    "finished_at": n.finished_at,
                }
                for nid, n in dag.nodes.items()
            },
            "timestamp": time.time(),
        }
        filepath = os.path.join(self.storage_path, f"{context.workflow_id}.json")
        with open(filepath, 'w') as f:
            json.dump(checkpoint, f, ensure_ascii=False, indent=2)
    
    def load(self, workflow_id: str) -> Optional[Dict[str, Any]]:
        import json
        import os
        filepath = os.path.join(self.storage_path, f"{workflow_id}.json")
        if not os.path.exists(filepath):
            return None
        with open(filepath, 'r') as f:
            return json.load(f)
    
    def _serialize(self, obj: Any) -> Any:
        if obj is None:
            return None
        if isinstance(obj, (str, int, float, bool)):
            return obj
        if isinstance(obj, dict):
            return {k: self._serialize(v) for k, v in obj.items()}
        if isinstance(obj, list):
            return [self._serialize(v) for v in obj]
        return str(obj)
    
    def restore(self, workflow_id: str, dag: DAGBuilder) -> Optional[WorkflowContext]:
        """从检查点恢复工作流状态"""
        data = self.load(workflow_id)
        if data is None:
            return None
        context = WorkflowContext(
            workflow_id=workflow_id,
            inputs=data.get("inputs", {}),
            node_outputs=data.get("node_outputs", {}),
        )
        for nid, state in data.get("nodes", {}).items():
            if nid in dag.nodes:
                node = dag.nodes[nid]
                node.status = NodeStatus(state["status"])
                node.output = state.get("output")
                node.error = state.get("error")
                node.attempts = state.get("attempts", 0)
                node.started_at = state.get("started_at")
                node.finished_at = state.get("finished_at")
        return context


class WorkflowMonitor:
    """工作流监控器:记录执行日志和指标"""
    
    def __init__(self):
        self.events: List[Dict[str, Any]] = []
        self.lock = threading.Lock()
    
    def on_event(self, event_type: str, node: Node):
        with self.lock:
            self.events.append({
                "type": event_type,
                "node_id": node.id,
                "timestamp": time.time(),
                "status": node.status.value,
                "attempts": node.attempts,
                "error": node.error,
            })
    
    def get_summary(self) -> Dict[str, Any]:
        with self.lock:
            total = len(set(e["node_id"] for e in self.events))
            succeeded = sum(1 for e in self.events if e["type"] == "node_success")
            failed = sum(1 for e in self.events if e["type"] == "node_failed")
            compensations = sum(1 for e in self.events if "compensation" in e["type"])
            return {
                "total_nodes": total,
                "succeeded": succeeded,
                "failed": failed,
                "compensations": compensations,
                "events": list(self.events),
            }

检查点管理器CheckpointManager负责将工作流状态序列化到磁盘,支持断点续跑。监控器WorkflowMonitor作为事件监听器接入引擎,记录所有节点事件的日志,提供执行摘要统计。

七、完整示例:构建一个数据处理工作流

下面用一个完整的示例展示如何使用上述组件构建和执行一个工作流:

# 定义任务处理器
class DataFetcher:
    def __call__(self, source: str):
        print(f"从 {source} 获取数据...")
        return {"raw_data": [1, 2, 3, 4, 5], "count": 5, "quality_score": 0.85}
    
    def compensate(self, output, context):
        print(f"补偿:清理已获取的数据 {output}")

class DataValidator:
    def __call__(self, data):
        print(f"校验数据: {data}")
        return {"validated": True, "quality_score": data.get("quality_score", 0)}

class DataProcessor:
    def __call__(self, data, mode):
        print(f"以 {mode} 模式处理数据: {data}")
        return {"processed": True, "result": [x * 2 for x in data.get("raw_data", [])]}

class ReportGenerator:
    def __call__(self, result):
        print(f"生成报告: {result}")
        return {"report_url": "/reports/123", "summary": "处理完成"}

class RejectionHandler:
    def __call__(self, reason):
        print(f"拒绝处理,原因: {reason}")
        return {"rejected": True}

# 构建并执行工作流
yaml_spec = """
name: demo
inputs:
  - name: data_source
    type: string
    required: true
  - name: analysis_mode
    type: string
    default: full
nodes:
  - id: fetch_data
    type: task
    handler: DataFetcher
    params:
      source: "${inputs.data_source}"
    retry:
      max_attempts: 3
      backoff: exponential
      base_delay: 500
  - id: validate_data
    type: task
    handler: DataValidator
    depends_on: [fetch_data]
    params:
      data: "${fetch_data.output}"
  - id: check_quality
    type: condition
    depends_on: [validate_data]
    branches:
      - when: "${validate_data.output.quality_score} > 0.8"
        next: process_data
      - default: reject_data
  - id: process_data
    type: task
    handler: DataProcessor
    depends_on: [check_quality]
    params:
      data: "${fetch_data.output}"
      mode: "${inputs.analysis_mode}"
  - id: generate_report
    type: task
    handler: ReportGenerator
    depends_on: [process_data]
    params:
      result: "${process_data.output}"
  - id: reject_data
    type: task
    handler: RejectionHandler
    depends_on: [check_quality]
    params:
      reason: "质量分数过低"
outputs:
  report: "${generate_report.output}"
"""

handlers = {
    "DataFetcher": DataFetcher(),
    "DataValidator": DataValidator(),
    "DataProcessor": DataProcessor(),
    "ReportGenerator": ReportGenerator(),
    "RejectionHandler": RejectionHandler(),
}

dag = DAGBuilder().from_yaml(yaml_spec)
engine = WorkflowEngine(dag, handlers, max_workers=4)
monitor = WorkflowMonitor()
engine.add_listener(monitor)
context = engine.run(inputs={"data_source": "database://test", "analysis_mode": "full"})
summary = monitor.get_summary()
print(f"执行摘要: 成功{summary['succeeded']}个, 失败{summary['failed']}个")

这个示例定义了一个完整的数据处理工作流,包含数据获取(带重试)、校验、条件分支(质量分数判断)、处理、报告生成等节点。引擎会自动并行执行无依赖的节点,处理条件分支,在失败时执行补偿操作。

八、性能优化与扩展方向

工作流引擎在实际生产中还需要考虑以下优化方向:

分布式执行。当单机的线程池无法满足吞吐量需求时,可以将任务分发到多个worker节点执行。这需要一个任务队列(如Redis、Kafka)来分发任务,各worker从队列消费任务并回报结果。引擎的主节点负责调度和状态管理,worker节点只负责执行。

动态DAG。某些场景下,工作流的结构需要在运行时动态调整——例如根据前置节点的输出决定是否添加新的任务节点。这需要在DSL中支持动态节点生成,引擎在执行过程中动态修改DAG结构。

超时控制。每个任务节点应支持超时配置,超时后自动取消执行并触发重试或补偿。Python中可以通过concurrent.futures.Future.result(timeout=...)实现,但需要注意线程的中断问题——Python的线程无法被强制中断,需要任务处理器内部配合检查取消标志。

可视化与调试。将DAG渲染为图形化界面,实时显示每个节点的执行状态、耗时、输入输出,对于调试和监控非常有价值。可以将DAG导出为Graphviz格式生成流程图,或集成到Web界面中实时展示。

版本管理与灰度发布。工作流定义应该支持版本管理,新旧版本可以共存。通过灰度发布策略,逐步将流量切换到新版本工作流,降低变更风险。

九、总结

工作流编排引擎是Agent系统的核心基础设施。本文从DSL设计出发,详细介绍了DAG构建、条件分支、并行执行、失败补偿等关键机制的实现。核心设计思路是:用声明式DSL描述流程,用DAG管理依赖关系,用线程池实现并行调度,用补偿机制保证最终一致性,用检查点支持断点续跑。

这套引擎的设计遵循了关注点分离原则——流程定义与执行逻辑分离,调度策略与业务逻辑分离,正常执行与容错恢复分离。这种分离使得每个部分可以独立演进和测试,也使得引擎可以适配不同的业务场景。

在实际项目中,可以根据具体需求对引擎进行裁剪或扩展。例如,简单的场景可以去掉补偿机制和检查点,只保留DAG调度;复杂的分布式场景可以增加任务队列和分布式锁。关键是要理解每个机制解决的问题和适用场景,做到按需取用。

【声明】本内容来自华为云开发者社区博主,不代表华为云及华为云开发者社区的观点和立场。转载时必须标注文章的来源(华为云社区)、文章链接、文章作者等基本信息,否则作者和本社区有权追究责任。如果您发现本社区中有涉嫌抄袭的内容,欢迎发送邮件进行举报,并提供相关证据,一经查实,本社区将立刻删除涉嫌侵权内容,举报邮箱: cloudbbs@huaweicloud.com
  • 点赞
  • 收藏
  • 关注作者

评论(0

0/1000
抱歉,系统识别当前为高风险访问,暂不支持该操作

全部回复

上滑加载中

设置昵称

在此一键设置昵称,即可参与社区互动!

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。

*长度不超过10个汉字或20个英文字符,设置后3个月内不可修改。