ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

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不用等B,B挂了消息还在队列里等着。代价是引入了队列这个中间件,调试的时候传递环节多了,出了问题得查队列状态。适合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_agent""type":"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":{# 可选,回复某条消息时填原消息ID"type":"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:str=None)->dict:"""构造一条符合Schema的消息"""msg={"message_id":str(uuid.uuid4()),# 生成唯一ID"from_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:"""简单校验消息格式,缺必填字段就返回False"""required=["message_id","from_agent","to_agent","msg_type","content","timestamp"]forfieldinrequired:iffieldnotinmsg:returnFalsereturnTrue# ---------- 消息总线 ----------classMessageBus:"""消息总线,每个Agent有一个收件箱,发消息就是往收件箱里塞"""def__init__(self):self.queues={}# Agent名 -> 该Agent的队列self.lock=threading.Lock()# 保护queues字典的线程锁defregister(self,agent_name:str):"""注册一个Agent,给它分配一个收件箱"""withself.lock:self.queues[agent_name]=queue.Queue()defsend(self,msg:dict,timeout:float=5.0)->bool:"""发送消息到目标Agent的收件箱,带超时"""ifnotvalidate_message(msg):raiseValueError("消息格式不合法")# 格式不对直接拒绝target=msg["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:float=30.0)->Optional[dict]:"""从收件箱取消息,带超时"""ifagent_namenotinself.queues:returnNonetry:returnself.queues[agent_name].get(timeout=timeout)exceptqueue.Empty:returnNone# 超时没消息返回None# ---------- Agent基类 ----------classBaseAgent:"""Agent基类,封装了收发消息和重试逻辑"""def__init__(self,name:str,bus:MessageBus,max_retries:int=3,retry_interval:float=1.0):self.name=name# Agent名字self.bus=bus# 挂着的消息总线self.max_retries=max_retries# 最大重试次数self.retry_interval=retry_interval# 重试间隔秒数self.bus.register(name)# 在总线上注册自己defsend_and_wait(self,to_agent:str,msg_type:str,content:dict,timeout:float=60.0)->Optional[dict]:"""同步模式,发消息后等回复,带重试"""msg=create_message(self.name,to_agent,msg_type,content)forattemptinrange(self.max_retries):# 最多重试max_retries次self.bus.send(msg)# 发出去# 等回复,reply_to要匹配我发出的message_idreply=self._wait_reply(msg["message_id"],timeout)ifreplyisnotNone:returnreply# 拿到回复就返回print(f"[{self.name}] 第{attempt+1}次重试,等待{self.retry_interval}秒")time.sleep(self.retry_interval)# 没回复,等一会再试returnNone# 重试全失败返回Nonedef_wait_reply(self,msg_id:str,timeout:float)->Optional[dict]:"""等待匹配的回复消息,过滤掉不相关的"""deadline=time.time()+timeoutwhiletime.time()<deadline:remaining=deadline-time.time()msg=self.bus.receive(self.name,timeout=remaining)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_to=msg["message_id"])deflisten(self):"""阻塞监听收件箱,收到消息就处理并回复,适合异步模式"""whileTrue:msg=self.bus.receive(self.name,timeout=60.0)ifmsgisNone:continueifmsg["msg_type"]=="error":print(f"[{self.name}] 收到错误消息:{msg['content']}")continuereply=self.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_desc=msg["content"].get("task","")# 这里假装调用大模型生成代码,实际项目里换成你的LLM调用code=f"def solve():\n # 实现:{task_desc}\n return 'done'"returncreate_message(self.name,msg["from_agent"],"result",{"code":code,"language":"python"},reply_to=msg["message_id"])returnsuper().on_message(msg)# 其他类型走默认逻辑classTestAgent(BaseAgent):"""测试Agent,收到代码就跑测试"""defon_message(self,msg:dict)->Optional[dict]:ifmsg["msg_type"]=="result":code=msg["content"].get("code","")# 假装执行测试,实际项目里用subprocess跑passed="def "incode# 简单判断有没有函数定义returncreate_message(self.name,msg["from_agent"],"result",{"test_passed":passed,"detail":"检查通过"ifpassedelse"没有函数定义"},reply_to=msg["message_id"])returnsuper().on_message(msg)# ---------- 跑起来 ----------if__name__=="__main__":bus=MessageBus()dev=CodeAgent("dev_agent",bus)# 创建开发Agenttester=TestAgent("test_agent",bus)# 创建测试Agent# 模拟一个外部调用者,让开发Agent写代码caller=BaseAgent("caller",bus)# 第一步,caller让dev写一个排序函数print("=== 第一步,请求开发Agent写代码 ===")reply=caller.send_and_wait("dev_agent","task",{"task":"实现一个快速排序函数"})ifreply:print(f"开发Agent返回代码:\n{reply['content']['code']}")code=reply["content"]["code"]# 第二步,caller把代码发给test_agent测试print("\n=== 第二步,请求测试Agent测试代码 ===")test_reply=caller.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发到dev,dev处理完返回,caller再发给tester,tester返回结果,整条传递流程通了。

判断成功的标准很简单,每一步的reply都不为None,content里的字段跟预期一致。如果某一步返回None,看终端里打印的重试日志,基本能定位是哪个Agent没回消息。

常见报错有几种。消息格式不合法会抛ValueError,检查必填字段是不是都填了。目标Agent不存在send方法返回False,检查Agent有没有注册。一直重试都失败,大概率是接收方的on_message逻辑有bug,处理的时候抛异常了没回消息,给on_message加个try except兜底就好。

踩坑记录

第一个坑,广播消息死循环。最早我写广播的时候,from_agent也收到了自己发的消息,处理完又广播出去,无限循环。后来加了判断,广播时跳过from_agent,问题解决。这个坑挺隐蔽的,本地测试数据量小的时候感觉不到,一上量消息爆炸才发现。

第二个坑,回复消息匹配错。多个请求并发的时候,A发了消息1和B发了消息2,B的回复先回来了,A拿去一匹配发现reply_to对不上。我一开始的receive逻辑是来什么收什么,没做过滤。后来改成按reply_to匹配,收到不相关的消息先放一边,等匹配的那条。如果你用队列做总线,这个匹配逻辑一定要写,不然并发场景下乱套。

第三个坑,重试导致重复处理。A发消息给B,B处理完了但回复丢了,A超时重发,B又处理了一遍。如果是写数据库的操作,重复执行可能出问题。解决办法是给消息加幂等性,接收方记录已处理过的message_id,重复消息直接返回上次的结果。我现在的代码里没加这层,生产环境记得补上。

延伸与判断

这套通信机制本质是个轻量级的消息中间件,换个场景也能用。比如Agent和外部系统通信,把外部API包装成一个Agent注册到总线上,其他Agent就能跟它通信。再比如做Agent的灰度发布,新版Agent和旧版Agent都注册到总线上,按比例把消息分给新旧版本,慢慢切流量。

局限性也有。这个总线是进程内的,Agent分布在多台机器上就用不了,得换Redis或者RabbitMQ做消息中间件。没有持久化,进程重启消息就丢了,重要消息得落库。并发量大了单个队列会成为瓶颈,得考虑分队列或者上专业消息中间件。

我的建议,学习和原型阶段用这套进程内总线完全够用,理解了通信的原理和坑,上生产再换Redis或者RabbitMQ,代码结构基本不用动,把MessageBus的实现换掉就行。

结尾

Agent通信这件事,核心就三件事,消息怎么传、长什么样、出了问题怎么办。把这三件事想清楚,选对传递方式,定义好消息格式,做好重试和容错,多Agent协作的基础就稳了。下一篇拿这些概念上手搭一个完整的多Agent项目团队。

返回列表