ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

DeepSeek-Agent-Harness-2026终极指南-第12章第56节-子Agent与编排-并行子Agent:asyncio与任务分派

DeepSeek-Agent-Harness-2026终极指南-第12章第56节-子Agent与编排-并行子Agent:asyncio与任务分派 DeepSeek Agent Harness 2026终极指南 - 第12章第56节 并行子Agentasyncio与任务分派第55节的子 Agent 能干活了但现在是串行——一个接一个。如果三个子 Agent审查、测试、优化能同时跑效率翻倍。这节做并行子 Agent用 asyncio 并发执行多个子 Agent控制并发数、聚合部分失败的结果。三个 Agent 并行审查代码实测耗时对比——串行 9 秒 vs 并行 3 秒。本文导航串行的瓶颈一个等一个asyncio.gather 并行执行并发数控制信号量限流部分失败的结果聚合完整实现parallel_agents.py实测串行 vs 并行耗时对比小结串行的瓶颈一个等一个第55节的子 Agent 是串行执行的。如果你要并行审查三个文件# 串行一个接一个总耗时 三个之和result1reviewer.run(审查 file1.py)# 3秒result2reviewer.run(审查 file2.py)# 3秒result3reviewer.run(审查 file3.py)# 3秒# 总耗时9秒问题在哪子 Agent 的大部分时间在等待模型响应网络 IOCPU 是闲的。串行执行浪费了这段等待时间。子Agent3子Agent2子Agent1主程序子Agent3子Agent2子Agent1主程序串行9秒审查file13秒后返回审查file23秒后返回审查file33秒后返回并行化的思路三个子 Agent 同时发起请求各自等待总耗时 ≈ 最长的一个3秒而不是三个之和9秒。asyncio.gather 并行执行Python 的 asyncio 是处理 IO 并发的标准工具。核心是asyncio.gather()——同时跑多个协程等全部完成。importasyncioasyncdefrun_parallel(tasks:list[tuple[str,str]])-list[str]: 并行执行多个子 Agent 任务。 tasks: [(agent_name, task), ...] 返回结果列表与输入顺序一致 asyncdef_run_one(agent:SubAgent,task:str)-str:# 用 asyncio.to_thread 让同步的 run 在子线程跑不阻塞事件循环returnawaitasyncio.to_thread(agent.run,task)# 并发执行所有任务coroutines[_run_one(agent,task)foragent,taskintasks]resultsawaitasyncio.gather(*coroutines)returnlist(results)关键点asyncio.to_thread()把同步的agent.run()放到子线程执行避免阻塞事件循环。因为子 Agent 的run()是同步代码内部用同步的client.chat()直接 await 会卡死事件循环。并发数控制信号量限流如果一次要跑 100 个子 Agent全并发会导致API 限流超出 DeepSeek 的并发限制2500但实际受账号配额限制内存爆炸100 个上下文同时存在不可控某个慢任务拖垮整体所以要用**信号量Semaphore**限流importasyncioasyncdefrun_parallel_limited(tasks:list,max_concurrent:int5)-list[str]:并行执行但最多 max_concurrent 个同时跑semaphoreasyncio.Semaphore(max_concurrent)asyncdef_run_with_limit(agent,task):asyncwithsemaphore:# 获取信号量满了就等returnawaitasyncio.to_thread(agent.run,task)coroutines[_run_with_limit(a,t)fora,tintasks]returnawaitasyncio.gather(*coroutines)信号量的作用当同时运行的任务达到上限时新的任务要排队等。这样并发数可控不会打爆 API。部分失败的结果聚合并行执行时某个子 Agent 可能失败超时、异常。不能让一个失败拖垮整个批次。用return_exceptionsTrue让gather收集异常而不是抛出asyncdefrun_parallel_safe(tasks:list,max_concurrent:int5)-list[dict]: 并行执行单个失败不影响整体。 返回 [{status: ok/failed, result/error}, ...] semaphoreasyncio.Semaphore(max_concurrent)asyncdef_run_one(agent,task):asyncwithsemaphore:try:resultawaitasyncio.to_thread(agent.run,task)return{status:ok,result:result}exceptExceptionase:return{status:failed,error:str(e)}coroutines[_run_one(a,t)fora,tintasks]returnawaitasyncio.gather(*coroutines)这样即使某个子 Agent 失败其他子 Agent 的结果也能正常聚合。失败的信息也会返回供上层判断。完整实现parallel_agents.py# deep_pilot/parallel_agents.py —— 并行子Agent v0.7from__future__importannotationsimportasynciofromtypingimportAnyfromdeep_pilot.loggerimportget_loggerfromdeep_pilot.sub_agentimportSubAgent loggerget_logger(__name__)classParallelAgentRunner:并行子 Agent 运行器def__init__(self,max_concurrent:int5):self.max_concurrentmax_concurrentdefrun(self,tasks:list[tuple[SubAgent,str]])-list[dict[str,Any]]: 并行执行多个子 Agent 任务同步入口。 tasks: [(agent, task), ...] 返回 [{status, result/error}, ...]顺序与输入一致。 returnasyncio.run(self._run_async(tasks))asyncdef_run_async(self,tasks:list[tuple[SubAgent,str]])-list[dict[str,Any]]:semaphoreasyncio.Semaphore(self.max_concurrent)asyncdef_run_one(agent:SubAgent,task:str)-dict[str,Any]:asyncwithsemaphore:try:resultawaitasyncio.to_thread(agent.run,task)return{status:ok,result:result}exceptExceptionase:logger.error(f子 Agent{agent.name}失败:{e})return{status:failed,error:str(e)}coroutines[_run_one(agent,task)foragent,taskintasks]returnawaitasyncio.gather(*coroutines)# 全局单例parallel_runnerParallelAgentRunner(max_concurrent5)实测串行 vs 并行耗时对比uv run python-c import time from deep_pilot.sub_agent import SubAgent from deep_pilot.parallel_agents import ParallelAgentRunner # 创建三个审查子Agent不同审查重点 reviewer_security SubAgent(security, 你是安全审查员专注找安全漏洞。, [read_file, grep]) reviewer_quality SubAgent(quality, 你是代码质量审查员专注找代码坏味道。, [read_file, grep]) reviewer_perf SubAgent(performance, 你是性能审查员专注找性能瓶颈。, [read_file, grep]) tasks [ (reviewer_security, 审查 tools.py 的安全问题), (reviewer_quality, 审查 tools.py 的代码质量), (reviewer_perf, 审查 tools.py 的性能问题), ] runner ParallelAgentRunner(max_concurrent3) # 串行 start time.time() for agent, task in tasks: agent.run(task) serial_time time.time() - start # 并行 start time.time() results runner.run(tasks) parallel_time time.time() - start print(f串行耗时: {serial_time:.2f}s) print(f并行耗时: {parallel_time:.2f}s) print(f加速比: {serial_time/parallel_time:.1f}x) print(f结果数: {len(results)}, 成功: {sum(1 for r in results if r[\status\]\ok\)}) 控制台输出精简2026-09-14 10:00:01 | INFO | sub_agent | [子Agent:security] 第 1 轮 2026-09-14 10:00:01 | INFO | sub_agent | [子Agent:quality] 第 1 轮 2026-09-14 10:00:01 | INFO | sub_agent | [子Agent:performance] 第 1 轮 ... 串行耗时: 9.21s 并行耗时: 3.18s 加速比: 2.9x 结果数: 3, 成功: 3三个子 Agent 同时启动注意日志时间戳相同并行耗时 3.18 秒串行 9.21 秒加速比 2.9 倍。三个结果全部成功。小结串行的瓶颈是 IO 等待子 Agent 大部分时间在等模型响应CPU 闲着。asyncio.gather 并行同时跑多个协程总耗时 ≈ 最长的一个。asyncio.to_thread把同步的agent.run()放子线程避免阻塞事件循环。信号量限流asyncio.Semaphore控制并发数防止打爆 API。部分失败聚合return_exceptions让单个失败不影响整体失败的也返回状态。实测加速比 2.9x串行 9.21 秒 vs 并行 3.18 秒。DeepPilot v0.7 并行子 Agent 完成——从串行单兵到并行军团大任务效率翻倍。下节预告并行能同时跑多个子 Agent 了但每个子 Agent 都是通用的。更专业的做法是给不同子 Agent 不同的角色——研究员负责查资料、工程师负责写代码、审查员负责挑毛病。下一节做角色系统researcher/coder/reviewer 三剑客各自有不同的提示词和工具权限。如果觉得本文对你有帮助欢迎点赞、收藏、关注三连本系列持续更新中关注不迷路~
返回列表