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

如何通过LangChain的astream_event()实现带依赖工具的顺序执行?

工具依赖与异步流式执行解决方案

针对你的需求——让有依赖关系的工具按顺序执行,同时保留astream_events的流式输出能力,提供以下两种可行方案:

方案一:细粒度工具依赖控制(推荐)

通过给工具添加依赖标记,并在工具执行前等待依赖完成,实现精准的顺序控制,同时不影响无依赖工具的并发执行。

步骤1:给工具添加依赖元数据

为需要依赖的工具添加dependencies属性,声明它依赖的工具名称:

import asyncio
from langchain.tools import tool

@tool
async def tool1():
    """tool1"""
    print('start1...')
    cmd = {'cmd': 'Scrape'}
    await asyncio.sleep(5)
    go2_command_queue.put(cmd)
    print('time1....')

@tool
async def tool2():
    """tool2"""
    print('start2...')
    print(go2_command_queue.qsize())
    print('time2....')

# 标记tool2无依赖
tool2.dependencies = []

@tool
async def tool3():
    """tool3:依赖tool2执行完成"""
    # 等待tool2完成的信号
    await tool_completed_events['tool2'].wait()
    print('start3...')
    # 这里编写tool3的业务逻辑

# 标记tool3依赖tool2
tool3.dependencies = ['tool2']

步骤2:修改流式事件处理逻辑

在agent_completion函数中,维护工具完成状态,当工具启动前检查依赖是否完成:

async def agent_completion(
    agent_executor,
    message: str,
    tools: list,
    tool_completed_events: dict
) -> AsyncGenerator:
    """基于查询决策工具调用,支持异步流式输出"""
    tool_names = [tool.name for tool in tools]

    async for event in agent_executor.astream_events(
        {
            "input": message,
            "tools": tools,
            "tool_names": tool_names,
            "agent_scratchpad": lambda x: format_to_openai_tool_messages(x["intermediate_steps"]),
        },
        version='v2'
    ):
        kind = event['event']
        if kind == "on_chain_start":
            if event["name"] == "Agent":
                yield f"\n### Agent: `{event['name']}`,Agent Input: `{event['data'].get('input')}`\n"
        elif kind == "on_chat_model_stream":
            content = event["data"]["chunk"].content
            if content:
                yield content
        elif kind == "on_tool_start":
            tool_name = event['name']
            # 获取当前工具的依赖列表
            current_tool = next(t for t in tools if t.name == tool_name)
            dependencies = getattr(current_tool, 'dependencies', [])
            # 等待所有依赖工具完成
            for dep_name in dependencies:
                await tool_completed_events[dep_name].wait()
            yield f"\n### Tool: `{tool_name}`,Tool Input: `{event['data'].get('input')}`\n"
        elif kind == "on_tool_end":
            tool_name = event['name']
            # 标记当前工具完成,触发后续依赖工具的执行
            tool_completed_events[tool_name].set()
            # 重置事件,以便下一次查询复用
            tool_completed_events[tool_name].clear()
            yield f"\n### Tool Finished: `{tool_name}`,Tool Results: \n"
            yield f"`{event['data'].get('output')}`\n"
        elif kind == "on_chain_end":
            if event["name"] == "Agent":
                yield f"\n### Agent Finished: `{event['name']}`,Agent Results: \n"
                yield f"{event['data'].get('output')['output']}\n"

步骤3:主函数中初始化事件

每次用户查询时,重新初始化工具完成事件,避免会话间的状态干扰:

if __name__ == '__main__':
    llm = ChatOpenAI(
        model='gpt-4o',
        temperature=0.,
        api_key=OPENAI_API_KEY,
        streaming=True,
        max_tokens=None,
    )
    tools = [tool1, tool2, tool3]
    agent = create_openai_tools_agent(llm, tools, AGENT_PROMPT)
    agent_executor = AgentExecutor(
        agent=agent,
        tools=tools,
        verbose=False,
        return_intermediate_steps=True
    )

    while True:
        user_message = input('Query: ')
        # 每次查询初始化工具完成事件
        tool_completed_events = {tool.name: asyncio.Event() for tool in tools}
        responses = agent_completion(agent_executor, user_message, tools, tool_completed_events)
        async for response in responses:
            print(response, flush=True)

方案二:全局串行控制(简单但不灵活)

如果所有工具都需要按顺序执行,可直接设置AgentExecutor的max_concurrent_tools=1参数,强制工具串行执行,同时保留流式输出:

agent_executor = AgentExecutor(
    agent=agent,
    tools=tools,
    verbose=False,
    return_intermediate_steps=True,
    max_concurrent_tools=1  # 限制同时执行的工具数量为1
)

该方案优点是实现简单,但会让所有工具失去并发能力,仅适合所有工具都有严格顺序依赖的场景。


内容的提问来源于stack exchange,提问作者Xiang

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 00:08:10