ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

转型实战项目七:从零实现一个分布式多 Agent 协作工作流引擎

转型实战项目七:从零实现一个分布式多 Agent 协作工作流引擎 转型实战项目七从零实现一个分布式多 Agent 协作工作流引擎在传统后端工程师转型为 AI 智能体架构师的进阶征程中“不依赖任何现成开源框架如 LangChain / AutoGen / CrewAI纯手工从零实现一个轻量级、分布式、基于有向无环图DAG的多 Agent 异步协作工作流引擎”是检验你对多智能体拓扑编排、异步协程、状态传递与并发容错掌握深度的终极硬核毕业攻坚项目。通过亲手编写这个引擎你将深刻洞悉任务有向无环图DAG的拓扑排序Topological Sort与依赖解析算法如何利用 Pythonasyncio实现无依赖子任务的极限并行并发执行节点产物Artifacts在多 Agent 之间的安全类型流转与状态共享。本文将带领大家**“使用纯 Python 原生协程从零手写一个生产级、跨节点异步协同的多 Agent 工作流引擎完整核心源码”**。一、轻量级分布式多 Agent 工作流引擎架构全景拓扑[ 用户提交多步骤业务目标 (自动编译为 Task DAG) ] │ ▼ ┌────────────────────────────────────────────────────────┐ │ 分布式多 Agent 协作工作流引擎核心中枢 │ ├────────────────────────────────────────────────────────┤ │ ├── 1. 拓扑解析器 (DAG Topological Resolver): │ │ │ • 解析各节点依赖关系: Step_3 依赖 Step_1 与 Step_2│ │ ├── 2. 异步协程调度池 (Async Task Dispatcher): │ │ │ • 发现 Step_1 与 Step_2 无相互依赖 ──►【并发抢跑!】│ │ └── 3. 全局产物上下文总线 (Artifacts Shared Bus) │ └───────────────────────┬────────────────────────────────┘ │ ┌──────────────┴──────────────┐ ▼ (并发并行执行) ▼ (并发并行执行) ┌─────────────────┐ ┌─────────────────┐ │ Task 1: 市场调研│ │ Task 2: 财务核算│ │ (Researcher) │ │ (Quant Engine) │ └────────┬────────┘ └────────┬────────┘ │ (产出 Artifact A) │ (产出 Artifact B) └──────────────┬──────────────┘ │ (依赖全部就绪) ▼ ┌────────────────────────────────────────────────────────┐ │ Task 3: 终态战略研报撰写 (Master Executive Writer) │ │ 聚合 Artifact A B ──► 生成最终交付方案! │ └────────────────────────────────────────────────────────┘二、从零纯手写的分布式多 Agent 工作流引擎完整实现实操创建micro_agent_workflow_engine.pyimport asyncio import time from typing import Dict, Any, List, Set, Callable from pydantic import BaseModel, Field class WorkflowTaskNode(BaseModel): task_id: str role_name: str action_handler: Callable[[Dict[str, Any]], Any] depends_on: Set[str] Field(default_factoryset) class Config: arbitrary_types_allowed True class MicroAgentWorkflowEngine: def __init__(self): self.nodes: Dict[str, WorkflowTaskNode] {} self.artifacts_bus: Dict[str, Any] {} # 共享产物总线 def add_node(self, task_id: str, role_name: str, handler: Callable[[dict], Any], depends_on: List[str] None): deps set(depends_on) if depends_on else set() node WorkflowTaskNode( task_idtask_id, role_namerole_name, action_handlerhandler, depends_ondeps ) self.nodes[task_id] node return self async def execute_workflow_async(self, initial_input: dict) - Dict[str, Any]: print(f 【启动多 Agent 协作工作流 】总节点数: {len(self.nodes)}) self.artifacts_bus initial_input.copy() completed_tasks: Set[str] set() pending_nodes self.nodes.copy() t0 time.time() # 循环推进直到所有 DAG 节点执行完毕 while pending_nodes: # 1. 寻找当前所有依赖已满足的可执行就绪节点 (Ready Nodes) ready_tasks: List[WorkflowTaskNode] [] for t_id, node in list(pending_nodes.items()): if node.depends_on.issubset(completed_tasks): ready_tasks.append(node) del pending_nodes[t_id] if not ready_tasks: raise RuntimeError( 【检测到 DAG 循环依赖死锁或无法解析的依赖节点】) print(f\n⚡ [调度并发执行] 当前就绪可并发执行的节点: {[t.task_id for t in ready_tasks]}) # 2. 并发执行当前层的所有就绪节点 (Asyncio Gather 并行抢跑!) async def _run_single_node(n: WorkflowTaskNode): print(f ▶ [{n.role_name}] 开始执行任务 [{n.task_id}]...) # 注入前置依赖产物 res await n.action_handler(self.artifacts_bus) return n.task_id, res results await asyncio.gather(*[_run_single_node(n) for n in ready_tasks]) # 3. 收集产物并更新状态 for t_id, output in results: self.artifacts_bus[t_id] output completed_tasks.add(t_id) print(f ✅ [{t_id}] 任务圆满完成并回填产物。) elapsed_ms int((time.time() - t0) * 1000) print(f\n 【工作流全链路执行成功 ✅】总耗时仅: {elapsed_ms}ms) return self.artifacts_bus三、真实多 Agent 协作业务演练实操# 1. 定义 3 个独立的业务 Agent 协程函数 async def researcher_agent(bus: dict) - str: await asyncio.sleep(0.3) # 模拟调研网络 IO return 市场情报: 2026 年新能源渗透率已突破 55% async def financial_quant_agent(bus: dict) - dict: await asyncio.sleep(0.3) # 模拟财务量化测算 return {estimated_roi: 3.8, risk_index: LOW} async def executive_writer_agent(bus: dict) - str: # 依赖前两个任务的产物 market_info bus[task_research] finance_info bus[task_finance] return f【战略内参】结合【{market_info}】与财务预测【ROI: {finance_info[estimated_roi]}】建议全力加大投入 # 2. 组装 DAG 工作流 async def main(): engine MicroAgentWorkflowEngine() engine.add_node(task_research, 市场调研专家, researcher_agent)\ .add_node(task_finance, 财务量化专家, financial_quant_agent)\ .add_node(task_report, 战略主笔专家, executive_writer_agent, depends_on[task_research, task_finance]) # 3. 启动执行 final_bus await engine.execute_workflow_async({company: 头部车企}) print(f\n 最终交付成果:\n{final_bus[task_report]}) # asyncio.run(main())四、写在最后通过纯手写实现这套不到 100 行的现代化异步多 Agent 工作流引擎你彻底摆脱了对笨重庞大第三方框架黑盒的恐惧你真正掌握了现代多智能体系统在并发调度、DAG 解析与状态流转方面的底层核心源码机理具备了在极端严苛场景下为企业自研定制高可用、低开销专用 AI 编排中枢的顶级硬核实力
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进