资讯动态

消息代理层设计与实践:Hermes-Agent 如何统一消息路由与投递

发布时间:2026/9/9 2:10:25 来源:尧图企业网站定制
1. Hermes这个名字透漏出来的定位我为什么需要一个消息代理层1.1 从信使这个意象说起Hermes 在希腊神话里是奥林匹斯的信使负责在各个神域之间传递信息。当时我拿到这个项目标题时第一反应是——这个名字其实已经把核心定位说得很明白了它不是一个业务系统不是一个数据平台而是一个负责把消息可靠地从生产者送达消费者的中间代理层。换句话说你的服务之间如果已经出现了消息满天飞、回调绕成迷宫的情况那你就需要一个 hermes-agent 这样的角色来收口。这类模块在实际工程里常见的落点有几个作为轻量消息代理屏蔽底层 MQ 差异、作为跨服务异步事件的路由分发层、作为定时任务/延迟任务的调度代理甚至可以作为内部 API 网关的补充对请求做协议转换和动态路由。我这边实践的版本更偏向第一种和第二种的合体——一个挂在业务服务和基础消息组件之间的薄薄一层把消息发到哪、怎么发、失败怎么办这件事统一收口。1.2 我接手时的系统现状接手之前我所在的系统大概处于这种状态业务服务直接调用各种 SDK 去连消息队列有的用 RocketMQ有的用 Kafka甚至有两条核心链路还在用 Redis 的 List 做简单消息队列。生产者发送消息的姿势五花八门有人同步发有人异步发有人不配置重试有人把消息体直接序列化成 JSON 塞进去然后消费者那边又各自解析一遍。表面上看系统还能跑但一旦出现消息堆积或者消费失败排查链路的路径长到让人崩溃——你得从生产者一路查到消费者中间隔着三套不同的消息组件、至少五六种消息体格式以及各团队完全不一致的命名规范。hermes-agent 要解决的核心痛点就是这件事对所有上游提供统一的生产者入口对所有下游提供统一的消费者接入模型把路由、序列化、重试、死信、监控这些横切能力全部沉淀到代理层。于是业务团队不再关心这条消息走的是 Kafka 还是 RocketMQ只需要关心我发出去的是什么语义的事件。下面我会按我自己实际落地这个模块的顺序把架构设计、核心 API、踩坑记录和调优经验完整拆开讲。2. hermes-agent 的核心架构拆解从收件到投递的完整链路2.1 路由层发件匹配规则怎么设计代理层的第一步永远是这条消息该去哪。我最初在设计路由时犯过一个典型错误——把路由规则写死在代码里每次新接入一个业务方都要改代码重新发版这完全违背了代理层的初衷。后面重构为规则驱动之后整个模块才真正变得好用。具体来说每一条消息都带一个逻辑主题Topic业务方只认这个逻辑主题不感知物理目的地。路由层维护一张规则表大致结构是规则项说明示例值逻辑主题业务方发送时使用的主题名order.created物理目标实际的消息组件类型和主题名rocketmq:order_topic_prod路由属性根据消息头或消息体字段做条件路由regioncn-hangzhou投递策略同步投递、异步投递或批量投递async-batch降级目标主目标不可用时的备选路径local-file-fallback路由匹配的顺序也有讲究。我采用的是先精确匹配、再前缀匹配、最后默认路由的三级策略避免业务方在规则配置上花太多心思去学习匹配语法。这个设计和计算机网络里的路由表逻辑是相通的——精确匹配保证最高优先级默认路由兜底保证不会出现消息因找不到路由而被静默丢弃。2.2 队列与并发模型为什么吞吐和时延能兼顾路由确定了目的地之后接下来就是投递。hermes-agent 在内存中为每个逻辑主题维护一组待发送队列队列底层采用有界环形缓冲避免生产者过快导致内存无限增长。发送线程池的大小、队列的容量、单次批量拉取的条数这三个参数是性能的关键变量。我在压测中常用的初始配置是每个主题队列容量 10000发送线程数为 CPU 核数的两倍单次批量拉取 200 条消息。这套配置在单主题每秒 5000 条消息体大小约 1KB 的场景下可以做到平均投递时延 8msP99 时延 35ms 左右。如果你需要更高的吞吐优先调大批量拉取条数而不是线程数因为线程数过多反而会引起锁竞争和上下文切换开销。这里有个很多人容易忽视的点消费确认模型。生产者把消息投递到 MQ 成功之后hermes-agent 并不认为整个流程结束了它会等待消费者消费成功的确认信号。如果消费失败消息要回到本地缓冲重试队列而不是直接让 MQ 的重试机制介入。这种本地先重试、本地不行再交给 MQ 重试的分层策略避免了无意义的事务外重试和重复消费风暴。2.3 失败兜底本地缓冲、重试和死信代理层必须做好最坏情况的准备。外部 MQ 组件不是 100% 可靠的网络抖动、Broker 重启、磁盘写满任何故障都可能让投递失败。hermes-agent 的兜底设计分为三层。第一层是内存重试队列投递失败的消息会进入带退避策略的重试队列初始延迟 1 秒每次失败重试间隔指数递增最大重试间隔不超过 5 分钟。第二层是本地磁盘缓冲如果内存重试次数超过阈值默认 5 次消息会追加写入本地文件由后台进程定期扫描并重新投递。第三层才是真正意义上的死信——如果本地文件里的消息超过 24 小时仍未投递成功系统会把这个消息写入死信表并告警通知到值班人。这三层兜底设计里第二层是最容易被忽略但实际最管用的。我经历过一次 RocketMQ 集群全挂的情况持续了将近半小时。由于本地磁盘缓冲的存在服务重启后消息自动续投递业务方全程没有感知到消息丢失。所以如果你在自研消息代理层本地缓冲那二十行代码绝对值得认真写。3. 把它落到代码里配置接入与核心 API 实战3.1 一份典型配置文件的逐行拆解先给出一份 hermes-agent 的配置文件这份配置我经过了多次线上验证可以直接作为起步模板hermes: routes: - logical-topic: order.created physical-target: rocketmq:order_topic_prod delivery-policy: async-batch fallback-target: local-file-fallback - logical-topic: order.paid physical-target: rocketmq:order_paid_topic_prod delivery-policy: sync transport: thread-pool-size: 16 queue-capacity: 10000 batch-size: 200 retry-initial-delay-ms: 1000 retry-max-delay-ms: 300000 retry-max-attempts: 5 consumer: auto-ack: false ack-timeout-ms: 60000 fallback: local-file-path: ./data/hermes-fallback scan-interval-ms: 60000 max-pending-hours: 24配置里几个关键项我解释一下。delivery-policy: async-batch告诉代理该逻辑主题使用异步批量投递模式适用于对时延不敏感但吞吐量大的场景比如日志采集、行为上报delivery-policy: sync适用于需要确认消息已到达服务端的场景比如订单支付成功后的通知类消息。auto-ack: false强制消费者手动确认消费完成这能避免消息还在处理中就被 ack导致异常时丢失。ack-timeout-ms: 60000是消费超时阈值超过这个时间未 ackhermes-agent 会判定消费失败并触发重投递。3.2 核心接口的使用三个必会姿势第一个是同步发送。这适合对结果有强感知的场景比如用户在页面上触发了某个操作需要立刻知道操作结果。HermesAgent agent HermesAgent.builder() .configPath(hermes-agent.yml) .build(); SendResult result agent.send(order.created, Jsons.of(orderId, 20240911001, amount, 399.00, region, cn-hangzhou)); if (result.isSuccess()) { // 业务侧记录日志即可 } else { // 这里的失败意味着代理层已经尝试过本地重试仍然失败 // 需要业务侧决定是否转人工处理 }第二个是异步监听消费。业务方只需要在处理方法上标注HermesListener注解指定监听的逻辑主题代理层会自动完成消息反序列化、消费确认和异常处理。public class OrderNotifyConsumer { HermesListener(queue order.paid) public void onOrderPaid(OrderPaidEvent event) { // 这里写业务逻辑发短信、推送 App 通知、更新订单状态等 // 方法正常返回hermes-agent 自动 ack // 方法抛异常hermes-agent 捕获并触发重投递 } }这里有个细节值得注意消费方法一定要保证幂等。因为 hermes-agent 在消费超时或异常时会重新投递同一个事件你的业务方法很可能被调用多次。上例中如果一个订单支付成功通知被消费了两次你不能给用户发两次短信。所以监听方法内第一件事应该是查重或加分布式锁。第三个是批量消费。适合高吞吐场景比如收集埋点日志后批量写入数仓HermesListener(queue user.behavior.tracked, batch true) public void onUserBehaviorTracked(ListUserBehaviorEvent events) { ListLogRecord records events.stream() .map(e - new LogRecord(e.getUserId(), e.getAction(), e.getTimestamp())) .collect(Collectors.toList()); logBatchSaver.save(records); }批量消费模式下hermes-agent 会攒够一定数量由batch-size配置控制或者等待超过一个时间窗口后主动推送一批事件给消费者。这样避免了单条处理带来的频繁网络开销和线程上下文切换。3.3 事务内发送与异步确保一个容易踩坑的组合在实际业务里更新数据库和发送消息这两件事经常需要放在同一个事务里。很多新手会这样写Transactional public void createOrder(OrderDTO order) { orderDao.insert(order); agent.send(order.created, Jsons.of(orderId, order.getId())); }这个写法存在隐患。事务是在方法返回后才提交的如果agent.send在事务提交之前就把消息发出去消费者收到消息后去查订单数据很可能还查不到。反过来如果消息发送成功但事务最终回滚了消费者就会收到一条根本不存在的订单事件。正确的做法是利用 hermes-agent 提供的事务发送器TransactionalSender它内部会先注册一条待发送消息等事务提交后再真正投递事务回滚则自动丢弃该消息。public void createOrder(OrderDTO order) { transactionTemplate.execute(status - { orderDao.insert(order); agent.sendTransactional(order.created, Jsons.of(orderId, order.getId())); }); }如果你在用的消息组件本身支持事务消息比如 RocketMQhermes-agent 会优先走底层的事务消息机制如果不支持代理层会在本地维护一张待发送表配合本地消息表 定时兜底扫描来实现事务一致性。无论走哪种机制业务方只需要记住一点凡是先写库再发消息的地方一律使用sendTransactional不要用裸的send。4. 上线后踩过的坑和帮别人排掉的坑4.1 消息丢失的三个隐蔽位置代理层上线前问题往往出在用起来不顺上线后问题则集中在消息不见了。第一类消息丢失发生在消费者反序列化阶段。老系统里有些消息体不是标准 JSON有的是Map嵌套List再嵌套自定义对象有的甚至直接塞了 Java 原生序列化后的byte[]。一旦 hermes-agent 的反序列化器遇到不兼容的格式默认配置下会丢弃消息并记录 error 日志不会走死信逻辑。我的修复方案是给自定义配置添加容错型反序列化器解析失败的消息进入死信队列并带上原始字节内容而不是直接丢弃。不要相信错误消息丢了就丢了这种说法线上 90% 的幽灵问题根源都是某条消息被静默丢弃导致下游数据不一致。第二类消息丢失发生在消费者进程崩溃的瞬间。正常流程中消费者处理完消息后主动 ack但如果消费者处理完还没来得及 ack 就宕机了hermes-agent 会判定消息未消费成功并重新投递。这时如果业务方没有做幂等就会出现重复处理。我在 3.2 节反复强调幂等就是因为这个场景太容易发生而且一旦发生排查成本极高。第三类最容易忽视本地缓冲文件损坏。本地磁盘缓冲文件在服务被强制 kill 时可能写了一半恢复后读取该文件会抛出异常如果不做容错这部分消息就永久丢失。我采用的解决办法是每条消息追加写入时带 CRC32 校验值读取时先校验不通过的单条消息丢弃并告警其余正常消息继续投递。这个改动成本极低却能把不可控的丢消息场景压缩到几乎为零。4.2 重试风暴一个配置引发的线上事故上线后第二周我收到告警下游订单服务接口超时率突然飙升。查了半天最后定位到原因是我们上线时把retry-initial-delay-ms配成了 200msretry-max-attempts配成了 20 次。当下游服务出现短暂抖动时hermes-agent 会在 200ms 内启动第一轮重试失败后继续指数退避重试。下游只是抖了几秒钟却承受了来自代理层的高频重试轰炸导致抖动被放大成不可用。这个事故的核心教训是代理层的重试策略应该以下游的承受能力为约束而不是以上游的焦急程度为约束。后续我把初试延迟调整为 2 秒并且在全链路配置了熔断开关——当连续失败率达到阈值时hermes-agent 会主动暂停向该目标投递新消息只保留已有消息的重试等下游恢复后再继续。这样既能保证消息最终可达又避免重试本身成为雪崩的导火索。4.3 路由漏配导致的黑洞消息还有一类问题特别隐蔽就是路由规则漏配。我在 2.1 节里面设计了默认路由作为兜底但早期版本里默认路由是没有的。有一次业务团队新加了一个user.level.up的逻辑主题但运维平台上的路由规则没有同步添加。结果就是生产者发送消息时路由层找不到匹配规则代理层默认把这个消息丢弃了。这个问题的尴尬之处在于发送接口返回的是成功因为消息成功进入了 hermes-agent但消息实际没有到达任何 MQ消费者自然永远收不到。业务方排查两天都没头绪最后查到是路由缺失。后续我做了两个改进一是路由缺失时send方法返回失败而不是成功让业务方第一时间感知到异常二是新增告警机制——当同一个逻辑主题在 5 分钟内出现 10 次以上的路由缺失时立即通知值班人员。5. 性能基准与调优经验什么样的配置适合什么场景5.1 一套压测基准数据代理层这类组件最怕的是看起来能用压一压就崩。我基于 4 核 8G 的普通云主机用 2000 条/秒到 20000 条/秒的递增压力测过一组数据可以作为参考起点发送速率条/秒平均投递时延msP99 时延ms内存占用MB表现2000412180稳定5000835265稳定100001588480稳定GC 频率升高2000046210780出现性能拐点时延抖动明显实测中最值得关注的是 10000 条/秒以上的场景。此时 CPU 占用率已经超过 70%GC 频率快速上升代理层成为瓶颈的概率显著增大。如果你的业务预估会长期超过这个水位优先考虑水平扩容多实例部署后按逻辑主题做一致性哈希分区把负载分散到多个 hermes-agent 节点上。5.2 两个提升吞吐的有效调整第一个调整是批量投递的阈值联动。我最初把batch-size固定为 200后来改为动态调整低峰期按条数触发高峰期按时间窗口触发达到 10ms 窗口即投递一次。这个改动的效果是在混合负载场景下平均时延下降了约 20%批次的填充率提升到 92% 以上没有因为凑不够批量而拖延消息。第二个调整是发送线程池的任务窃取机制。JDK 的ForkJoinPool天然支持工作窃取比固定线程池在多个主题队列并发消费的场景下更能榨干 CPU。默认情况下每个主题一个队列不同主题的消息量差距可能很大有的队列积压严重而有的线程空闲。换成工作窃取池后空闲线程会主动从繁忙队列中拿任务执行整体吞吐大约提升了 15%。这个优化代码改动不到十行收益却很直观。5.3 什么时候不该用 hermes-agent写了这么多它有多好用也得说清楚它的边界。hermes-agent 本质上是一个通用代理层它带来的统一、可控、可观测是以一层抽象开销为代价的。如果你的系统里只有一两个微服务消息量每天不到几万条那直接用现成的 RocketMQ/Kafka 客户端就够了不需要额外引入这个代理层——多一层就多一个故障点这是架构上永远要记住的守恒定律。另外如果你的团队已经有稳定且运转良好的消息治理体系没有多个消息组件并存的混乱局面也没有跨团队主题命名不统一的问题那 hermes-agent 对你来说就是过度设计。自己写一个代理层和引入一个开源消息网关类似解决的都是规模上来之后治理难的问题规模没上来就硬上只会增加维护成本。我在实际落地过程中最大的体会是像 hermes-agent 这样的信使模块它的价值不在于某个发送接口多好用而在于它能让你在出问题时快速回答三个问题这条消息从哪来、现在在哪、为什么还没到。只要把这三个问题变成几张表和几条告警代理层就算真正合格了。如果你正在被消息链路上的混乱折磨不妨照着这个思路自己搭一个先从最小的路由收口开始一步步把这层信使做扎实。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价