如何通过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
相关产品推荐
相关产品推荐

