
多智能体协作开发中真正的难点往往不在某个大模型能不能回答而在于如何把一件复杂任务拆分给多个子智能体、让它们各司其职再用一套可观测的调度机制把结果串联起来。DeepAgents 这类多智能体实践把模型调用、上下文管理、工具执行、任务队列、重试与结果回收统一封装成工程组件这套组件在社区里通常被称为 harness。本文不绑定某个闭源框架而是用原生 Python 从子智能体开始逐步实现一个支持异步任务编排的最小多智能体协作系统。读完可以直接照做也可以把这里的任务模型、并发模型和排查思路迁移到 LangGraph、AutoGen 或自己封装的 Agent 框架中。这套知识适合正在学习 AI 大模型应用开发、多智能体系统核心架构以及相关工程实践的开发者。全文会围绕一个完整案例展开自动生成三篇技术专题文章。每个专题由“资料查证 Agent”“写作 Agent”“审校 Agent”分工完成调度层负责并发执行、超时控制和失败重试。运行成功后你会看到多个子智能体如何被组织成一个可扩展的协作流程。1. 多智能体协作中的 DeepAgents 到底解决了什么问题1.1 一个任务拆给多个智能体问题从“模型能不能答”变成了“协作会不会断”先从一个常见场景说起。假设你要让大模型生成一篇技术文章比如“异步任务编排原理”。单 Agent 的做法是把整个需求一股脑丢给模型得到一个完整回答。它的缺点是明显的资料查证、结构设计、文字表达、事实校验全被压缩在一次模型调用中一旦内容变长模型很容易出现开头精彩、后面空洞或者把不准确的技术概念写得像真的一样。多智能体的做法是拆。把任务拆成多个子任务每个子任务交给一个职责单一的智能体子智能体核心职责接收输入输出结果资料查证 Agent检索方向确认、来源类型判断、待验证问题整理专题关键词结构化提纲写作 Agent根据提纲撰写初稿关键词 提纲Markdown 草稿审校 Agent检查事实表述、结构完整性、空泛表述关键词 草稿审校意见与修改建议拆分之后每个 Agent 的上下文更短、职责边界更清晰Prompt 也更容易维护。但问题立刻转移到另一端三个 Agent 之间如何传数据谁先执行如果其中一个失败怎么办多个专题任务同时进来如何控制并发这些都属于协作机制而不是模型能力问题。1.2 harness 到底在封装什么社区里出现的 DeepAgents、Codex Harness、DeepSeek Harness 等说法都把“大模型调用之外的工程代码”叫成 harness。它不是一个具体的插件而是指一套把 Agent 运行过程工程化的骨架。一个标准的 harness 至少需要管理下面这些环节模型 API 的异步调用与超时控制每个子智能体的 System Prompt 维护任务输入、中间结果、最终结果的传递并发控制、队列、重试和失败隔离日志、追踪和结果回收。如果没有 harness开发者的代码里会到处散落openai.ChatCompletion调用、手动拼接上下文字典、裸try except处理超时。这样的系统在单 Agent 演示时还能跑一旦进入异步编排几乎无法排查“任务到底卡在哪一步”。所以可以这样理解模型是大脑harness 是神经系统。大脑负责生成内容神经系统负责把任务送到正确的位置、把结果带回正确的地方。DeepAgents 实践的核心不是一个神秘框架而是把这条“神经系统”设计清楚。1.3 DeepAgents 的最小组成任务、智能体、编排器无论框架名称怎么变多智能体系统的最小模型都可以抽象成三个部分任务Task一次具体执行的描述。至少包含任务 ID、目标智能体、输入参数和重试次数。智能体Agent一个拥有角色、System Prompt、工具调用能力和输入处理逻辑的执行单元。编排器Coordinator负责任务投递、并发调度、超时控制、结果回收和异常重试。它们的运行流程可以简化为任务创建 - 队列投递 - worker 消费 - 子智能体执行 - 结果回收 - 进入下一阶段后面的代码会按照这个最小模型展开。先不要急着引入 Redis、Kafka、分布式链路追踪先把单机异步编排跑通。理解了这个模型生产环境只是把内存队列替换成消息队列、把内存状态替换成数据库状态而已。2. 环境准备与项目骨架2.1 运行环境与依赖建议使用 Python 3.10 及以上版本。原因很简单本文大量使用asyncio队列和asyncio.wait_for3.10 之后异步语法更稳定类型标注也更友好。依赖建议固定到以下范围openai1.30.0 pydantic2.5.0 PyYAML6.0.1安装命令mkdir deepagents-demo cd deepagents-demo python -m venv .venv source .venv/bin/activate pip install openai1.30.0 pydantic2.5.0 PyYAML6.0.1这里使用 OpenAI 官方 Python SDK但并不是必须使用 OpenAI 官方模型。SDK 支持通过base_url指向任何兼容 OpenAI 接口的服务本地部署模型、第三方模型网关都可以接入。学习阶段也可以先用一个响应速度快的轻量模型减少调试等待时间。注意如果你的模型 API 不支持流式返回或并发较高建议在测试时限制 worker 数量避免接口限流。2.2 项目目录与核心模块一个可运行的最小项目结构如下deepagents-demo/ ├── requirements.txt ├── config.yaml ├── main.py ├── llm_client.py ├── agent.py ├── jobs.py ├── coordinator.py └── outputs/各文件职责如下文件职责llm_client.py封装模型 API 异步调用agent.py定义 Agent 基类与子智能体jobs.py定义专题任务 Job 与内部多阶段流程coordinator.py实现任务队列、worker 池与重试main.py入口程序负责加载配置并启动config.yaml模型参数与编排参数配置outputs/结果输出目录这个结构没有强依赖任何具体框架。以后如果要迁移到 LangGraphagent.py中的子智能体可以直接对应到 LangGraph 的不同节点coordinator.py则对应到状态图之外的执行调度层。2.3 模型接口配置在config.yaml中配置模型和编排参数llm: base_url: https://api.openai.com/v1 api_key: ${OPENAI_API_KEY} model: gpt-4o-mini temperature: 0.3 max_tokens: 2048 workflow: max_workers: 3 max_retries: 2 per_task_timeout: 60 queue_size: 100 agents: research: system_prompt: 你负责资料查证。只输出查证结果不要写完整文章。要求输出3个可查证方向每个方向写明来源类型和需要验证的问题。 draft: system_prompt: 你负责根据研究提纲撰写技术文章初稿。使用 Markdown控制在800字以内结构要清晰。 review: system_prompt: 你负责审校文章。检查事实表述是否严谨、结构是否完整、是否存在空泛表述输出审校意见。这里的核心参数解释如下参数含义注意事项max_workers并发 worker 数调大提高吞吐但会增加模型 API 并发压力max_retries任务失败重试次数建议 1 到 3过大会放大故障时间per_task_timeout单个任务超时时间模型响应慢时容易触发需要按模型速度调整queue_size队列容量背压机制避免任务无限积压占内存temperature采样温度查证与审校建议偏低写作可略高api_key写成${OPENAI_API_KEY}只是一种示例约定实际读取环境变量即可。不要在仓库中提交真实密钥。3. 先实现子智能体角色、上下文与任务结构3.1 Agent 基类与任务结构子智能体需要一个统一接口否则编排器很难用一种方式去调度所有 Agent。先定义任务和结果的数据结构# agent.py from dataclasses import dataclass, field from typing import Any dataclass class Task: task_id: str target: str payload: dict field(default_factorydict) retry_count: int 0 dataclass class AgentResult: sender: str content: str meta: dict field(default_factorydict)注意retry_count字段。它会跟随任务在队列中流转当消费者捕获异常后判断是否允许再次入队。接着定义 Agent 基类# agent.py class BaseAgent: name base def __init__(self, llm, system_prompt: str): self.llm llm self.system_prompt system_prompt async def run(self, task: Task) - AgentResult: user_prompt self._build_user_prompt(task) content await self.llm.complete( systemself.system_prompt, useruser_prompt, ) return AgentResult( senderself.name, contentcontent, meta{task_id: task.task_id}, ) def _build_user_prompt(self, task: Task) - str: raise NotImplementedError这段代码解决的一个关键问题是模型调用被统一封装在llm.complete中子智能体只需要关心如何构造 User Prompt以及如何消费返回值。这样编排器和子智能体之间就有了稳定契约。3.2 三个可运行的子智能体示例资料查证 Agent 接收专题关键词输出查证提纲class ResearchAgent(BaseAgent): name research def _build_user_prompt(self, task: Task) - str: return ( f专题关键词{task.payload[keyword]}\n 请列出3个可查证方向每个方向说明来源类型和要验证的关键问题。 )写作 Agent 接收关键词和查证结果输出文章初稿class DraftAgent(BaseAgent): name draft def _build_user_prompt(self, task: Task) - str: return ( f专题关键词{task.payload[keyword]}\n f查证提纲\n{task.payload[research_result]}\n 请基于提纲写一篇 Markdown 技术短文控制在800字以内。 )审校 Agent 接收草稿输出修改意见class ReviewAgent(BaseAgent): name review def _build_user_prompt(self, task: Task) - str: return ( f专题关键词{task.payload[keyword]}\n f文章草稿\n{task.payload[draft_result]}\n 请审校以上文章指出事实表述、结构和空泛表达问题并给出修改建议。 )每个子智能体的输入都来自上一阶段输出字段名必须一致。这里统一使用research_result、draft_result作为中间结果字段。3.3 子智能体的验证方式不要一上来就启动完整流程。先单独验证每个子智能体是否能正常返回python -m asyncio在测试脚本中直接调用import asyncio from llm_client import LLMClient from agent import ResearchAgent async def test(): cfg {base_url: https://api.openai.com/v1, api_key: your-key, model: gpt-4o-mini, temperature: 0.3, max_tokens: 512} llm LLMClient(cfg) agent ResearchAgent(llm, system_prompt你负责资料查证。) result await agent.run(Task(T1, research, {keyword: 异步任务编排原理})) print(result.content) asyncio.run(test())这一步常见的坑有两个payload中字段名写错下游 Agent 取不到值。解决方式是在每个 Agent 的_build_user_prompt中先做字段存在性校验而不是直接调用task.payload[xxx]导致 KeyError。System Prompt 没有覆盖角色边界导致查证 Agent 写出整篇文章。解决方式是 Prompt 中明确“只输出什么不要输出什么”。4. 从顺序编排到异步任务编排4.1 顺序链路研究、写作、审校子智能体就绪后最简单的协作方式是顺序执行# sync_demo.py from agent import ResearchAgent, DraftAgent, ReviewAgent, Task async def run_sync_pipeline(keyword: str): payload {keyword: keyword} research_result await research_agent.run( Task(T1, research, payload) ) payload[research_result] research_result.content draft_result await draft_agent.run( Task(T2, draft, payload) ) payload[draft_result] draft_result.content review_result await review_agent.run( Task(T3, review, payload) ) return review_result.content顺序编排的优点是简单、结果稳定适合理解阶段依赖。但它的缺点在生产场景中非常明显一次完整流程的耗时等于三个 Agent 耗时之和并且同一时间只能处理一个专题。4.2 为什么生产场景必须切换成异步编排当任务数量变大顺序执行会暴露几个问题模型接口延迟高三个 Agent 串行会导致单任务耗时很长用户请求进入系统后需要等待全部子智能体完成才能返回多个专题任务无法并发处理资源利用率低单个 Agent 失败会导致整个链路中断缺少重试和隔离能力。同步和异步编排的差异如下维度顺序同步编排异步任务编排单任务耗时各阶段累加各阶段累加但任务可并行多任务支持不支持支持并发处理故障隔离一个失败全链路失败单任务失败可重试或单独补偿资源利用率低受 worker 数控制可观测性弱只能看 print可围绕 task_id 聚合日志实现复杂度低中等这里的“异步”并不是指某个 Agent 内部使用了异步代码而是指任务投递和任务执行解耦调用方提交任务后立即返回worker 池在后台消费队列。这正是 DeepAgents 异步任务编排的核心思想。4.3 异步编排中的关键参数进入实现前先明确几个参数的作用。这些参数在后续coordinator.py中会直接使用参数默认值建议调大影响调小影响max_workers3并发提升API 压力增大吞吐下降max_retries2容错更强故障时间更长失败容易丢弃per_task_timeout60 秒容忍慢模型长任务频繁超时queue_size100积压更多任务提交方更快被阻塞这些参数没有“万能值”。实际项目里要根据模型平均响应时间、API 限流策略和业务容忍度综合调整。学习阶段建议max_workers从 3 开始超时从 60 秒开始。5. 异步任务编排的完整实现5.1 Coordinator 与内存任务队列设计思路是队列中不再放单个智能体任务而是放“一个完整专题 Job”。每个 Job 内部包含多阶段的子智能体调用。worker 从队列中取出 Job 后依次执行其内部流程。这样做的好处是任务边界与业务边界一致一个 worker 整条链路负责到底。# coordinator.py import asyncio import logging logger logging.getLogger(coordinator) class Coordinator: def __init__(self, cfg: dict, agents: dict): self.queue asyncio.Queue(maxsizecfg.get(queue_size, 100)) self.agents agents self.max_workers cfg.get(max_workers, 3) self.max_retries cfg.get(max_retries, 2) self.timeout cfg.get(per_task_timeout, 60) self.done [] async def submit(self, job): await self.queue.put(job) async def _worker(self, worker_index: int): while True: job await self.queue.get() try: logger.info(worker%s receive job%s keyword%s, worker_index, job.job_id, job.keyword) result await asyncio.wait_for( job.execute(self.agents), timeoutself.timeout, ) self.done.append(result) logger.info(worker%s finish job%s, worker_index, job.job_id) except Exception as exc: logger.warning(worker%s job%s error%s, worker_index, job.job_id, exc) if job.retry_count self.max_retries: job.retry_count 1 await self.queue.put(job) finally: self.queue.task_done() async def run(self): workers [ asyncio.create_task(self._worker(i)) for i in range(self.max_workers) ] await self.queue.join() for w in workers: w.cancel() await asyncio.gather(*workers, return_exceptionsTrue)这里有一个非常容易写错的点finally中的queue.task_done()对应的是本次queue.get()。如果任务重试会再次put后续get还会再调用一次task_done()。queue.join()会等待所有已写入任务被消费完毕因此不会提前退出。5.2 Job把多个子智能体编排成一个工作单元在jobs.py中定义专题任务# jobs.py from agent import Task class TopicJob: def __init__(self, job_id: str, keyword: str, stages: list[str]): self.job_id job_id self.keyword keyword self.stages stages self.retry_count 0 self.outputs {} async def execute(self, agents: dict) - dict: payload {keyword: self.keyword} for stage in self.stages: agent agents[stage] task Task( task_idf{self.job_id}-{stage}, targetstage, payloaddict(payload), ) result await agent.run(task) self.outputs[stage] result.content payload[f{stage}_result] result.content payload[job_id] self.job_id return payloadJob 内部是顺序链路外部是并发调度。这正好解决了阶段依赖问题同一个 Job 内三个 Agent 有先后顺序不同 Job 之间由多个 worker 并行消费。这里的stages是子智能体名称列表从配置或 main.py 中注入。以后如果要插入新的子智能体比如“配图 Agent”只需要在列表中加一项并保证数据流转字段不冲突。5.3 并发 worker、超时与重试worker 并发数直接决定同一时刻最多有多少个 Job 在执行。以max_workers3为例三个任务可以并行处理每个任务内部串行执行三个阶段。这样既避免了多个 Agent 同时抢占同一个上下文的混乱又提升了整体吞吐。超时使用asyncio.wait_for包裹整个 Job。当某个 Agent 的模型调用卡住超过per_task_timeout后抛出TimeoutError该 Job 进入重试逻辑。重试需要特别注意一个策略问题当前实现只要 Job 抛出异常就整体重试会重新执行已经成功的阶段造成 API 调用浪费。生产环境更合理的方式是记录 Job 已经完成到哪个阶段从失败的阶段开始重试为每个阶段设置独立超时。学习阶段可以先用整体重试理解了流程后再升级为阶段级重试。6. 运行验证与日志观测6.1 准备一份测试输入在main.py中准备三个专题关键词# main.py import asyncio import json import logging import yaml from llm_client import LLMClient from agent import ResearchAgent, DraftAgent, ReviewAgent from coordinator import Coordinator from jobs import TopicJob def load_agents(cfg: dict): llm_cfg cfg[llm] llm LLMClient(llm_cfg) agent_cfgs cfg[agents] agents { research: ResearchAgent(llm, agent_cfgs[research][system_prompt]), draft: DraftAgent(llm, agent_cfgs[draft][system_prompt]), review: ReviewAgent(llm, agent_cfgs[review][system_prompt]), } return agents async def main(): logging.basicConfig( levellogging.INFO, format%(asctime)s %(levelname)s %(name)s %(message)s, ) cfg yaml.safe_load(open(config.yaml, encodingutf-8)) agents load_agents(cfg) coordinator Coordinator(cfg[workflow], agents) topics [ (J001, 异步任务编排原理), (J002, 多智能体系统核心架构), (J003, harness工程实践), ] for job_id, keyword in topics: job TopicJob( job_idjob_id, keywordkeyword, stages[research, draft, review], ) await coordinator.submit(job) await coordinator.run() with open(outputs/results.json, w, encodingutf-8) as f: json.dump(coordinator.done, f, ensure_asciiFalse, indent2) if __name__ __main__: asyncio.run(main())注意LLMClient的__init__需要接收一个 dict 配置内部解析base_url、api_key、model、temperature和max_tokens。api_key如果使用${OPENAI_API_KEY}占位符需要先解析环境变量或者直接在配置中填入可用值。6.2 运行命令与预期输出启动程序export OPENAI_API_KEYyour-key python main.py预期日志大致如下2026-01-06 10:00:01 INFO coordinator worker0 receive jobJ001 keyword异步任务编排原理 2026-01-06 10:00:01 INFO coordinator worker1 receive jobJ002 keyword多智能体系统核心架构 2026-01-06 10:00:02 INFO coordinator worker2 receive jobJ003 keywordharness工程实践 2026-01-06 10:00:35 INFO coordinator worker0 finish jobJ001 2026-01-06 10:00:38 INFO coordinator worker1 finish jobJ002 2026-01-06 10:00:41 INFO coordinator worker2 finish jobJ003运行结束后检查outputs/results.json每个 Job 的结果都应包含job_idkeywordresearch_resultdraft_resultreview_result缺少任何一项都说明某个阶段没有正常写入结果。6.3 用日志串联 task_id日志是所有异步系统排查的基础。当前示例中job_id已经在日志中出现但task_id还没有。建议在每个 Agent 内部把task.task_id注入到日志字段中logger.info(agent%s task%s start, self.name, task.task_id)生产环境推荐使用结构化日志{time: 2026-01-06T10:00:01Z, level: INFO, job_id: J001, stage: research, message: task start}有了job_idstage可以用一条 grep 命令还原一个完整 Job 的执行轨迹grep J001 outputs/app.log这是异步编排中最便宜、最有效的可观测手段。7. 常见问题排查清单7.1 三个高频故障第一个高频故障是“任务提交了但 worker 不消费”。现象是日志中没有 worker 接收记录。先检查coordinator.run()是否被调用再看 worker 是否被asyncio.create_task创建。如果启动时忘记await queue.join()主协程可能直接退出。第二个高频故障是字段链断裂。典型表现为写作 Agent 的 Prompt 中research_result是空的但查证 Agent 明明已经执行成功。原因通常是查证结果没有写回到 payload或者写回的字段名与下游读取字段名不一致。排查方法是打印每个阶段执行完后的 payload 快照。第三个高频故障是模型接口限流。现象是任务频繁进入重试日志中出现 429 或限流错误。解决方式不是调大max_retries而是降低max_workers或增大请求间隔。学习环境可以串行跑生产环境需要做接口限流退避。7.2 从现象倒推原因的检查顺序遇到异常时按下面的顺序排查不要直接猜任务有没有进入队列在submit前后打日志确认 job 被写入worker 有没有收到任务查“receive job”日志模型调用有没有返回看 Agent 内部是否抛出超时或 API 错误中间结果有没有写入 payload打印每个阶段完成后的 payload 字段结果有没有被写入done确认输出文件内容是否触发了重试看retry_count日志。这个顺序的本质是从数据流上游往下游查。任务不进队列查提交方进队列没消费查 worker消费了没执行查 Agent执行了没结果查结果回收。7.3 常见问题速查表问题现象常见原因检查方式处理建议程序启动后立即退出没有await queue.join()检查 main 协程结构等待队列消费完成worker 不消费任务worker 未创建或未启动检查 create_task 是否执行启动 worker 池KeyError: payload 字段中间结果字段名不一致打印 payload 快照统一字段命名任务一直重试API 限流或超时查看异常日志降低并发或增大超时模型返回 JSON 解析失败模型未按格式输出输出原始字符串Prompt 中约束格式或加解析兜底结果文件缺少某个阶段阶段 Agent 抛异常被吞掉查看 warn 日志不要在 except 后静默注意不要用裸except Exception吞掉所有异常。至少要记录异常类型、任务 ID、Agent 名称和当前阶段否则限流、超时、字段错误会全部混在一起无法排查。8. 从示例到生产工程化扩展与学习建议8.1 内存队列替换为消息队列本文实现使用asyncio.Queue适合单机演示和开发调试。一旦需要多实例部署、任务重启恢复就必须替换为真正的消息队列。常见选择包括Redis Stream轻量适合中小规模任务RabbitMQ支持复杂路由和确认机制Kafka适合高吞吐事件流。替换时需要注意消息内容应该是 Job 的序列化描述而不是 Python 对象。Worker 反序列化后重新构造TopicJob再执行相同的 Agent 流程。8.2 状态持久化、幂等与监控生产环境必须解决三个问题。状态持久化Job 的每一次状态变化都要落库。状态可以是READY、RUNNING、SUCCEEDED、FAILED。重启后从数据库中恢复未完成 Job。幂等消费模型调用可能因为网络原因重复执行。设计上要有任务 ID下游消费时判断是否已经处理过避免重复写入结果。监控指标至少采集四类指标队列积压长度、Job 成功率、平均执行时长、最大执行时长。这些指标能直接反映编排层是否健康。8.3 下一步学习路径本文覆盖的是多智能体协作的最小闭环。继续深入可以从以下方向展开为子智能体接入工具调用比如搜索、代码执行、数据库查询引入阶段级重试避免整个 Job 从头执行学习 LangGraph 这类专门的状态图编排框架将 Agent 建模成图节点增加评估集用自动化方式验证多智能体输出质量研究上下文压缩策略当子智能体数量增加时控制上下文长度。学习建议是先把本示例跑通然后刻意制造一个故障比如把模型接口 URL 写错、把字段名改错再通过日志还原问题。能自己完成现象到根因的推导多智能体异步编排这一关才算真正过了。一个多智能体系统能不能跑起来取决于模型调用通不通能不能稳定跑下去取决于编排层的队列、超时、重试、状态恢复和日志是否完整。DeepAgents 不是玄学它就是把这些工程问题一个个解决掉。理解了任务、智能体和编排器三者之间的关系你就能在任意框架下设计出清晰、可观测、可扩展的多智能体协作流程。