
工作流平台的架构演进全记录从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(levellogging.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_factorydict) node_results: Dict[str, Any] field(default_factorydict) metadata: Dict[str, Any] field(default_factorydict) 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_idworkflow_id, variablesinitial_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单进程同步阻塞主线程无法并行处理—V2Celery异步任务积压峰值QPS不足吞吐量提升5xV3微服务拆分服务间耦合部署粒度问题独立扩缩容V4事件驱动事件溯源复杂跨服务编排解耦80%依赖V5多租户隔离租户数据隔离企业客户需求支持SaaS化每次重构的决策依据V1→V2当单日工作流执行量超过1000条时同步模式开始出现超时。V2→V3当需要独立升级LLM调用服务而不影响其他模块时微服务拆分成为必然。V3→V4当跨工作流的依赖关系越来越复杂时同步RPC调用的链式失败问题严重。V4→V5当第一个企业客户要求数据物理隔离时多租户架构正式提上日程。仍在讨论的开放问题是否需要引入工作流定义DSL还是继续使用JSON/YAML配置状态存储从Redis迁移到PostgreSQL的时机和风险评估是否引入Saga模式处理分布式事务补偿五、总结工作流平台的五次重构反映了创业项目中技术架构演进的典型路径从简单够用到逐步复杂化每一次重构都是对业务需求变化的响应。核心原则始终未变保持状态与执行分离、保证节点执行的幂等性、坚持通过事件解耦服务依赖。下一步的重点是完善可观测性分布式追踪和业务监控以及工作流的可视化编排能力。