Agent通信机制,Agent之间怎么交流以及消息协议设计

📅 2026/8/9 15:23:44
Agent通信机制,Agent之间怎么交流以及消息协议设计
Agent通信机制Agent之间怎么交流以及消息协议设计上周帮一个团队调多Agent系统两个Agent之间传数据经常对不上。开发Agent把代码片段发给测试Agent测试Agent拿到的字符串里混着Markdown标记解析的时候直接报错。两个人盯了半天才发现发送方和接收方对消息格式根本没约定各写各的。这就是Agent通信要解决的核心问题。多个Agent要协作得有一套大家都认的消息格式和传递规则。今天这篇就把Agent通信这件事拆开讲从几种传递方式到消息协议设计再到错误处理最后给一套能跑的完整代码。Agent通信要解决的几件事Agent之间传消息看着简单真做起来要回答三个问题。消息怎么传。A调用B是A直接喊B还是A把消息扔到一个中间地方B自己去取还是俩人读写同一块共享数据。这三种方式各有各的适用场景选错了后面全是别扭。消息长什么样。A发给B的内容B得能看懂。得约定好字段叫什么、什么类型、哪些必须有哪些可选。没有这个约定每加一个Agent就得跟所有人重新对一遍格式。传出去对方没收到怎么办。网络会断进程会挂对方处理会超时。通信层得有重试和容错不能消息发出去就当完事了。三种传递方式对比直接调用最简单。A调B的函数拿到返回值完事。写起来跟普通函数调用没区别。缺点是A和B得在同一进程里而且A得一直等B返回耦合很紧。适合Agent数量少、调用快的场景。消息队列解耦做得好。A把消息扔到队列里就不管了B有空了自己去取。A不用等BB挂了消息还在队列里等着。代价是引入了队列这个中间件调试的时候传递环节多了出了问题得查队列状态。适合Agent多、处理慢、需要异步的场景。共享状态适合那种多个Agent都要读写的公共数据。比如一个项目看板产品Agent写需求开发Agent改状态测试Agent加测试结果。大家读写同一份数据谁需要谁去取。难点在并发控制两个Agent同时改同一条记录容易打架。适合数据需要多方共同维护的场景。我自己的经验Agent数量在三个以内逻辑简单直接调用就够了。到了四五个Agent还有异步需求上消息队列。共享状态这招我用得少除非真的有那种大家都要读写的公共数据。消息格式设计格式设计我推荐用JSON加一份JSON Schema做约束。JSON好读好写所有语言都支持。Schema把字段定义钉死发送方照着填接收方照着验格式对不上当场报错。一条消息至少要有这几个字段。消息ID用来追踪发送方和接收方标识身份消息类型说明这是干什么的内容体放实际数据时间戳记录发送时间还有个可选的关联ID用来串起一组相关的消息。下面这份Schema我实际项目里在用你拿去改改就能用。# 消息格式的JSON Schema定义MESSAGE_SCHEMA{type:object,# 顶层必须是个对象required:[# 这些字段必须有缺一个就拒收message_id,from_agent,to_agent,msg_type,content,timestamp],properties:{message_id:{# 唯一标识用uuid生成方便追踪type:string,description:消息唯一ID},from_agent:{# 发送方名字比如dev_agenttype:string,description:发送方Agent标识},to_agent:{# 接收方名字test_agent或*表示广播type:string,description:接收方Agent标识},msg_type:{# 消息类型约定好枚举值type:string,enum:[task,result,query,error,heartbeat],description:消息类型},content:{# 实际内容结构由msg_type决定type:object,description:消息内容体},timestamp:{# 发送时间ISO格式字符串type:string,description:发送时间戳},reply_to:{# 可选回复某条消息时填原消息IDtype:string,description:关联的原始消息ID}}}这里有个设计取舍。content字段我用了object类型而不是string。好处是结构化接收方直接按字段取值。坏处是每种消息类型的content结构不一样得额外约定。我后来给每种msg_type单独写了一份子Schema验证的时候先看msg_type再套对应的子Schema。同步还是异步同步通信就是A发完消息死等B的回复拿到结果再往下走。逻辑直观代码好写。问题是B慢的话A一直卡着资源浪费。异步通信是A发完消息就干别的去了B处理完了通过回调或者队列把结果送回来。A不会卡住吞吐量高。代价是代码逻辑碎了你得处理回调、处理超时、处理结果到达时A的上下文还在不在。我的建议Agent之间调用快、逻辑简单用同步。涉及大模型生成的步骤动不动几秒十几秒用异步。一个需求从分析到实现到测试走完整个流程中间好几个Agent串着用异步能并行处理多个需求。错误处理和重试通信层最容易出问题的地方有三个。消息发出去对方没收到对方收到了但处理报错了对方处理太慢一直不回。第一种靠重试解决。发完消息设个超时超时没确认就重发。重试次数设个上限比如3次都失败就标记为发送失败往上抛异常让调用方决定怎么办。第二种靠错误消息解决。接收方处理出错应该回一条error类型的消息把错误信息带回来。发送方收到error就知道这事没成可以重试或者走兜底逻辑。第三种靠超时和熔断解决。给每次调用设个超时时间超时就认为失败。某个Agent连续失败好几次暂时别再调它了等它恢复。下面是完整代码把上面这些设计都实现了。一个消息总线加上Agent基类支持同步和异步两种模式带重试和错误处理。importjsonimportuuidimporttimeimportqueueimportthreadingfromdatetimeimportdatetime,timezonefromtypingimportOptional,Callable,Any# ---------- 消息构造和验证 ----------defcreate_message(from_agent:str,to_agent:str,msg_type:str,content:dict,reply_to:strNone)-dict:构造一条符合Schema的消息msg{message_id:str(uuid.uuid4()),# 生成唯一IDfrom_agent:from_agent,# 记录发送方to_agent:to_agent,# 记录接收方*表示广播msg_type:msg_type,# 消息类型content:content,# 实际内容timestamp:datetime.now(timezone.utc).isoformat(),# UTC时间戳}ifreply_to:msg[reply_to]reply_to# 如果是回复带上原消息IDreturnmsgdefvalidate_message(msg:dict)-bool:简单校验消息格式缺必填字段就返回Falserequired[message_id,from_agent,to_agent,msg_type,content,timestamp]forfieldinrequired:iffieldnotinmsg:returnFalsereturnTrue# ---------- 消息总线 ----------classMessageBus:消息总线每个Agent有一个收件箱发消息就是往收件箱里塞def__init__(self):self.queues{}# Agent名 - 该Agent的队列self.lockthreading.Lock()# 保护queues字典的线程锁defregister(self,agent_name:str):注册一个Agent给它分配一个收件箱withself.lock:self.queues[agent_name]queue.Queue()defsend(self,msg:dict,timeout:float5.0)-bool:发送消息到目标Agent的收件箱带超时ifnotvalidate_message(msg):raiseValueError(消息格式不合法)# 格式不对直接拒绝targetmsg[to_agent]withself.lock:iftargetnotinself.queuesandtarget!*:returnFalse# 目标Agent不存在iftarget*:# 广播模式发给所有Agentwithself.lock:forname,qinself.queues.items():ifname!msg[from_agent]:q.put(msg)returnTrueself.queues[target].put(msg)# 点对点发送returnTruedefreceive(self,agent_name:str,timeout:float30.0)-Optional[dict]:从收件箱取消息带超时ifagent_namenotinself.queues:returnNonetry:returnself.queues[agent_name].get(timeouttimeout)exceptqueue.Empty:returnNone# 超时没消息返回None# ---------- Agent基类 ----------classBaseAgent:Agent基类封装了收发消息和重试逻辑def__init__(self,name:str,bus:MessageBus,max_retries:int3,retry_interval:float1.0):self.namename# Agent名字self.busbus# 挂着的消息总线self.max_retriesmax_retries# 最大重试次数self.retry_intervalretry_interval# 重试间隔秒数self.bus.register(name)# 在总线上注册自己defsend_and_wait(self,to_agent:str,msg_type:str,content:dict,timeout:float60.0)-Optional[dict]:同步模式发消息后等回复带重试msgcreate_message(self.name,to_agent,msg_type,content)forattemptinrange(self.max_retries):# 最多重试max_retries次self.bus.send(msg)# 发出去# 等回复reply_to要匹配我发出的message_idreplyself._wait_reply(msg[message_id],timeout)ifreplyisnotNone:returnreply# 拿到回复就返回print(f[{self.name}] 第{attempt1}次重试等待{self.retry_interval}秒)time.sleep(self.retry_interval)# 没回复等一会再试returnNone# 重试全失败返回Nonedef_wait_reply(self,msg_id:str,timeout:float)-Optional[dict]:等待匹配的回复消息过滤掉不相关的deadlinetime.time()timeoutwhiletime.time()deadline:remainingdeadline-time.time()msgself.bus.receive(self.name,timeoutremaining)ifmsgisNone:returnNone# 检查是不是对我那条消息的回复ifmsg.get(reply_to)msg_id:returnmsg# 不是回复我的放回去处理或者丢弃这里简单丢弃print(f[{self.name}] 收到无关消息类型{msg[msg_type]}已忽略)returnNonedefon_message(self,msg:dict)-Optional[dict]:收到消息后的处理逻辑子类重写这个方法# 默认实现直接回一个收到确认returncreate_message(self.name,msg[from_agent],result,{status:ok},reply_tomsg[message_id])deflisten(self):阻塞监听收件箱收到消息就处理并回复适合异步模式whileTrue:msgself.bus.receive(self.name,timeout60.0)ifmsgisNone:continueifmsg[msg_type]error:print(f[{self.name}] 收到错误消息:{msg[content]})continuereplyself.on_message(msg)# 调子类的处理逻辑ifreplyandmsg[msg_type]!result:self.bus.send(reply)# 把回复发回去# ---------- 一个具体例子 ----------classCodeAgent(BaseAgent):开发Agent收到需求就返回一段代码defon_message(self,msg:dict)-Optional[dict]:ifmsg[msg_type]task:task_descmsg[content].get(task,)# 这里假装调用大模型生成代码实际项目里换成你的LLM调用codefdef solve():\n # 实现:{task_desc}\n return donereturncreate_message(self.name,msg[from_agent],result,{code:code,language:python},reply_tomsg[message_id])returnsuper().on_message(msg)# 其他类型走默认逻辑classTestAgent(BaseAgent):测试Agent收到代码就跑测试defon_message(self,msg:dict)-Optional[dict]:ifmsg[msg_type]result:codemsg[content].get(code,)# 假装执行测试实际项目里用subprocess跑passeddef incode# 简单判断有没有函数定义returncreate_message(self.name,msg[from_agent],result,{test_passed:passed,detail:检查通过ifpassedelse没有函数定义},reply_tomsg[message_id])returnsuper().on_message(msg)# ---------- 跑起来 ----------if__name____main__:busMessageBus()devCodeAgent(dev_agent,bus)# 创建开发AgenttesterTestAgent(test_agent,bus)# 创建测试Agent# 模拟一个外部调用者让开发Agent写代码callerBaseAgent(caller,bus)# 第一步caller让dev写一个排序函数print( 第一步请求开发Agent写代码 )replycaller.send_and_wait(dev_agent,task,{task:实现一个快速排序函数})ifreply:print(f开发Agent返回代码:\n{reply[content][code]})codereply[content][code]# 第二步caller把代码发给test_agent测试print(\n 第二步请求测试Agent测试代码 )test_replycaller.send_and_wait(test_agent,result,{code:code,language:python})iftest_reply:print(f测试结果:{test_reply[content]})else:print(测试Agent没有回复可能超时了)else:print(开发Agent没有回复重试3次都失败了)效果验证运行上面这段代码你会看到三段输出。第一步打印开发Agent返回的代码第二步打印测试Agent的测试结果最后显示测试通过。如果一切正常说明消息从caller发到devdev处理完返回caller再发给testertester返回结果整条传递流程通了。判断成功的标准很简单每一步的reply都不为Nonecontent里的字段跟预期一致。如果某一步返回None看终端里打印的重试日志基本能定位是哪个Agent没回消息。常见报错有几种。消息格式不合法会抛ValueError检查必填字段是不是都填了。目标Agent不存在send方法返回False检查Agent有没有注册。一直重试都失败大概率是接收方的on_message逻辑有bug处理的时候抛异常了没回消息给on_message加个try except兜底就好。踩坑记录第一个坑广播消息死循环。最早我写广播的时候from_agent也收到了自己发的消息处理完又广播出去无限循环。后来加了判断广播时跳过from_agent问题解决。这个坑挺隐蔽的本地测试数据量小的时候感觉不到一上量消息爆炸才发现。第二个坑回复消息匹配错。多个请求并发的时候A发了消息1和B发了消息2B的回复先回来了A拿去一匹配发现reply_to对不上。我一开始的receive逻辑是来什么收什么没做过滤。后来改成按reply_to匹配收到不相关的消息先放一边等匹配的那条。如果你用队列做总线这个匹配逻辑一定要写不然并发场景下乱套。第三个坑重试导致重复处理。A发消息给BB处理完了但回复丢了A超时重发B又处理了一遍。如果是写数据库的操作重复执行可能出问题。解决办法是给消息加幂等性接收方记录已处理过的message_id重复消息直接返回上次的结果。我现在的代码里没加这层生产环境记得补上。延伸与判断这套通信机制本质是个轻量级的消息中间件换个场景也能用。比如Agent和外部系统通信把外部API包装成一个Agent注册到总线上其他Agent就能跟它通信。再比如做Agent的灰度发布新版Agent和旧版Agent都注册到总线上按比例把消息分给新旧版本慢慢切流量。局限性也有。这个总线是进程内的Agent分布在多台机器上就用不了得换Redis或者RabbitMQ做消息中间件。没有持久化进程重启消息就丢了重要消息得落库。并发量大了单个队列会成为瓶颈得考虑分队列或者上专业消息中间件。我的建议学习和原型阶段用这套进程内总线完全够用理解了通信的原理和坑上生产再换Redis或者RabbitMQ代码结构基本不用动把MessageBus的实现换掉就行。结尾Agent通信这件事核心就三件事消息怎么传、长什么样、出了问题怎么办。把这三件事想清楚选对传递方式定义好消息格式做好重试和容错多Agent协作的基础就稳了。下一篇拿这些概念上手搭一个完整的多Agent项目团队。