ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Agent点对点通信协议实战:hermes peer架构与可靠性设计

Agent点对点通信协议实战:hermes peer架构与可靠性设计 做了好几年Agent开发最让我头疼的从来不是模型怎么调而是Agent之间怎么好好说话。早期项目里每个Agent都连着消息队列和中控服务链路越拉越长日志散落在五六个系统里排查一次超时问题能从下午翻到天黑。后来接触到 hermes peer 这套点对点通信协议算是把这块理顺了不少Agent之间直接交换消息业务流量不再绕道中心枢纽链路变短、延迟降低、排查也直观。这篇文章不打算做概念堆砌我会直接从协议设计思路讲起拆解 hermes peer 的消息格式、节点发现、握手和可靠性机制再带一个全栈协作案例完整演示三个Agent如何配合跑通一条订单处理链路。适合正在做Agent框架选型、或者想自己设计一套轻量通信协议的开发者参考用得上就拿走。1. 为什么Agent之间需要一套点对点通信协议从中心化调度的痛点说起1.1 中心化调度撑不住的协作场景很多Agent项目起步阶段都走中心化调度一个总控Agent负责接收任务拆解子任务分发给其他Worker Agent再汇总结果。这么做的好处是逻辑好写、状态统一问题出在规模上来之后。先说一个我实际遇到的数字。某个内部系统里接了大概五十个Agent每个Agent承担不同的领域职能有的做意图识别有的做翻译有的做数据库查询。所有消息都经过中心消息队列中转。一开始没觉得有什么问题但到了业务高峰期队列积压能把平均响应时间从200毫秒拉到3秒而且只要中控进程一抖动整条链路直接断掉所有Agent集体失联。更麻烦的是Agent之间如果要做一次多轮协作每轮都要经过中控转发一条看似简单的对话背后可能走了七八次队列跳转。这种拓扑本质上就是一个星型结构中控既要管任务编排又要管消息路由还要对外提供状态查询接口任何一项负载过高都会拖垮整体。点对点通信的思路很简单把消息通道从中心节点里解放出来每个Agent既是客户端又是服务端A想跟B说话直接建立连接不用什么消息都跑到中控那里绕一圈。这也是 hermes peer 最核心的设计立场不做Agent业务逻辑只做Agent之间的那条“直连通道”。1.2 hermes peer 的定位与核心设计目标第一次看到 hermes peer 这个名字我以为是某个信使服务后来发现还挺贴切——Hermes 在希腊神话里本来就是传递讯息的信使。它的定位非常聚焦解决Agent与Agent之间的传输问题不掺和Agent内部怎么思考、怎么调用模型。换句话说Agent框架负责“大脑”hermes peer 负责“神经系统”。设计目标可以归纳为三条。第一是轻量协议不会强制依赖外部中间件一个Agent进程内嵌协议栈就能工作部署成本很低。第二是去中心化业务消息不经过中心服务转发A和B之间是真正意义上的点到点。第三是可观测每条消息都有唯一ID和完整链路信息排查问题的时候能像查快递一样把消息路径串起来。我自己在选型的时候最看中的就是最后一点。Agent项目跟普通Web服务不一样消息在多个智能体之间流转一旦出问题没有消息级别的追踪能力几乎没法定位到底是谁丢了消息、谁处理超时、谁返回了错误结果。hermes peer 在这块的设计思路是把“信使”这件事做透协议层就带上追踪、确认、重试这些机制。2. hermes peer协议拆解消息格式、节点发现与可靠性设计2.1 统一消息信封为什么每条消息都要包一层点对点通信第一个要解决的是消息格式。Agent的类型五花八门有的是Python写的有的是Node.js写的有的跑在容器里有的跑在边缘设备上。你不能要求各个Agent用同一种语言的数据结构互相调用所以协议必须定义一种跨语言的消息信封。hermes peer 的做法是设计一个通用的Envelope信封结构业务数据装在Payload里通信元数据放在信封外层。一个典型的信封长这样{ message_id: a3f2c8e1-9b4d-4f6a-b7e2-1c8d3f5a2e90, sender_id: agent-order-collector, receiver_id: agent-inventory, message_type: order.created, sequence: 12, timestamp: 1713000000123, ttl: 60000, payload: { order_id: ORD-20240413-001, items: [ {sku: SKU-1001, quantity: 2} ] } }message_id 是全链路唯一的消息IDsender_id 和 receiver_id 标记通信双方message_type 表示业务语义sequence 用来处理乱序ttl 是存活时间超过这个时间还没被处理完的消息可以直接丢弃。为什么非要加这一层统一信封我举一个实际例子。项目里有两个Agent一个返回JSON带下划线字段另一个内部用驼峰命名如果没有统一信封这两个Agent接口对接的时候光是字段映射就要写一堆胶水代码。有了信封业务方只需要关注 message_type 和 payload通信的通用能力去重、追踪、过期全部由协议层承担各Agent内部爱用什么字段风格都无所谓。从维护角度看统一信封最大的价值在于可追踪。每一条消息都有message_id从发送到确认的完整路径都能记录。线上排查问题的时候只需要grep这个ID所有环节的日志就都串起来了不用再去猜是谁发出的消息。2.2 节点发现与寻址不靠中心也能找到彼此点对点通信有个天然问题A怎么知道B在哪里Web服务可以查注册中心但hermes peer 的设计理念是尽量不去依赖中心化组件否则又绕回星型结构的老路。它提供了三种节点发现方式按场景选。第一种是静态配置适合节点少、地址固定的私有化部署。在每个Agent的配置里直接写清楚对端节点的主机和端口例如 agent-inventory 的地址是 192.168.1.20:8201。这种方式最简单适合刚起步的项目。第二种是动态发现基于 mDNS/DNS-SD 协议Agent启动时广播自己的能力其他Agent在同一个局域网内可以自动发现它。这类似你用手机投屏到电视电视广播“我能投屏”手机收到广播就把它列出来。好处是零配置坏处是只限于局域网跨网段就不行了。第三种是指定引导节点Bootstrap Node。注意引导节点不做业务消息的中转它的职责只是在Agent启动时告诉大家“谁在哪”类似电话簿而不是电话交换机。Agent启动后先问引导节点要一份节点列表然后自己跟目标Agent直连后续业务消息完全不经过引导节点。这个方案兼顾了部署灵活性和去中心化是我在实际项目里最常用的方式。寻址机制上hermes peer 用“能力名 节点ID”两层结构。每个Agent启动时声明自己有哪些能力比如 inventory.check、order.create。发送方不需要关心目标Agent具体在哪个IP只需要声明“我要给谁发”协议层根据能力名解析出目标节点再建立连接。这个设计跟微服务里的服务发现思路一脉相承但更轻量没有引入一套完整的服务网格。2.3 握手、心跳与连接管理可靠性的第一道防线协议不是HTTP那种简化的一问一答为了在不可靠的网络环境里维持长连接hermes peer 设计了一套连接生命周期管理机制。建立连接时发起方发一个 HELLO 包包含自己的节点ID、能力列表和协议版本。接收方收到后回 HELLO_ACK确认协议版本兼容并返回自己的能力列表。这里有个容易忽略的细节能力列表的交换非常关键。如果A要给B发 inventory.check 消息但B根本不具备这个能力握手阶段就能发现直接报错而不是等业务消息发过去之后才收到“无法处理”的响应。这个设计把很多问题提前暴露在连接建立阶段比运行时再去探测要省事得多。连接建立之后双方每30秒互发一次心跳包 PING/PONG连续3次没有收到对端心跳就判定节点离线触发重连和告警。心跳间隔这个参数其实很有讲究设太短会浪费带宽设太长又会让故障发现变慢。我之前在自己的环境里测过一组对比数据心跳间隔故障发现延迟单节点额外流量24小时5秒最快15秒感知约34KB30秒最快90秒感知约5.7KB60秒最快180秒感知约2.8KB对于一个内部Agent集群来说30秒是比较平衡的选择故障发现延迟可以接受流量开销也几乎可以忽略。如果你的业务对故障恢复要求极高可以收紧到10秒代价是多一点点心跳流量问题不大。2.4 消息确认、幂等去重与乱序处理协议层的可靠性设计点对点传输和HTTP请求一样面对不可靠网络消息可能丢失、重复或者乱序。hermes peer 在协议层内置了三套机制来解决这些问题。消息确认机制。发送方发出业务消息后接收方处理完要回一个ACK。注意协议有两种确认策略。一种是“收到即ACK”接收方把消息放进处理队列后就确认适合消息量大的场景另一种是“处理完再ACK”等业务逻辑执行成功后才确认适合对可靠性要求极高的资金、订单类场景。如果发送方在超时时间内没收到ACK会按指数退避策略重发第一次1秒后重试第二次2秒第三次4秒最多重试5次。超过重试次数消息进入死信队列由人工或补偿任务处理。幂等去重机制。网络重试会导致同一个消息被接收方收到多次。hermes peer 在设计上要求每个接收方维护一个最近已处理消息ID的集合。收到消息后先查这个集合如果已经处理过就直接丢掉。这个机制对库存扣减这种操作尤其重要——一条消息被重发两次不能扣两次库存。乱序处理机制。协议给每条消息都带了一个 sequence 序号接收方维护一个滑动窗口只接受窗口内的序号。不过我这里要说句实在话Agent之间的协作大多数是事件驱动的两个独立事件之间的前后关系未必那么严格。如果业务里确实存在强顺序依赖比如必须先“创建订单”再“扣库存”我的建议是把这种顺序依赖放在业务逻辑层面控制而不是完全依赖协议层帮你排序。协议层的乱序处理是兜底不能当主方案用。3. 全栈协作实战三Agent订单处理链路的完整落地3.1 场景设计与Agent职责划分理论说得再多不如看一个能跑的案例。我以一个拼团电商的订单自动处理作为场景设计一条完整链路用户下单之后系统需要检查库存、冻结库存、生成发货单然后再发通知给用户。这条链路我拆成三个AgentAgent A订单采集Agent接收前端提交的订单数据做格式校验和清洗生成标准订单然后发给库存Agent。 Agent B库存Agent检查每个SKU的库存是否足够够的话冻结库存生成发货单再通知Agent C。 Agent C通知Agent收到发货单后负责发送站内信和短信通知。为什么这样划分关键原则是“每个Agent只干一类事Agent之间通过消息协作不直接调用对方内部方法”。订单Agent不关心库存怎么算它只负责把标准化订单发出去库存Agent不关心用户用什么渠道接收通知它只负责把发货结果发出去。职责边界越清晰Agent之间的消息类型就越稳定后续替换任何一个Agent的实现都不影响其他部分。3.2 消息类型与全链路消息流定义三个Agent之间需要定义四个消息类型消息类型发送方接收方语义order.createdAgent AAgent B新订单已创建请求库存处理inventory.processedAgent BAgent A库存处理完成回传结果fulfillment.generatedAgent BAgent C发货单已生成请求发送通知notification.sentAgent CAgent B通知发送结果回执链路是这样走的Agent A 收到下单请求发出 order.createdAgent B 监听并处理成功后回 inventory.processed 给 Agent A同时发出 fulfillment.generated 给 Agent CAgent C 发送通知回通知结果给 Agent B。每个消息都带 message_id整条链路可以串成一条追踪链。这里有个设计细节值得分享Agent A 和 Agent B 之间是同步等待关系Agent A 发出 order.created 后必须等 Agent B 返回库存处理结果才能给前端一个明确答复。但 Agent B 和 Agent C 之间是异步关系B 发出 fulfillment.generated 之后不需要等 C 返回结果就可以继续处理下一个订单。同步和异步混合使用既保证了核心流程的确定性又避免了不必要的阻塞。3.3 核心代码实现发送、接收与回调处理我自己用的是Python实现代码写得很精简核心就三个模块消息定义、发送封装、接收回调。# message.py import json import time import uuid class Envelope: def __init__(self, sender_id, receiver_id, message_type, payload, ttl60000): self.message_id str(uuid.uuid4()) self.sender_id sender_id self.receiver_id receiver_id self.message_type message_type self.payload payload self.timestamp int(time.time() * 1000) self.ttl ttl def to_json(self): return json.dumps({ message_id: self.message_id, sender_id: self.sender_id, receiver_id: self.receiver_id, message_type: self.message_type, timestamp: self.timestamp, ttl: self.ttl, payload: self.payload }) staticmethod def from_json(data): obj json.loads(data) return Envelope( obj[sender_id], obj[receiver_id], obj[message_type], obj[payload], obj.get(ttl, 60000) )发送端封装# sender.py import asyncio from hermes_peer import PeerClient class OrderSender: def __init__(self, peer_client): self.client peer_client self.pending_orders {} async def send_order(self, order_data): env Envelope( sender_idagent-order-collector, receiver_idagent-inventory, message_typeorder.created, payloadorder_data ) # 注册等待回调收到 inventory.processed 后解除阻塞 future asyncio.get_event_loop().create_future() self.pending_orders[env.message_id] future await self.client.send(env) # 等待库存处理结果最多等10秒 result await asyncio.wait_for(future, timeout10.0) return result def handle_inventory_result(self, env): 接收方返回的 inventory.processed 回调 if env.message_id in self.pending_orders: future self.pending_orders.pop(env.message_id) future.set_result(env.payload)接收端监听# receiver.py from hermes_peer import PeerServer class InventoryAgent: def __init__(self, peer_server): self.server peer_server self.server.register_handler(order.created, self.handle_order) self.server.register_handler(notification.sent, self.handle_notification) async def handle_order(self, env): 处理订单创建消息 order env.payload # 1. 校验库存 ok await self.check_stock(order[items]) if not ok: return self.reply_error(env, 库存不足) # 2. 冻结库存 await self.freeze_stock(order[items]) # 3. 生成发货单 fulfillment await self.generate_fulfillment(order) # 4. 回传库存处理结果给订单Agent await self.server.send(Envelope( sender_idagent-inventory, receiver_idenv.sender_id, message_typeinventory.processed, payload{order_id: order[order_id], status: success} )) # 5. 通知Agent C生成发货通知 await self.server.send(Envelope( sender_idagent-inventory, receiver_idagent-notification, message_typefulfillment.generated, payload{fulfillment_id: fulfillment[id], user_id: order[user_id]} )) return True这里要特别说明一个我踩过的坑回调处理函数必须是幂等的。我在早期版本里把“冻结库存”的操作放在了消息处理函数里后来因为网络抖动一条 order.created 被重发了两次库存直接冻了两遍最后靠对账才发现。正确做法是把“冻结库存”操作设计成幂等的同一个订单ID冻结两次不会产生副作用这样即使消息重复到达也不会出错。3.4 前端可视化与后端编排一个“全栈”协作系统长什么样Agent协作系统光有后端协议还不够用户和运维都希望能直观看到任务跑到哪一步了。我在这套系统里加了一个轻量级前端面板完整链路是前端页面通过 WebSocket 连接后端编排服务后端编排服务订阅 hermes peer 的消息事件把链路上每个步骤推给前端。前端我用的是 Vue 3 WebSocket组件展示三列分别对应三个Agent的状态卡片。每一列显示当前Agent处理的消息数量、最近处理的message_id、当前状态空闲/忙碌/异常。链路流转的时候卡片之间有流动动画能很直观地看到消息从订单Agent跳到库存Agent再到通知Agent。template div classagent-board AgentCard v-foragent in agents :keyagent.node_id :nameagent.name :statusagent.status :message-countagent.message_count :last-message-idagent.last_message_id / /div /template script setup import { onMounted, onUnmounted, ref } from vue const agents ref([]) let ws null onMounted(() { ws new WebSocket(ws://localhost:8080/ws/agent-status) ws.onmessage (event) { const data JSON.parse(event.data) // 更新对应Agent的状态卡片 const idx agents.value.findIndex(a a.node_id data.node_id) if (idx ! -1) { agents.value[idx] { ...agents.value[idx], ...data } } } }) onUnmounted(() { ws?.close() }) /script后端编排服务是Node.js写的订阅协议层抛出的状态事件再通过Socket.IO推给前端。关键代码如下const { PeerNode } require(hermes-peer-sdk) const { Server } require(socket.io) const peer new PeerNode({ nodeId: agent-orchestrator, bootstrapUrl: http://localhost:9000/bootstrap }) const io new Server(8080, { cors: { origin: * } }) peer.on(message.processed, (message) { // 消息处理成功推送给前端 io.emit(agent-status, { node_id: message.receiver_id, status: processed, last_message_id: message.message_id, message_count: message.sequence }) }) peer.on(message.failed, (message, error) { io.emit(agent-status, { node_id: message.receiver_id, status: error, last_message_id: message.message_id, error: error.message }) })这个“全栈”体现在整条链路的打通协议层负责Agent间通信后端服务负责状态聚合与事件推送前端负责可视化呈现。任一层独立替换都不影响其他层这也是我比较推荐的分层设计方式。如果你只想看文本日志后端把同一套事件处理逻辑打到终端或者日志收集系统就行前端不是必须的。4. 部署与排查实录hermes peer落地中的常见坑与对策4.1 节点之间连不上先分清网络层还是协议层问题点对点通信最常遇到的第一类问题是“两个Agent进程都启动了但就是连不上”。我排查这个问题的思路是自上而下分层排。第一步查网络层。两个节点之间是不是真的能通信用 telnet 或者 nc 直接测目标端口是否可达。很多情况是防火墙策略挡了端口或者云安全组忘了放行。一个Agent既要做服务端监听端口又要做客户端去连接别人所以防火墙至少要放行它监听的那个端口。第二步查协议握手。如果端口通了但连不上抓包看HELLO包有没有发出去、有没有收到HELLO_ACK。如果HELLO发出去了没响应大概率是协议版本不兼容或者对端的安全校验拒绝了连接。升级Agent版本的时候特别容易出这个问题两个节点版本不一致握手直接失败。第三步查引导节点。用了引导节点模式的话看看Agent启动时是否成功拿到了节点列表。实际工作中我发现一个很容易忽略的场景Agent A 的静态配置里写的是内网IP但 Agent B 跑在另一个网络环境A 拿到 B 的地址之后发现根本路由不过去。解决方式是在引导节点上注册对外可达的地址而不是容器内网地址。4.2 消息丢失与重试风暴一次线上事故的复盘有一次线上系统出现了一个很有意思的现象订单Agent发出的库存检查消息总是超时日志里全是重试记录库存Agent那台机器的CPU被打满消息队列堆积了几万条。查了半天才发现问题出在确认机制的配置上。库存Agent配置成了“收到即ACK”消息一进处理队列就回ACK订单Agent以为处理成功了于是继续发下一批消息。但库存Agent其实已经积压了大量未处理的消息处理线程池被打满新的消息排队时间越来越长慢慢就变成了雪崩。重试机制在这个场景下反而帮了倒忙每一条超时重发的消息又堆回队列队列越来越长处理越来越慢。这件事给我的教训很简单确认机制的选择必须跟业务特性匹配。如果下游处理能力是瓶颈就一定要用“处理完再ACK”让上游感知到下游的真实负载。同时要给消息设置合理的TTL不要允许消息无限期重试超过阈值的消息直接进死信队列留给人工处理。流控机制也很重要可以在协议层加一个简单的令牌桶限制发送方单位时间内的消息量避免一瞬间把对端打爆。4.3 消息重复怎么设计一个真正幂等的处理函数点对点通信的超时重试机制必然带来消息重复的问题。网络抖动一下同一条消息发了两遍如果接收方处理函数不是幂等的就会产生严重副作用。举个最典型的场景订单Agent发出 order.created库存Agent处理完了但ACK在网络里丢了订单Agent超时后重发同一条消息。库存Agent收到重复消息如果直接把库存再冻结一遍账面就对不上了。我在实践里总结了一套幂等设计三板斧。第一板斧接收方必须做消息ID去重已经处理过的message_id直接跳过。第二板斧业务操作尽量用业务单号做唯一约束比如订单ID在库存冻结表里建唯一索引重复插入直接报错被捕获不会产生第二次扣减。第三板斧对于不能天然幂等的操作比如发送短信通知可以用“处理状态表”记录该订单是否已发送过通知发送前先查状态。这套组合方案在大部分Agent协作场景下都够用了。不要追求纯理论上的绝对幂等工程上能做到“重复消息不产生重复业务结果”就是合格。4.4 Agent状态不一致如何设计与处理补偿机制点对点通信没有中心事务管理器三个Agent协作过程中如果中间某个环节失败就可能出现状态不一致。最常见的例子库存Agent扣了库存、发了 fulfillment.generated但通知Agent发送短信失败。这时候用户的单子有发货单却没收到通知链路状态是碎的。解决这个问题不能依赖单一Agent要在设计阶段就想好补偿策略。我的做法是引入一个“反向消息”机制。通知Agent发送失败后不再沉默而是发一条 notification.failed 给库存Agent。库存Agent收到之后回滚之前冻结的库存并给订单Agent发一条 order.failed订单Agent把订单标记为处理异常由人工介入或者后续自动重试。class NotificationAgent: async def handle_fulfillment_generated(self, env): try: await send_sms(env.payload[user_id], env.payload[fulfillment_id]) await self.server.send(Envelope( sender_idagent-notification, receiver_idagent-inventory, message_typenotification.sent, payload{fulfillment_id: env.payload[fulfillment_id], status: success} )) except Exception as e: # 关键失败也要告诉上游触发补偿 await self.server.send(Envelope( sender_idagent-notification, receiver_idagent-inventory, message_typenotification.failed, payload{fulfillment_id: env.payload[fulfillment_id], reason: str(e)} ))这个设计的关键在于失败不是一个“终点事件”而是链路上一个可以被补偿的消息。Agent之间通过消息互相传递结果包括失败结果整条链路才能自愈。这种模式比中心化事务要灵活得多也不难落地——无非是每个Agent多处理一种“失败消息类型”而已。4.5 常见问题速查表症状可能原因排查思路节点之间ping不通防火墙、安全组、NAT配置先用nc测端口再抓包看包是否到达能ping通但业务消息收不到握手阶段能力列表不匹配检查HELLO/HELLO_ACK日志核对消息类型注册消息间歇性丢失ACK超时时间设置过短延长处理超时时间检查处理线程池负载消息重复处理处理函数不是幂等的加message_id去重、业务单号唯一约束链路状态不一致缺少失败通知与补偿机制增加失败消息类型实现回滚逻辑重试导致消息堆积确认机制不合理缺少流控改“处理完再ACK”限制单位时间发送量我在实际落地 hermes peer 这套方案之后最大的感受是排查链路问题的成本明显降下来了。以前一条消息丢了要翻三个系统的日志去猜现在拿着 message_id 一查从谁发出、到谁处理、在哪里失败一清二楚。如果你也在做Agent协作我真心建议先把点对点通信这套基础扎扎实实搞明白。不用一上来就上一个特别重的框架先理清消息格式、节点发现、确认重试、幂等去重这几个核心问题再根据业务场景引入具体的协议实现很多看似复杂的问题都会变得可控得多。最后再分享一个小技巧Agent之间协作不要追求所有消息都能100%送达而是要保证每条消息的“结果”都有处可去、有据可查这样系统在不可靠网络上才能活得久。
返回列表