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

相关新闻

Subdominator健康检查与故障排除:确保API密钥与数据源可用性

Subdominator健康检查与故障排除:确保API密钥与数据源可用性

Subdominator健康检查与故障排除:确保API密钥与数据源可用性 【免费下载链接】Subdominator SubDominator helps you discover subdomains associated with a target domain efficiently and with minimal impact for your Bug Bounty 项目地址: https://gitcode.…

2026/7/27 14:10:29 阅读更多 →
嵌入式I2C与1-Wire总线实战:从寄存器配置到中断与DMA优化

嵌入式I2C与1-Wire总线实战:从寄存器配置到中断与DMA优化

1. 项目概述与总线协议核心价值在嵌入式开发的日常里,我们总在和各种各样的外设打交道,传感器、EEPROM、实时时钟、身份认证芯片……它们就像一个个沉默的“小弟”,等着主控MCU这个“老大”发号施令。而“老大”和“小弟”之间说悄悄话的通道…

2026/7/27 14:09:28 阅读更多 →
5分钟快速部署ChatTTS-ui:本地语音合成解决方案完全指南

5分钟快速部署ChatTTS-ui:本地语音合成解决方案完全指南

5分钟快速部署ChatTTS-ui:本地语音合成解决方案完全指南 【免费下载链接】ChatTTS-ui 一个简单的本地网页界面,使用ChatTTS将文字合成为语音,同时支持对外提供API接口。A simple native web interface that uses ChatTTS to synthesize text …

2026/7/27 14:09:28 阅读更多 →

最新新闻

【AI邮件撰写黄金法则】:20年资深IT专家亲授7大高转化率写作心法

【AI邮件撰写黄金法则】:20年资深IT专家亲授7大高转化率写作心法

更多请点击: https://kaifayun.com 第一章:AI邮件撰写的核心认知与底层逻辑 AI邮件撰写并非简单地将自然语言生成(NLG)模型套用于模板填充,其本质是语义理解、意图建模与上下文协同的复合过程。真正的效能提升来源于对…

2026/7/27 14:31:39 阅读更多 →
166、自动对焦(AF)技术全览:CDAF、PDAF、双像素对焦与深度学习AF的实战对比

166、自动对焦(AF)技术全览:CDAF、PDAF、双像素对焦与深度学习AF的实战对比

166、自动对焦(AF)技术全览:CDAF、PDAF、双像素对焦与深度学习AF的实战对比 一个让我熬夜三天的对焦问题 2019年某款旗舰机项目,客户反馈暗光下拍照总是“拉风箱”——对焦来回抽动,最终出片模糊。我盯着log看了三天,发现CDAF在低照度场景下对比度信号几乎淹没在噪声里,…

2026/7/27 14:31:39 阅读更多 →
ADC117x评估模块实战指南:从硬件连接到性能测试优化

ADC117x评估模块实战指南:从硬件连接到性能测试优化

1. 项目概述:为什么需要评估模块?在嵌入式系统、通信设备或者任何需要处理模拟信号的数字系统中,模数转换器(ADC)的性能往往是决定整个系统“天花板”的关键。你选了一颗参数看起来不错的ADC芯片,但把它焊到…

2026/7/27 14:31:39 阅读更多 →
深度学习与计算机视觉:从原理到工业应用实践

深度学习与计算机视觉:从原理到工业应用实践

1. 深度学习技术概述:从实验室到产业落地的跨越 深度学习作为机器学习的重要分支,其核心在于构建多层神经网络来模拟人脑的学习机制。与传统的机器学习方法相比,深度学习最大的突破在于能够自动从原始数据中提取多层次的特征表示,…

2026/7/27 14:31:39 阅读更多 →
QuaterNet数据预处理全解析:从原始动作捕捉数据到四元数序列转换

QuaterNet数据预处理全解析:从原始动作捕捉数据到四元数序列转换

QuaterNet数据预处理全解析:从原始动作捕捉数据到四元数序列转换 【免费下载链接】QuaterNet Proposes neural networks that can generate animation of virtual characters for different actions. 项目地址: https://gitcode.com/gh_mirrors/qu/QuaterNet …

2026/7/27 14:31:39 阅读更多 →
PLM在精细化工领域的布局:2026年主要厂商技术特色

PLM在精细化工领域的布局:2026年主要厂商技术特色

一、精细化工PLM赛道:技术迭代特征与整体行业格局近半年,国内精细化工数字化转型迈入深耕落地阶段,传统通用型PLM仅能实现基础文件、物料管控,已无法适配行业配方保密、批次差异化、合规严苛、产品快速迭代的专属特性。2026年精细…

2026/7/27 14:30:38 阅读更多 →

日新闻

【JAVA毕设源码分享】基于SpringBoot的社区智能垃圾管理系统的设计与实现(程序+文档+代码讲解+一条龙定制)

【JAVA毕设源码分享】基于SpringBoot的社区智能垃圾管理系统的设计与实现(程序+文档+代码讲解+一条龙定制)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围:&am…

2026/7/27 0:00:54 阅读更多 →
SPI实战指南:从时钟模式到寄存器配置,解决嵌入式通信难题

SPI实战指南:从时钟模式到寄存器配置,解决嵌入式通信难题

1. 项目概述:从寄存器手册到实战指南 如果你手头有一份类似德州仪器(TI)TMS320x240xA系列DSP的SPI模块技术手册,看着里面密密麻麻的寄存器位定义、时序图和公式,是不是感觉头大?这份资料虽然权威&#xff0…

2026/7/27 0:00:54 阅读更多 →
【JAVA毕设源码分享】基于springboot的水果购物管理系统的设计与实现(程序+文档+代码讲解+一条龙定制)

【JAVA毕设源码分享】基于springboot的水果购物管理系统的设计与实现(程序+文档+代码讲解+一条龙定制)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围:&am…

2026/7/27 0:00:54 阅读更多 →

周新闻

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 数据集6000张 完整源码已标注数据集训练好的模型环境配置教程程序运行说明文档,可以直接使用!系统支持图片、视频、摄像头等多种方式检测裂缝,功能强大实用。 1数据集6000张 8各类别

2026/7/27 4:33:59 阅读更多 →
深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

pubg数据集 精选原图1.42万数据 1.49万标签 无任何重复、算法增强或冗余图像! pubg绝地求生目标检测数据集 1分类:e_body,14905个标签,txt格式 共计14244张图,99%为640*640尺寸图像 适合yolo目标检测、AI训练关键词&am…

2026/7/27 6:31:56 阅读更多 →
Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex检测数据集数据集详情检测类别: allies enemy tag图片总量:7247张训练集:5139张验证集:1425张测试集:683张标注状态:全部已标注,即拿即用数据格式:支持YOLO格式及其他格式&#…

2026/7/27 4:01:12 阅读更多 →

月新闻