LangChain异步工具调用链断裂问题及并发并行实现方案咨询
解决LangChain链式异步调用断裂问题的方案
核心思路
Celery异步任务默认返回task_id,但LangChain链式调用需要前一步的实际输出作为输入,所以核心是让LangChain能够等待异步任务完成并获取真实结果,而非直接拿到任务ID。以下提供两种落地方案:
方案一:自定义AsyncTool适配Celery异步任务
通过封装Celery任务的结果获取逻辑,让LangChain工具直接返回任务执行结果,而非task_id。
1. Celery基础配置(tasks.py)
from celery import Celery import time celery_app = Celery( 'langchain_tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/0' ) # 模拟web搜索长耗时任务 @celery_app.task def web_search_task(query: str) -> str: time.sleep(3) return f"搜索结果:{query}的最新行业动态..." # 模拟图像生成长耗时任务 @celery_app.task def image_gen_task(prompt: str) -> str: time.sleep(5) return "https://example.com/generated_mountain.png" # 模拟图像描述长耗时任务 @celery_app.task def image_caption_task(image_url: str) -> str: time.sleep(2) return f"图像内容:这是一张展示{image_url.split('/')[-1].split('.')[0]}的高清风景图"
2. 自定义LangChain异步工具(tools.py)
from langchain.tools import AsyncTool from tasks import celery_app, web_search_task, image_gen_task, image_caption_task class AsyncWebSearchTool(AsyncTool): name = "async_web_search" description = "异步执行web搜索,输入为搜索关键词" async def _run(self, query: str) -> str: task = web_search_task.delay(query) # 等待任务完成,设置超时避免请求挂死 result = task.get(timeout=30) return result class AsyncImageGenTool(AsyncTool): name = "async_image_gen" description = "根据文本提示异步生成图像,输入为图像描述prompt" async def _run(self, prompt: str) -> str: task = image_gen_task.delay(prompt) result = task.get(timeout=60) return result class AsyncImageCaptionTool(AsyncTool): name = "async_image_caption" description = "异步生成图像的文本描述,输入为图像URL" async def _run(self, image_url: str) -> str: task = image_caption_task.delay(image_url) result = task.get(timeout=30) return result
3. 构建异步LangChain链并集成FastAPI(main.py)
from fastapi import FastAPI from langchain.agents import initialize_agent, AgentType from langchain.llms import OpenAI from tools import AsyncWebSearchTool, AsyncImageGenTool, AsyncImageCaptionTool app = FastAPI() # 初始化LLM(可替换为自定义LLM) llm = OpenAI(temperature=0.7) # 加载异步工具 tools = [ AsyncWebSearchTool(), AsyncImageGenTool(), AsyncImageCaptionTool() ] # 初始化异步代理链,支持链式调用 async_agent = initialize_agent( tools, llm, agent=AgentType.CHAIN_REACT_DESCRIPTION, verbose=True, handle_parsing_errors=True ) @app.post("/chain-execution") async def execute_chain(input_query: str): # 执行完整链式调用:搜索→生成图像→描述图像 final_result = await async_agent.arun(input=input_query) return {"output": final_result} if __name__ == "__main__": import uvicorn uvicorn.run(app, host="0.0.0.0", port=8000)
方案二:并行执行独立任务(非严格链式场景)
如果部分任务无依赖关系,可通过asyncio.gather实现并行执行,提升整体效率:
# 在main.py中新增并行执行接口 import asyncio @app.post("/parallel-execution") async def parallel_tasks(query: str): # 并行启动搜索和图像生成任务 search_task = AsyncWebSearchTool()._run(query) image_gen_task = AsyncImageGenTool()._run(f"关于{query}的自然风景图") # 等待两个并行任务完成 search_result, image_url = await asyncio.gather(search_task, image_gen_task) # 串行执行图像描述任务(依赖图像URL) caption_result = await AsyncImageCaptionTool()._run(image_url) return { "search_result": search_result, "generated_image": image_url, "image_caption": caption_result }
关键注意事项
- 超时控制:在
task.get()中设置合理超时时间,避免FastAPI请求因Celery任务挂起而超时 - 异常处理:添加try-except捕获Celery任务失败、超时等异常,避免链中断
- Celery优化:根据任务类型配置worker并发数,比如CPU密集型任务减少并发,IO密集型任务增加并发
- 版本兼容:确保使用LangChain v0.1+版本,该版本对异步工具的支持更完善
内容的提问来源于stack exchange,提问作者Hassan Baig
相关产品推荐
相关产品推荐

