工作流平台的架构演进全记录:从MVP到企业级的五次重大重构
构建一个支撑企业级Agent的工作流平台,是过去一年技术工作的核心。从最初200行的Python脚本到现在数万行的分布式系统,经历了五次重大架构重构。每一次重构都源于对系统瓶颈的深刻认知和对业务需求的前瞻判断。本文还原这五次重构的关键决策和技术细节。
一、引言
工作流引擎是Agent产品的核心基础设施。它负责编排LLM调用、工具调用、条件判断和人工审批等环节,形成可执行的业务工作流。一个合格的工作流平台需要满足三个核心要求:高可靠性(工作流不能丢)、高扩展性(支持自定义节点类型)和高性能(端到端延迟可控)。
项目从去年7月的MVP版本起步,到今年6月演进为企业级平台,经历了单进程脚本、异步任务队列、微服务拆分、事件驱动架构、多租户隔离五次重构。每次重构都解决了前一个版本的瓶颈,但也引入了新的复杂度。以下是完整的技术演进记录。
二、原理:工作流引擎的核心抽象
在讨论具体架构之前,先定义工作流引擎的核心抽象。一个通用工作流平台包含以下关键概念:
核心设计原则:
- 状态与执行分离:工作流的状态持久化在外部存储中,执行器是无状态的。这样任意执行器宕机不会丢失工作流状态。
- 节点可扩展:通过插件机制支持自定义节点类型,包括LLM调用、HTTP请求、代码执行、人工审批等。
- 事件驱动:工作流之间的依赖通过事件总线解耦,避免同步等待造成的资源浪费。
- 幂等执行:每个节点的执行必须支持重试且不产生副作用,这是分布式环境下可靠性的基础保证。
三、代码:第五版架构核心实现
以下是第五次重构后的核心工作流引擎实现,采用事件驱动架构:
import asyncio import json import logging from abc import ABC, abstractmethod from dataclasses import dataclass, field from datetime import datetime from enum import Enum from typing import Any, Callable, Dict, List, Optional from uuid import uuid4 logging.basicConfig(level=logging.INFO) logger = logging.getLogger(__name__) class NodeType(Enum): LLM = "llm_call" HTTP = "http_request" CODE = "code_execution" CONDITION = "condition" APPROVAL = "human_approval" PARALLEL = "parallel_fork" class WorkflowStatus(Enum): PENDING = "pending" RUNNING = "running" SUSPENDED = "suspended" COMPLETED = "completed" FAILED = "failed" class NodeStatus(Enum): IDLE = "idle" EXECUTING = "executing" SUCCEEDED = "succeeded" FAILED = "failed" SKIPPED = "skipped" @dataclass class ExecutionContext: """工作流执行上下文""" workflow_id: str variables: Dict[str, Any] = field(default_factory=dict) node_results: Dict[str, Any] = field(default_factory=dict) metadata: Dict[str, Any] = field(default_factory=dict) def get_variable(self, key: str, default: Any = None) -> Any: return self.variables.get(key, default) def set_variable(self, key: str, value: Any) -> None: self.variables[key] = value class StateStore(ABC): """状态存储抽象接口""" @abstractmethod async def save_workflow_state( self, workflow_id: str, state: Dict ) -> None: pass @abstractmethod async def load_workflow_state( self, workflow_id: str ) -> Optional[Dict]: pass @abstractmethod async def save_node_result( self, workflow_id: str, node_id: str, result: Dict ) -> None: pass class InMemoryStateStore(StateStore): """内存状态存储实现""" def __init__(self): self._store: Dict[str, Dict] = {} async def save_workflow_state( self, workflow_id: str, state: Dict ) -> None: self._store[workflow_id] = state async def load_workflow_state( self, workflow_id: str ) -> Optional[Dict]: return self._store.get(workflow_id) async def save_node_result( self, workflow_id: str, node_id: str, result: Dict ) -> None: key = f"{workflow_id}:{node_id}" self._store[key] = result class NodeExecutor(ABC): """节点执行器基类""" def __init__(self, max_retries: int = 3): self.max_retries = max_retries @abstractmethod async def execute( self, context: ExecutionContext, config: Dict ) -> Dict: pass async def execute_with_retry( self, context: ExecutionContext, config: Dict ) -> Dict: """带重试的执行逻辑""" last_error = None for attempt in range(1, self.max_retries + 1): try: result = await self.execute(context, config) logger.info(f"节点执行成功, 尝试次数: {attempt}") return result except Exception as e: last_error = e logger.warning( f"节点执行失败 (第{attempt}次): {e}" ) if attempt < self.max_retries: await asyncio.sleep(2 ** attempt) raise RuntimeError( f"节点执行失败, 已重试{self.max_retries}次: {last_error}" ) class WorkflowEngine: """工作流引擎核心""" def __init__(self, state_store: StateStore): self.state_store = state_store self.executors: Dict[NodeType, NodeExecutor] = {} self._event_handlers: Dict[str, List[Callable]] = {} def register_executor( self, node_type: NodeType, executor: NodeExecutor ) -> None: """注册节点执行器""" self.executors[node_type] = executor def on( self, event: str, handler: Callable ) -> None: """注册事件处理器""" if event not in self._event_handlers: self._event_handlers[event] = [] self._event_handlers[event].append(handler) async def _emit_event( self, event: str, data: Dict ) -> None: """触发事件""" handlers = self._event_handlers.get(event, []) tasks = [handler(data) for handler in handlers] if tasks: await asyncio.gather(*tasks) async def execute_workflow( self, workflow_def: Dict, initial_vars: Optional[Dict] = None ) -> ExecutionContext: """执行工作流""" workflow_id = uuid4().hex context = ExecutionContext( workflow_id=workflow_id, variables=initial_vars or {} ) await self.state_store.save_workflow_state(workflow_id, { 'status': WorkflowStatus.RUNNING.value, 'started_at': datetime.now().isoformat() }) await self._emit_event('workflow.started', { 'workflow_id': workflow_id }) try: nodes = workflow_def.get('nodes', []) for node in nodes: node_id = node['id'] node_type = NodeType(node['type']) config = node.get('config', {}) executor = self.executors.get(node_type) if not executor: raise ValueError( f"未注册的执行器类型: {node_type}" ) result = await executor.execute_with_retry( context, config ) context.node_results[node_id] = result await self.state_store.save_node_result( workflow_id, node_id, result ) await self.state_store.save_workflow_state(workflow_id, { 'status': WorkflowStatus.COMPLETED.value, 'completed_at': datetime.now().isoformat() }) await self._emit_event('workflow.completed', { 'workflow_id': workflow_id, 'node_count': len(nodes) }) except Exception as e: await self.state_store.save_workflow_state(workflow_id, { 'status': WorkflowStatus.FAILED.value, 'error': str(e), 'failed_at': datetime.now().isoformat() }) logger.error(f"工作流执行失败 {workflow_id}: {e}") raise return context # 使用示例 async def main(): engine = WorkflowEngine(InMemoryStateStore()) # 注册事件处理器 async def on_completed(data: Dict): logger.info(f"工作流完成: {data['workflow_id']}") engine.on('workflow.completed', on_completed) # 定义并执行工作流 workflow_def = { 'nodes': [ { 'id': 'node_1', 'type': 'llm_call', 'config': {'prompt': '分析用户输入'} } ] } try: context = await engine.execute_workflow(workflow_def) logger.info(f"执行结果: {context.node_results}") except Exception as e: logger.error(f"工作流执行失败: {e}") if __name__ == "__main__": asyncio.run(main())四、五次重构的关键权衡
| 版本 | 架构模式 | 核心问题 | 重构动机 | 收益 |
|---|---|---|---|---|
| V1 | 单进程同步 | 阻塞主线程 | 无法并行处理 | — |
| V2 | Celery异步 | 任务积压 | 峰值QPS不足 | 吞吐量提升5x |
| V3 | 微服务拆分 | 服务间耦合 | 部署粒度问题 | 独立扩缩容 |
| V4 | 事件驱动 | 事件溯源复杂 | 跨服务编排 | 解耦80%依赖 |
| V5 | 多租户隔离 | 租户数据隔离 | 企业客户需求 | 支持SaaS化 |
每次重构的决策依据:
- V1→V2:当单日工作流执行量超过1000条时,同步模式开始出现超时。
- V2→V3:当需要独立升级LLM调用服务而不影响其他模块时,微服务拆分成为必然。
- V3→V4:当跨工作流的依赖关系越来越复杂时,同步RPC调用的链式失败问题严重。
- V4→V5:当第一个企业客户要求数据物理隔离时,多租户架构正式提上日程。
仍在讨论的开放问题:
- 是否需要引入工作流定义DSL,还是继续使用JSON/YAML配置?
- 状态存储从Redis迁移到PostgreSQL的时机和风险评估?
- 是否引入Saga模式处理分布式事务补偿?
五、总结
工作流平台的五次重构反映了创业项目中技术架构演进的典型路径:从简单够用到逐步复杂化,每一次重构都是对业务需求变化的响应。核心原则始终未变:保持状态与执行分离、保证节点执行的幂等性、坚持通过事件解耦服务依赖。下一步的重点是完善可观测性(分布式追踪和业务监控)以及工作流的可视化编排能力。