news 2026/7/27 14:10:26

工作流平台的架构演进全记录:从MVP到企业级的五次重大重构

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
工作流平台的架构演进全记录:从MVP到企业级的五次重大重构

工作流平台的架构演进全记录:从MVP到企业级的五次重大重构

构建一个支撑企业级Agent的工作流平台,是过去一年技术工作的核心。从最初200行的Python脚本到现在数万行的分布式系统,经历了五次重大架构重构。每一次重构都源于对系统瓶颈的深刻认知和对业务需求的前瞻判断。本文还原这五次重构的关键决策和技术细节。

一、引言

工作流引擎是Agent产品的核心基础设施。它负责编排LLM调用、工具调用、条件判断和人工审批等环节,形成可执行的业务工作流。一个合格的工作流平台需要满足三个核心要求:高可靠性(工作流不能丢)、高扩展性(支持自定义节点类型)和高性能(端到端延迟可控)。

项目从去年7月的MVP版本起步,到今年6月演进为企业级平台,经历了单进程脚本、异步任务队列、微服务拆分、事件驱动架构、多租户隔离五次重构。每次重构都解决了前一个版本的瓶颈,但也引入了新的复杂度。以下是完整的技术演进记录。

二、原理:工作流引擎的核心抽象

在讨论具体架构之前,先定义工作流引擎的核心抽象。一个通用工作流平台包含以下关键概念:

核心设计原则:

  1. 状态与执行分离:工作流的状态持久化在外部存储中,执行器是无状态的。这样任意执行器宕机不会丢失工作流状态。
  2. 节点可扩展:通过插件机制支持自定义节点类型,包括LLM调用、HTTP请求、代码执行、人工审批等。
  3. 事件驱动:工作流之间的依赖通过事件总线解耦,避免同步等待造成的资源浪费。
  4. 幂等执行:每个节点的执行必须支持重试且不产生副作用,这是分布式环境下可靠性的基础保证。

三、代码:第五版架构核心实现

以下是第五次重构后的核心工作流引擎实现,采用事件驱动架构:

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单进程同步阻塞主线程无法并行处理
V2Celery异步任务积压峰值QPS不足吞吐量提升5x
V3微服务拆分服务间耦合部署粒度问题独立扩缩容
V4事件驱动事件溯源复杂跨服务编排解耦80%依赖
V5多租户隔离租户数据隔离企业客户需求支持SaaS化

每次重构的决策依据:

  • V1→V2:当单日工作流执行量超过1000条时,同步模式开始出现超时。
  • V2→V3:当需要独立升级LLM调用服务而不影响其他模块时,微服务拆分成为必然。
  • V3→V4:当跨工作流的依赖关系越来越复杂时,同步RPC调用的链式失败问题严重。
  • V4→V5:当第一个企业客户要求数据物理隔离时,多租户架构正式提上日程。

仍在讨论的开放问题:

  • 是否需要引入工作流定义DSL,还是继续使用JSON/YAML配置?
  • 状态存储从Redis迁移到PostgreSQL的时机和风险评估?
  • 是否引入Saga模式处理分布式事务补偿?

五、总结

工作流平台的五次重构反映了创业项目中技术架构演进的典型路径:从简单够用到逐步复杂化,每一次重构都是对业务需求变化的响应。核心原则始终未变:保持状态与执行分离、保证节点执行的幂等性、坚持通过事件解耦服务依赖。下一步的重点是完善可观测性(分布式追踪和业务监控)以及工作流的可视化编排能力。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/7/27 14:09:54

2026年算法工程师最新必问面试题一

深度学习核心原理(7 题) 本部分聚焦于深度学习模型的核心原理、架构细节及训练技巧,是算法工程师面试中的重中之重。 1. CNN 中的卷积核是做什么的?1x1 卷积的作用是什么? 卷积核的作用: 特征提取器:卷积核是一个可学习的权重矩阵,通过在输入数据(如图像)上滑动,…

作者头像 李华
网站建设 2026/7/27 14:09:39

UCD90124电源时序控制器:系统复位、看门狗与数据记录实战解析

1. 项目概述与核心价值在服务器主板、通信基站或者高端工控设备的研发过程中&#xff0c;我们这些硬件工程师最头疼的问题之一&#xff0c;就是如何让十几个甚至几十个不同电压的电源轨&#xff08;Rail&#xff09;像一支训练有素的军队一样&#xff0c;井然有序地启动和关闭。…

作者头像 李华
网站建设 2026/7/27 14:08:10

AI 大模型日报 — 2026年7月27日(周一)

&#x1f916; AI 大模型日报 — 2026年7月27日&#xff08;周一&#xff09;过去一周&#xff08;7月20日–27日&#xff09;全球AI大模型领域重磅消息频出&#xff1a;Anthropic发布Claude Opus 5、月之暗面Kimi K3引爆业界、Google推出Gemini 3.6 Flash、阿里发布Qwen3.8-Ma…

作者头像 李华
网站建设 2026/7/27 14:07:20

改进boxinst_r50_fpn模型实现高效数字识别与定位

1. 数字识别与定位任务&#xff1a;Det改进实现boxinst_r50_fpn_ms-90k_coco模型训练 1.1. 模型概述 boxinst_r50_fpn_ms-90k_coco模型是基于Detectron2框架改进的数字识别与定位模型。这个模型的核心创新点在于将传统的两阶段检测流程优化为端到端的解决方案&#xff0c;显著…

作者头像 李华
网站建设 2026/7/27 14:05:26

rxjs-spy完全指南:从安装到高级调试的10个实用技巧

rxjs-spy完全指南&#xff1a;从安装到高级调试的10个实用技巧 【免费下载链接】rxjs-spy A debugging library for RxJS 项目地址: https://gitcode.com/gh_mirrors/rx/rxjs-spy rxjs-spy是一款专为RxJS打造的调试库&#xff0c;它能帮助开发者轻松追踪、分析和调试RxJ…

作者头像 李华
网站建设 2026/7/27 14:04:31

10个必学的TinyGo Drivers高级特性:中断处理、DMA与低延迟设计

10个必学的TinyGo Drivers高级特性&#xff1a;中断处理、DMA与低延迟设计 【免费下载链接】drivers TinyGo drivers for sensors, displays, wireless adaptors, and other devices that use I2C, SPI, GPIO, ADC, and UART interfaces. 项目地址: https://gitcode.com/gh_m…

作者头像 李华