异步与并发,用asyncio加速Agent执行效率

📅 2026/8/7 11:54:04
异步与并发,用asyncio加速Agent执行效率
异步与并发用asyncio加速Agent执行效率上一篇聊流式输出里面用到了astream和async当时说异步这东西在Agent里到处都是但很多人用着用着就踩坑。这篇就把asyncio在Agent里的门道从头捋一遍看看同步异步到底差在哪怎么用并行把执行效率提上去。为什么Agent需要异步Agent跑一次任务背后往往要调好几次模型调好几个工具还可能查向量库、读文件。每一步都是IO操作都要等网络响应。同步写法就是一步一步排队第一步没回来第二步干等着时间全耗在等上。举个具体的数。假设调一次模型要2秒Agent跑一个任务要调5次模型同步写法就是10秒串起来。但这5次调用之间如果没有依赖关系完全可以同时发出去2秒就全回来了。异步就是干这个的等待的时间从串着排队变成并着走。我自己踩过这个坑。做一个文档批处理脚本50个文档逐个调模型做摘要同步写法跑了快两分钟。当时没多想后来改异步并发十几秒跑完差了一个数量级。那之后我就记住批量IO场景里异步是刚需。同步invoke与异步ainvokeLangChain里所有Runnable都带两套方法同步的invoke和异步的ainvoke。先看同步写法。fromlangchain_openaiimportChatOpenAIfromlangchain_core.promptsimportChatPromptTemplatefromlangchain_core.output_parsersimportStrOutputParser modelChatOpenAI(modelgpt-4o-mini)promptChatPromptTemplate.from_template(总结这段文字{text})chainprompt|model|StrOutputParser()resultchain.invoke({text:一段待处理的文字})这段代码跑起来没问题单次调用该等多久还是等多久。换成ainvoke单次调用耗时几乎没差别区别在它能配合asyncio做并发。importasyncioasyncdefmain():resultawaitchain.ainvoke({text:一段待处理的文字})print(result)asyncio.run(main())ainvoke前面要加await整个函数得是async def。单看这一段好像只是把invoke换了名字再加个await看不出好处。好处要等到多个调用一起发的时候才显现。用asyncio.gather做并行调用并发的主角是asyncio.gather。它把多个协程同时丢进去一起等全部完成后一起返回结果。前面那个5次调用从10秒压到2秒靠的就是它。importasyncioasyncdefsummarize(text,chain):returnawaitchain.ainvoke({text:text})texts[文档一内容,文档二内容,文档三内容,文档四内容,文档五内容]asyncdefmain():tasks[summarize(t,chain)fortintexts]resultsawaitasyncio.gather(*tasks)print(results)asyncio.run(main())gather接收的是协程对象列表用星号展开传进去。它会让这些协程并发执行谁先回来谁的位先占好最后按传入顺序返回结果列表。这点很贴心结果顺序和输入顺序对得上不用自己再排。并发也不是无脑拉满。API有并发限制同时发50个请求很可能被限流甚至封号。实际用的时候我会配合asyncio.Semaphore控制并发数比如限制同时最多5个。semasyncio.Semaphore(5)asyncdefsummarize(text,chain):asyncwithsem:returnawaitchain.ainvoke({text:text})加一个信号量超出数量的请求排队等前面的释放稳当得多。还有个细节得提一句gather默认只要有一个协程抛异常整批就报错其余的结果也拿不到。批处理几十个文档时一个文件格式有问题就全盘崩溃挺烦的。给gather传个return_exceptionsTrue异常会被当成结果返回你拿到手再逐个判断哪些成功哪些失败跑完一遍心里有数。批量处理多个文档gather最常见的场景就是批处理。手里一堆文档要做摘要、做分类、做抽取逐个调太慢gather一把梭。importasynciofrompathlibimportPathasyncdefprocess_file(path,chain):textPath(path).read_text(encodingutf-8)summaryawaitchain.ainvoke({text:text})return{file:path,summary:summary}asyncdefbatch_summarize(folder,chain):fileslist(Path(folder).glob(*.txt))tasks[process_file(str(f),chain)forfinfiles]returnawaitasyncio.gather(*tasks)asyncio.run(batch_summarize(./docs,chain))这里读文件用的是同步的read_text。文件小、本地IO快的时候没啥感觉文件大或者走网络存储就会卡。遇到大文件读文件也得换异步的比如aiofiles。判断标准很简单凡是IO操作能异步就异步别给事件循环添堵。处理结果按文件顺序返回写回磁盘或者入库都方便。我一般会把结果存成jsonl一行一个文档的摘要后面检索或者分析都能直接用。异步在Web服务里的应用Agent真上量基本都跑在Web服务里。FastAPI这类框架本身就是异步的路由函数写成async里面调ainvoke、astream从入口到模型调用全程异步吞吐量比同步框架高一大截。fromfastapiimportFastAPIfromlangchain_openaiimportChatOpenAIfromlangchain_core.promptsimportChatPromptTemplatefromlangchain_core.output_parsersimportStrOutputParser appFastAPI()modelChatOpenAI(modelgpt-4o-mini)promptChatPromptTemplate.from_template(回答{question})chainprompt|model|StrOutputParser()app.get(/ask)asyncdefask(question:str):return{answer:awaitchain.ainvoke({question:question})}路由函数加async里面用await接ainvoke。这样多个用户同时请求框架能并发处理一个请求等模型的时候CPU去伺候别的请求不会互相堵。要是这里写成同步invoke整个进程同一时刻只能处理一个请求第二个用户排队等第一个跑完。并发一上来响应时间直线上升体验崩得很快。两个最常踩的坑异步听着美好坑也实打实多。说两个我自己栽过的。第一个是阻塞调用混进异步函数。async函数里塞了一个同步的requests.get或者time.sleep或者同步的数据库查询。这些操作不会让出事件循环整个循环被它卡住所有协程全停。表现就是服务突然卡死CPU没占用但请求全超时。解决办法就一条异步函数里只放异步操作。要调同步阻塞代码用asyncio.to_thread扔到线程池里跑别让它直接占着事件循环。importasyncioimportrequestsasyncdeffetch(url):returnawaitasyncio.to_thread(requests.get,url)第二个坑是事件循环套事件循环。脚本里用asyncio.run启了一个循环循环里又调了某个库那个库内部自己再起一个asyncio.run直接报错说already running event loop。这种情况常见于在异步上下文里调Jupyter或者一些老库。碰到这个错先查有没有在异步函数里嵌套跑循环把内层那个换成await对应协程就行。还有个隐蔽的Windows上asyncio默认事件循环策略和Mac、Linux不一样偶尔会报NotImplementedError。脚本开头加一句asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy())能消停。这个坑Windows用户碰到过都懂。小结这篇把asyncio在Agent里的用法捋了一遍。为什么需要异步因为Agent里IO操作多同步排队等太亏。ainvoke是异步版的invoke单看没差别配gather才显出威力。asyncio.gather让多个调用并行跑批量处理文档效率翻几倍。Web服务里异步贯通吞吐量高但前提是别混进阻塞调用别套事件循环。把这些门道摸清Agent跑起来又快又稳。不过快和稳还不够省钱。每次调用都重新算一遍重复的活儿一遍遍干token哗哗地烧。下一篇就聊缓存机制看看怎么把算过的结果存下来能省一截是一截。