ARTICLE DETAIL

资讯详情

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

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

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

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

上一篇聊流式输出,里面用到了astream和async,当时说异步这东西在Agent里到处都是,但很多人用着用着就踩坑。这篇就把asyncio在Agent里的门道从头捋一遍,看看同步异步到底差在哪,怎么用并行把执行效率提上去。


为什么Agent需要异步

Agent跑一次任务,背后往往要调好几次模型,调好几个工具,还可能查向量库、读文件。每一步都是IO操作,都要等网络响应。同步写法就是一步一步排队,第一步没回来第二步干等着,时间全耗在等上。

举个具体的数。假设调一次模型要2秒,Agent跑一个任务要调5次模型,同步写法就是10秒串起来。但这5次调用之间如果没有依赖关系,完全可以同时发出去,2秒就全回来了。异步就是干这个的,等待的时间从串着排队变成并着走。

我自己踩过这个坑。做一个文档批处理脚本,50个文档逐个调模型做摘要,同步写法跑了快两分钟。当时没多想,后来改异步并发,十几秒跑完,差了一个数量级。那之后我就记住,批量IO场景里异步是刚需。


同步invoke与异步ainvoke

LangChain里所有Runnable都带两套方法,同步的invoke和异步的ainvoke。先看同步写法。

fromlangchain_openaiimportChatOpenAIfromlangchain_core.promptsimportChatPromptTemplatefromlangchain_core.output_parsersimportStrOutputParser model=ChatOpenAI(model="gpt-4o-mini")prompt=ChatPromptTemplate.from_template("总结这段文字{text}")chain=prompt|model|StrOutputParser()result=chain.invoke({"text":"一段待处理的文字"})

这段代码跑起来没问题,单次调用该等多久还是等多久。换成ainvoke,单次调用耗时几乎没差别,区别在它能配合asyncio做并发。

importasyncioasyncdefmain():result=awaitchain.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]results=awaitasyncio.gather(*tasks)print(results)asyncio.run(main())

gather接收的是协程对象列表,用星号展开传进去。它会让这些协程并发执行,谁先回来谁的位先占好,最后按传入顺序返回结果列表。这点很贴心,结果顺序和输入顺序对得上,不用自己再排。

并发也不是无脑拉满。API有并发限制,同时发50个请求很可能被限流甚至封号。实际用的时候我会配合asyncio.Semaphore控制并发数,比如限制同时最多5个。

sem=asyncio.Semaphore(5)asyncdefsummarize(text,chain):asyncwithsem:returnawaitchain.ainvoke({"text":text})

加一个信号量,超出数量的请求排队等前面的释放,稳当得多。

还有个细节得提一句,gather默认只要有一个协程抛异常,整批就报错,其余的结果也拿不到。批处理几十个文档时,一个文件格式有问题就全盘崩溃,挺烦的。给gather传个return_exceptions=True,异常会被当成结果返回,你拿到手再逐个判断哪些成功哪些失败,跑完一遍心里有数。


批量处理多个文档

gather最常见的场景就是批处理。手里一堆文档要做摘要、做分类、做抽取,逐个调太慢,gather一把梭。

importasynciofrompathlibimportPathasyncdefprocess_file(path,chain):text=Path(path).read_text(encoding="utf-8")summary=awaitchain.ainvoke({"text":text})return{"file":path,"summary":summary}asyncdefbatch_summarize(folder,chain):files=list(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 app=FastAPI()model=ChatOpenAI(model="gpt-4o-mini")prompt=ChatPromptTemplate.from_template("回答{question}")chain=prompt|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哗哗地烧。下一篇就聊缓存机制,看看怎么把算过的结果存下来,能省一截是一截。

返回列表