尧图精选

AI Agent间通信实战:基于LangChain构建多Agent协作系统

🕒 发布时间:2026/9/1 22:41:24 📁 来源:尧图网络
如果你正在探索AI Agent开发可能已经发现了一个关键瓶颈单个Agent能力有限但让多个Agent协作起来却异常困难。每个Agent都有自己的“大脑”模型和“技能”工具如何让它们高效对话、传递信息、协同完成任务而不是各自为战这正是Agent间通信Agent-to-Agent Communication简称A2A要解决的核心问题。很多人以为A2A只是简单的消息传递就像在两个程序间发条字符串。但实际上它远不止于此。一个健壮的A2A机制需要处理状态同步、意图理解、任务分解、错误恢复和资源协调。没有它你的多Agent系统很可能变成一群“聋哑”的智能体空有算力却无法形成合力。本文将通过一个完整的、可运行的Demo带你彻底理解A2A的实战实现。我们将构建一个简单的“旅行规划”场景一个行程规划Agent负责生成旅行计划一个预算评估Agent负责审核成本两者通过明确的通信协议协作最终输出一个合理且经济的方案。你将看到从环境搭建、Agent定义、通信通道建立到任务编排的完整链路并理解背后的设计思想与常见陷阱。读完本文你将能清晰区分Agent间通信与普通API调用的本质不同。使用流行的Agent框架本文以LangChain为例快速搭建一个可工作的A2A原型。掌握设计Agent通信协议的关键要素消息格式、会话管理和超时处理。规避在A2A实践中最常见的几个错误比如循环依赖和状态混乱。1. A2A要解决的真实问题为什么不是简单的函数调用在深入代码之前我们必须先厘清一个根本问题既然Agent本质上也是代码模块为什么不让一个Agent直接调用另一个Agent的函数为什么需要专门的“通信”机制想象一个现实开发场景你需要构建一个智能客服系统其中包含查询理解Agent、知识检索Agent和回复生成Agent。如果采用简单的函数调用你会面临以下挑战强耦合理解Agent必须精确知道检索Agent的函数签名参数、返回值。一旦检索Agent升级接口变化所有调用方都要修改。阻塞与超时理解Agent调用检索Agent的函数时会同步等待其返回。如果检索Agent因网络或计算卡住整个调用链都会阻塞。状态管理困难一次用户对话可能涉及多个Agent的多次交互。函数调用模式难以维护跨Agent的会话上下文比如生成Agent如何知道当前回复是基于哪次检索的结果。缺乏灵活性你很难动态地插入一个新的Agent例如一个情感分析Agent到工作流中因为调用关系是硬编码的。A2A的核心思想是解耦与异步协作。它引入了几个关键概念消息总线Message Bus或通信通道Agent不直接调用彼此而是向一个公共的“通道”发送消息或从其中接收消息。标准化消息格式所有通信都遵循预定义的格式通常包含发送者、接收者、消息类型、内容、会话ID等就像网络协议一样。发布/订阅Pub/Sub或点对点P2P模式Agent可以订阅它关心的消息类型或者直接指定消息的接收者。这样做的好处是显而易见的系统更松耦合、易于扩展、支持异步处理、便于监控和调试通信流。我们的Demo将具体体现这些优势。2. 环境准备与核心工具选型我们将使用LangChain和OpenAI API来构建这个Demo。LangChain是一个强大的LLM应用开发框架它提供了构建和编排Agent的高级抽象非常适合演示A2A概念。即使你之前没有用过LangChain跟着步骤也能轻松上手。为什么选LangChain因为它内置了Agent、Tool、Chain等概念并且对Agent间的结构化通信有良好的支持通过AgentExecutor和自定义Tools模拟通信。这能让我们聚焦于通信逻辑而非从零搭建Agent基础设施。准备工作Python环境确保你已安装Python 3.8或更高版本。安装依赖创建一个新的虚拟环境然后安装以下包。pip install langchain langchain-openaiOpenAI API密钥你需要一个有效的OpenAI API密钥。将其设置为环境变量。# Linux/Mac export OPENAI_API_KEYyour-api-key-here # Windows (PowerShell) $env:OPENAI_API_KEYyour-api-key-here安全提示永远不要将API密钥硬编码在代码中提交到版本库。项目结构预览我们将创建以下文件travel_agent_demo/ ├── main.py # 主程序入口编排整个流程 ├── agents/ # Agent定义目录 │ ├── __init__.py │ ├── planner_agent.py # 行程规划Agent │ └── budget_agent.py # 预算评估Agent └── utils/ └── communication.py # 模拟通信通道的工具类3. 核心概念定义我们的Agent与通信协议在写代码前先定义清楚两个Agent的角色和它们之间的“通信协议”。行程规划Agent (PlannerAgent)职责根据用户输入如“我想去北京玩3天预算中等”生成一份详细的行程草案包括日期、景点、活动、住宿建议等。输出一个结构化的行程文本。需要咨询预算Agent吗是的。它在生成草案后需要将草案发送给预算评估Agent进行审核。预算评估Agent (BudgetAgent)职责接收行程草案对其中的各项花费门票、住宿、餐饮、交通进行估算和评估判断是否超出用户预算并提供优化建议如“某景点门票较贵可替换为...”。输出一份评估报告包含总估算费用、是否超预算、优化建议。行动将评估报告发回给行程规划Agent。通信协议我们自定义的简单格式 为了模拟A2A我们将设计一个简单的Message类。在实际的LangChain多Agent系统中可能会使用更复杂的AgentAction、AgentFinish或自定义Tools来传递信息。# utils/communication.py from typing import Dict, Any, Optional from dataclasses import dataclass dataclass class AgentMessage: 定义Agent间通信的消息格式 sender: str # 发送者Agent名称 receiver: str # 接收者Agent名称 message_type: str # 消息类型如 “query”, “response”, “error” content: Dict[str, Any] # 消息内容用字典存储结构化数据 session_id: str # 会话ID用于关联同一任务的多轮对话 timestamp: Optional[float] None # 时间戳 def to_dict(self): return { sender: self.sender, receiver: self.receiver, message_type: self.message_type, content: self.content, session_id: self.session_id, timestamp: self.timestamp }同时我们创建一个极简的CommunicationChannel类来模拟消息总线# utils/communication.py class SimpleCommunicationChannel: 一个简单的内存消息通道用于演示目的。生产环境应使用消息队列如RabbitMQ, Redis。 def __init__(self): self.messages [] # 存储所有消息 def send_message(self, message: AgentMessage): print(f[Channel] {message.sender} - {message.receiver}: {message.message_type}) self.messages.append(message) # 在实际场景中这里会触发接收者的消息处理回调 def get_messages_for_agent(self, agent_name: str): 获取发送给指定Agent的所有消息简易过滤 return [msg for msg in self.messages if msg.receiver agent_name]这个通道非常基础仅用于演示消息的发送和存储。在真实分布式系统中你会使用成熟的消息中间件。4. 构建行程规划Agent这个Agent的核心是使用LLMGPT来生成行程。我们将使用LangChain的create_react_agent模式并为其装备一个关键的“工具”Tool——consult_budget_agent这个工具就是它与其他Agent通信的接口。# agents/planner_agent.py from langchain.agents import create_react_agent, AgentExecutor from langchain.tools import Tool from langchain_openai import ChatOpenAI from langchain.prompts import PromptTemplate from utils.communication import AgentMessage, SimpleCommunicationChannel import uuid class PlannerAgent: def __init__(self, channel: SimpleCommunicationChannel, llmNone): self.name TravelPlanner self.channel channel self.llm llm or ChatOpenAI(modelgpt-3.5-turbo, temperature0.7) # 定义Agent的“工具”。这里只有一个咨询预算Agent。 tools [ Tool( nameConsultBudgetAgent, funcself._consult_budget_tool, # 工具的实际执行函数 description在生成初步行程草案后必须调用此工具将草案发送给预算评估AgentBudgetAgent进行审核。 输入应该是包含 draft 和 user_budget 的JSON格式字符串。 例如{{draft: ...行程文本..., user_budget: 中等预算}} ) ] # 使用ReAct模式的Prompt模板 planner_prompt PromptTemplate.from_template( 你是一个专业的旅行行程规划师。用户的需求是{user_input} 你的任务是 1. 根据用户需求生成一份详细、合理、有趣的旅行行程草案。 2. 草案生成后你必须调用 ConsultBudgetAgent 工具将草案和用户预算信息发送给预算评估Agent进行审核。 请开始你的工作。首先思考步骤然后生成草案最后调用工具。 ) # 创建ReAct Agent agent create_react_agent(llmself.llm, toolstools, promptplanner_prompt) self.agent_executor AgentExecutor(agentagent, toolstools, verboseTrue, handle_parsing_errorsTrue) def _consult_budget_tool(self, query: str) - str: 这是PlannerAgent与BudgetAgent通信的工具函数。 它并不直接调用BudgetAgent而是向通信通道发送一条消息。 import json try: # 解析输入来自LLM的调用 data json.loads(query) travel_draft data.get(draft, ) user_budget data.get(user_budget, ) # 构建一条消息 message AgentMessage( senderself.name, receiverBudgetEvaluator, message_typebudget_review_request, content{ travel_draft: travel_draft, user_budget_constraint: user_budget }, session_idself.current_session_id # 需要从外部传入 ) # 发送到通道 self.channel.send_message(message) return f消息已发送给BudgetEvaluator。请等待其评估结果。会话ID: {self.current_session_id} except Exception as e: return f调用咨询工具时出错{str(e)} def generate_plan(self, user_input: str, session_id: str) - str: 主执行方法处理用户输入运行Agent并返回最终结果。 self.current_session_id session_id # 运行AgentExecutor。Agent会根据Prompt思考并在适当时机调用我们定义的工具。 result self.agent_executor.invoke({user_input: user_input}) return result.get(output, 行程规划未完成。)关键点解析Tool作为通信桥梁ConsultBudgetAgent工具是PlannerAgent对外通信的“手”。当LLM决定需要预算审核时就会调用这个工具。异步通信模拟工具函数并不等待回复它只是发送一条消息到通道。这模拟了异步非阻塞的通信。实际的回复需要通过另一个机制如轮询或回调来获取我们在主流程中模拟这一点。会话IDsession_id至关重要它将同一任务的所有相关消息请求和响应关联起来尤其是在并发处理多个用户请求时。5. 构建预算评估Agent这个Agent相对独立它监听通道中发给自己的消息类型为budget_review_request处理消息内容然后生成回复消息。# agents/budget_agent.py from langchain_openai import ChatOpenAI from langchain.prompts import PromptTemplate from utils.communication import AgentMessage, SimpleCommunicationChannel import json class BudgetAgent: def __init__(self, channel: SimpleCommunicationChannel, llmNone): self.name BudgetEvaluator self.channel channel self.llm llm or ChatOpenAI(modelgpt-3.5-turbo, temperature0.2) # 温度低一些评估更稳定 def process_incoming_message(self, message: AgentMessage) - AgentMessage: 处理收到的消息核心业务逻辑在这里。 if message.message_type ! budget_review_request: return None # 忽略其他类型的消息 travel_draft message.content.get(travel_draft) user_budget message.content.get(user_budget_constraint) # 使用LLM进行预算评估 evaluation_prompt PromptTemplate.from_template( 你是一个严格的旅行预算评估师。 请审核以下旅行行程草案 --- 行程草案开始 --- {draft} --- 行程草案结束 --- 用户的预算要求是{budget}。 请执行以下任务 1. 估算此行程的大致总费用按人民币计算并分项住宿、交通、门票、餐饮简要说明。 2. 判断该估算是否明显符合或超出用户的预算要求。 3. 提供1-2条具体的优化建议以控制成本如替换某个高价项目。 请以JSON格式输出包含以下字段 - estimated_total_cost: 估算总费用字符串如“约3000元” - is_within_budget: 是否在预算内 (true/false) - cost_breakdown: 费用分项说明字符串 - optimization_suggestions: 优化建议字符串列表 ) chain evaluation_prompt | self.llm evaluation_result_str chain.invoke({draft: travel_draft, budget: user_budget}).content try: evaluation_data json.loads(evaluation_result_str) except json.JSONDecodeError: evaluation_data {error: Failed to parse evaluation result, raw_output: evaluation_result_str} # 构建回复消息 reply_message AgentMessage( senderself.name, receivermessage.sender, # 回复给发送者 message_typebudget_review_response, content{ original_session_id: message.session_id, evaluation: evaluation_data }, session_idmessage.session_id # 保持同一会话ID ) return reply_message def listen_and_process(self): 一个简易的监听方法用于演示。在实际中这应该是一个持续运行的服务或事件驱动。 # 从通道获取所有发给自己的消息 my_messages self.channel.get_messages_for_agent(self.name) replies [] for msg in my_messages: reply self.process_incoming_message(msg) if reply: self.channel.send_message(reply) # 将回复发送回通道 replies.append(reply) return replies关键点解析消息驱动BudgetAgent是被动触发的它通过listen_and_process方法或类似的事件监听器从通道中获取属于自己的消息进行处理。结构化输出我们要求LLM以JSON格式输出评估结果这便于程序化解析和后续处理。这是Agent间通信数据格式化的良好实践。保持会话上下文回复消息中携带了original_session_id确保PlannerAgent能将其与最初的请求关联起来。6. 主流程编排让两个Agent“对话”起来现在我们需要一个“导演”来协调这两个Agent模拟完整的A2A工作流。这个导演负责初始化Agent、启动任务、管理通信轮次。# main.py import uuid from utils.communication import SimpleCommunicationChannel from agents.planner_agent import PlannerAgent from agents.budget_agent import BudgetAgent def main(): # 1. 创建共享的通信通道 channel SimpleCommunicationChannel() # 2. 初始化两个Agent并注入同一个通道实例 planner PlannerAgent(channelchannel) budget_evaluator BudgetAgent(channelchannel) # 3. 模拟用户输入 user_request 我想在周末去杭州玩两天预算有限主要想看看西湖和灵隐寺体验一下杭帮菜。 session_id str(uuid.uuid4())[:8] # 生成一个简短的会话ID print(f 开始处理用户请求 ) print(f用户: {user_request}) print(f会话ID: {session_id}) print(- * 50) # 4. 启动PlannerAgent print(f[主流程] 启动PlannerAgent...) # 注意PlannerAgent在执行中会通过工具发送消息到通道但不会阻塞等待回复。 plan_result planner.generate_plan(user_request, session_id) print(f[主流程] PlannerAgent初步计划完成。) print(f初步行程草案:\n{plan_result}) print(- * 50) # 5. 模拟让BudgetAgent处理通道中的请求消息并回复 print(f[主流程] 启动BudgetAgent处理请求...) replies budget_evaluator.listen_and_process() if replies: print(f[主流程] BudgetAgent已处理 {len(replies)} 条请求并发出回复。) else: print(f[主流程] 通道中未找到待处理的请求。) print(- * 50) # 6. 模拟PlannerAgent获取回复在实际中是事件驱动的 # 这里我们简化处理直接从通道中查找回复消息。 print(f[主流程] 检查BudgetAgent的回复...) feedback_messages channel.get_messages_for_agent(planner.name) budget_feedback None for msg in feedback_messages: if msg.message_type budget_review_response and msg.session_id session_id: budget_feedback msg.content.get(evaluation) break # 7. 整合最终结果 print(f\n 最终整合结果 ) print(f【用户需求】{user_request}) print(f\n【生成的行程计划】) print(plan_result) if budget_feedback: print(f\n【预算评估反馈】) print(f- 估算总费用: {budget_feedback.get(estimated_total_cost, N/A)}) print(f- 是否在预算内: {是 if budget_feedback.get(is_within_budget) else 否}) print(f- 费用分项: {budget_feedback.get(cost_breakdown, N/A)}) print(f- 优化建议: {, .join(budget_feedback.get(optimization_suggestions, []))}) else: print(\n【预算评估反馈】未收到。) print(f\n 处理完成 ) if __name__ __main__: main()7. 运行与效果验证现在运行这个程序看看两个Agent是如何协作的。运行命令cd /path/to/travel_agent_demo python main.py预期输出示例 开始处理用户请求 用户: 我想在周末去杭州玩两天预算有限主要想看看西湖和灵隐寺体验一下杭帮菜。 会话ID: a1b2c3d4 -------------------------------------------------- [主流程] 启动PlannerAgent... [PlannerAgent思考过程...] 我需要先规划一个两天的杭州行程然后咨询预算。 草案第一天西湖第二天灵隐寺餐饮推荐楼外楼... 现在调用ConsultBudgetAgent工具。 [Channel] TravelPlanner - BudgetEvaluator: budget_review_request [主流程] PlannerAgent初步计划完成。 初步行程草案: 第一天游览西湖... 第二天参观灵隐寺... 餐饮建议楼外楼、知味观... -------------------------------------------------- [主流程] 启动BudgetAgent处理请求... [Channel] BudgetEvaluator - TravelPlanner: budget_review_response [主流程] BudgetAgent已处理 1 条请求并发出回复。 -------------------------------------------------- [主流程] 检查BudgetAgent的回复... 最终整合结果 【用户需求】我想在周末去杭州玩两天预算有限主要想看看西湖和灵隐寺体验一下杭帮菜。 【生成的行程计划】 第一天游览西湖... 第二天参观灵隐寺... 餐饮建议楼外楼、知味观... 【预算评估反馈】 - 估算总费用: 约1500元 - 是否在预算内: 是 - 费用分项: 住宿经济型400元交通公交/打车200元门票灵隐寺75元餐饮两日800元其他50元。 - 优化建议: 楼外楼人均较高可考虑新白鹿餐厅等性价比更高的杭帮菜馆西湖游船可选择公共交通船而非私人包船。 处理完成 如何验证成功流程验证观察控制台输出确认出现了两次[Channel] ...的日志这表示消息成功从PlannerAgent发送到BudgetAgent并且BudgetAgent成功回复。结果验证最终输出中包含了行程计划和预算评估两部分且评估内容是针对行程草案的具体分析这表明两个Agent的信息传递是有效的。会话隔离验证你可以尝试同时运行两个不同的用户请求需要扩展主程序观察不同session_id的消息是否被正确路由和处理没有混淆。8. 常见问题与排查思路在实现A2A时你可能会遇到以下典型问题问题现象可能原因排查方式解决方案Agent A发送了消息但Agent B没反应1. 消息格式不符合接收者的预期类型。2. 通信通道未正确共享两个Agent使用了不同的通道实例。3. 接收者的监听逻辑未启动或存在Bug。1. 打印发送的消息内容检查sender,receiver,message_type是否正确。2. 检查通道的send_message和get_messages_for_agent方法是否正常工作。3. 在接收者端添加日志确认其process_incoming_message方法被触发。1. 统一并严格定义消息协议。2. 确保所有Agent注入的是同一个通道单例。3. 将监听逻辑改为事件驱动或定时轮询并确保其持续运行。消息处理顺序错乱或丢失1. 使用简易内存通道在高并发下可能丢失消息。2. 没有处理消息的确认和去重机制。1. 检查通道实现是否线程安全。2. 为消息添加唯一ID并在处理逻辑中记录已处理ID防止重复。1.生产环境务必使用专业的消息队列如RabbitMQ、Kafka或Redis Pub/Sub它们提供持久化、确认和顺序保证。2. 在消息体中增加message_id和sequence_num字段。Agent陷入循环对话Agent A等待B的回复B又向A发起新请求形成死循环。在消息内容或会话上下文中添加对话轮次计数器。1. 设计清晰的对话流程和终止条件。2. 在Agent的决策逻辑中设置最大交互轮次限制。3. 使用AgentExecutor的max_iterations参数限制单个Agent的执行步数。LLM在Tool调用时参数格式错误Tool的描述不够清晰导致LLM生成的调用参数不符合func要求的格式。查看LangChain Agent执行的详细日志verboseTrue观察LLM决定调用Tool时生成的中间动作Action。1. 优化Tool的description用更清晰的例子说明输入格式。2. 在Tool的func内部添加更健壮的参数解析和错误处理并返回友好的错误信息给LLM让其重试。会话Session混淆多个用户请求同时处理时消息和回复匹配错误。检查每个消息的session_id是否在请求和回复中保持一致。1. 在任务发起时生成全局唯一的session_id并贯穿整个工作流。2. 在通信通道中提供按session_id过滤消息的查询方法。9. 最佳实践与工程化建议Demo跑通只是第一步。要将A2A应用于实际项目你需要考虑更多工程化因素通信层抽象与中间件化不要自己造轮子Demo中的SimpleCommunicationChannel仅用于教学。真实项目应直接集成成熟的消息队列。将通信层抽象为接口如MessageBus方便后续切换实现。使用异步处理利用asyncio或框架的异步支持如LangChain的异步API让Agent在等待回复时不阻塞提高系统吞吐量。设计健壮的消息协议定义Schema使用Pydantic或dataclass严格定义消息体的结构并进行验证。包含元数据消息中除了业务数据content还应包含message_id、timestamp、version、priority等系统字段。设计错误消息类型定义标准的error消息类型用于在Agent间传递处理失败的信息。Agent的职责与状态管理单一职责每个Agent应专注于一个明确、有限的领域如规划、评估、执行。避免创建“全能”Agent。无状态设计尽可能让Agent无状态其所需的所有上下文都来自传入的消息。状态应保存在外部存储如数据库、Redis中并通过session_id关联。超时与重试为Agent间的请求-响应设置超时机制。对于可重试的错误实现指数退避的重试逻辑。可观测性与调试全链路日志为每个session_id记录完整的消息流便于追踪问题。可视化工具考虑使用像LangSmith这样的平台来跟踪和可视化复杂的多Agent工作流。监控指标收集消息延迟、处理成功率、Agent调用次数等指标。安全与权限输入验证与清理每个Agent在处理来自其他Agent或外部的输入时都应进行验证和清理防止提示词注入或非法操作。权限控制并非所有Agent都能互相通信。可以设计一个简单的权限层在消息发送前检查发送者是否有权与接收者通信。通过这个小Demo你应该对Agent间通信A2A的核心价值——解耦、异步协作、灵活编排——有了直观感受。它不是一个炫技的概念而是构建复杂、可靠AI应用系统的必要基础设施。下一步你可以尝试扩展这个Demo增加第三个Agent如“酒店预订Agent”实现真正的异步事件循环或者将内存通道替换为Redis向生产环境迈出第一步。记住清晰的协议定义和健壮的通信层是让多个AI智能体真正发挥协同智能的关键。
上一篇/下一篇内容由系统自动关联 返回资讯列表 →