You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 02:52:31