转型实战项目七:从零实现一个分布式多 Agent 协作工作流引擎
在传统后端工程师转型为 AI 智能体架构师的进阶征程中,“不依赖任何现成开源框架(如 LangChain / AutoGen / CrewAI),纯手工从零实现一个轻量级、分布式、基于有向无环图(DAG)的多 Agent 异步协作工作流引擎”是检验你对多智能体拓扑编排、异步协程、状态传递与并发容错掌握深度的终极硬核毕业攻坚项目。
通过亲手编写这个引擎,你将深刻洞悉:
- 任务有向无环图(DAG)的拓扑排序(Topological Sort)与依赖解析算法;
- 如何利用 Python
asyncio实现无依赖子任务的极限并行并发执行; - 节点产物(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.py:
import 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_factory=set) 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_id=task_id, role_name=role_name, action_handler=handler, depends_on=deps ) 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 编排中枢的顶级硬核实力!