尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

轻量级消息代理 hermes-agent 实战指南:解耦、可靠投递与任务调度

轻量级消息代理 hermes-agent 实战指南:解耦、可靠投递与任务调度 做后端服务的同学大概率都有过这种经历业务模块之间的调用像一团乱麻A服务要同步等B服务返回结果B挂了A就跟着超时流量一上来数据库连接池先被打满定时任务散落在各个业务进程里没有统一的重试机制凌晨三点某个任务静默失败第二天数据对不上账才发现。我早期也被这类问题折磨过后来干脆自己动手维护了一个轻量级的消息代理组件也就是标题里提到的这个 hermes-agent。名字取自希腊神话里专门负责传信的赫尔墨斯定位很直白做一个默默跑在业务背后的“信使”把消息从生产者手里稳妥地送到消费者手里。这篇文章不是官方文档的复述而是我从实际项目里摸爬滚打出来的经验总结。内容会覆盖 hermes-agent 的核心模块拆解、部署接入的完整过程、那些文档里不会写的踩坑记录以及压测调优的思路。如果你正在为服务间同步耦合、任务投递不可靠、或者单纯想找一个比 Redis 列表更好用的消息中转方案而头疼这篇内容可以帮你少走不少弯路。1. 为什么我会选择维护一个 hermes-agent 这样的消息代理1.1 最初的问题同步调用把系统绑成了单点先说我最开始遇到的场景。当时手上有几个业务服务用户下单后需要同步调用库存服务、优惠券服务、积分服务。用户侧看着一个下单接口后端实际上是三个串行 HTTP 调用任何一个服务抖动下单就失败。这个架构的脆弱之处不用我多说典型的多米诺骨牌下游一个慢查询拖垮整条链路。这个痛点我估计很多人都有共鸣。同步调用的本质问题是“强耦合”——调用方必须等被调用方就绪双方必须同时在线流量高峰还必须让二者的容量精确匹配。可实际业务里库存服务和积分服务的处理速度根本不一样硬把它们绑在同一个事务里只会互相拖后腿。我当时需要的是一个能把“发消息”和“处理消息”拆开的中间层。生产者只管把“用户已下单”这个事件扔进去剩下的库存扣减、积分发放消费者自己慢慢消化。就算消费端某个服务临时挂掉消息还在队列里等着等它恢复之后继续处理。1.2 为什么不用现成的消息队列反而自己搞一个说到削峰填谷、异步解耦很多人第一反应是 Kafka、RabbitMQ 这类重型 MQ。这几个方案我都用过稳定是真稳定但重也是真重。Kafka 需要一个 Zookeeper 集群或者 KRaftRabbitMQ 要维护 Erlang 环境的版本兼容性如果只是几个内部服务之间传消息这些基础设施的运维成本比业务本身还高。相比之下hermes-agent 这种轻量代理的核心优势是“不挑环境、不抢资源”。它不需要独立的存储集群天然支持把消息转发到 HTTP 接口甚至可以直接推送到业务方的回调地址。这种形态特别适合内部服务数量在几个到几十个的中小型团队既想要消息队列的解耦能力又不想背上 Kafka 那套重运维包袱。我在选型时做了一个简单对比核心维度如下表对比项hermes-agentRedis StreamsRabbitMQKafka部署复杂度单个二进制无外部依赖依赖 Redis 版本依赖 Erlang 环境依赖 Kafka 协调组件消息投递方式HTTP推送/回调为主需自行实现推送逻辑AMQP 协议需要客户端 SDK 消费持久化能力支持本地积压存储依赖 Redis 持久化配置支持支持运维成本低中中高高适合场景中小规模服务解耦简单队列场景企业级复杂路由大数据量日志与流处理这个表不是否定 MQ 的价值而是想说选型得看场景。如果每秒消息量过万、需要复杂路由和回溯消费还是老实用 Kafka但如果是百来台机器规模、消息量每秒几百上千hermes-agent 这个体量的组件足以撑住。2. hermes-agent 的核心能力拆解从消息投递到任务生命周期管理2.1 连接管理与心跳保活先聊连接管理。消息代理的第一要务是保证生产端和消费端之间链路可用而链路可用靠的是心跳机制。hermes-agent 的处理方式比较直接生产端每间隔固定时间发送一个轻量心跳包代理在连续 N 个心跳周期没收到数据后判定对端失联触发连接回收和待发送消息的积压处理。这个机制跟我之前自己写 socket 长连接时的经验很一致核心是两个参数的配合心跳发送间隔和失联判定阈值。间隔太短会浪费带宽和 CPU太长则故障感知滞后。实际使用中业务允许的故障恢复时间大概是多少就把失联阈值设在多少以内心跳间隔取阈值的四分之一到三分之一比较合理。比如希望 30 秒内感知对端故障心跳间隔 10 秒一次连续 3 次未收到心跳就判定失联。2.2 消息路由基于 Topic 的订阅分发消息路由是消息代理的心脏。hermes-agent 采用了 Topic 加订阅分组的模型这个模型用一句话说就是生产者把消息发到某个 Topic代理根据订阅关系把消息副本推给每个订阅该 Topic 的消费组。这个设计模仿了消息队列里最成熟的那套思路。每个消费组内部消息按负载均衡策略分发给组内的多个消费者实例同一个消息在一个消费组内只会被一个实例处理。而不同的消费组之间则互为独立——比如订单服务创建一个“订单已支付”Topic库存服务和积分服务可以各自建一个消费组订阅它两边互不干扰地拿全量消息。我在实际配置时踩过一个认知误区把“消费组”和“消费者”混为一谈。一个消费组代表一个独立的业务处理方里面可以有一个或多个消费者实例如果你把同一个业务的两个实例放进了不同的消费组它们会各自消费一遍全量消息这在某些网关通知场景会造成重复推送。正确姿势是一个业务方对应一个消费组组内加实例才是扩容。2.3 消息生命周期与状态机机制消息的一生在 hermes-agent 里会经历这样一个状态流转新建、待投递、投递中、投递成功、投递失败、重试中、最终被放进死信队列。理解每个状态的含义对排查问题很有帮助。新建状态表示消息刚被代理接收尚未分配给任何消费组待投递状态表示消息已确认目标消费组等待被消费者拉取或等待代理主动推送投递中状态表示代理已通过网络把消息发给了消费者但还没收到确认回执收到消费者的成功应答后消息进入投递成功状态从待处理队列中移除反之进入投递失败状态根据配置决定是直接重试还是进入死信队列。这几个状态不是凭空设计的它们对应了消息系统最核心的“至少一次投递”语义。因为网络可能出问题消费者可能已处理但回执丢失所以消息被重复投递是不可避免的。能做的不是消灭重复投递而是通过状态机确保每条消息最多被“有效处理”一次——这部分我在后面踩坑那一章会展开聊。2.4 重试机制与死信队列重试机制是消息代理可靠性最关键的一环。hermes-agent 的重试采用了“指数退避加最大重试次数”的老牌策略第一次重试延迟 1 秒第二次延迟 2 秒然后 4 秒、8 秒、16 秒……直到达到配置的最大重试次数。超过最大次数后消息会被转入一个独立的死信 Topic不再骚扰消费者但会触发告警通知运维人员。指数退避的意义在于避免“重试风暴”。如果所有失败消息都立刻重试一旦消费端短暂故障代理就会在恢复瞬间收到海量重试请求直接把消费端再次打挂。这是踩坑踩出来的教训重试要像道歉一样第一次说对不起第二次隔一会儿再说表达诚意但不能把人烦死。死信队列并非处理完的终点。我通常会在死信 Topic 上挂一个专门的补偿消费者解析失败消息结合业务日志决定是修正数据后手动重投还是放弃。死信不能只存不处理不然等于把问题从“消息投递失败”拖延成“消息被静默丢弃”。3. 从零接入的完整实操让第一条消息从生产端到达消费端3.1 环境准备与部署形态部署 hermes-agent 的过程简单得有点不像一个“中间件”。只需要把编译好的二进制丢到服务器上写一个 YAML 配置文件然后启动进程就行。没有依赖数据库没有依赖单独的缓存服务消息的积压存储基于本地磁盘文件这也让整体部署非常轻量。我在测试环境用的是单节点部署生产环境采用两节点互为主备的形态。代理本身是无状态设计的消息落盘后即便节点重启也能从本地文件恢复未投递完的消息。这里注意一点消息是本地存储节点挂了但没有将文件转移到备节点的话备节点无法接续处理存量消息。所以更稳妥的做法是让生产端双写或者依赖上层业务的重试补偿。3.2 配置文件的语义详解以 YAML 配置为例核心配置项如下server: listen: 0.0.0.0:7788 # 代理对外服务端口生产端和消费端都连这里 max_connections: 2000 # 最大连接数超过后拒绝新连接保护代理自身 storage: data_dir: /var/lib/hermes-agent/data # 消息落盘目录注意预留磁盘空间 segment_size: 256MB # 每个数据分片的大小按消息量调整 sync_mode: async # async / fsyncasync 性能好fsync 更安全 route: topics: - name: order.created partitions: 8 # 分区数决定消息并行处理能力 consumers_mode: push # push / pull推送模式或拉取模式 delivery: ack_timeout: 30 # 消费者确认超时秒数超时判定失败 max_retries: 5 # 最大重试次数 retry_base_delay: 2 # 重试基础延迟秒数 dead_letter_topic: dlq.order.created # 死信 Topic几个容易理解错的地方说一下。partitions参数决定了一个 Topic 在代理内部被分成几个并行处理的分片分区数越多并行度越高但也会带来更大的文件句柄消耗和消息顺序性变化。sync_mode选 async 时消息先写内存缓冲再异步刷盘吞吐高但极端情况下可能丢刚写入的数据选 fsync每条消息都强制落盘安全但对磁盘性能要求高。我的建议是先从 async 开始跑压测之后如果吞吐达标且能接受极端场景下的少量消息丢失保持现状即可如果消息很重要再切换 fsync。3.3 生产端接入如何把一条消息安全投递出去生产端接入的代码模式非常干净。以 Python 为例核心就三个动作建立连接、发送消息、关闭连接。from hermes_agent import Producer producer Producer( endpointhermes://10.10.10.10:7788, topicorder.created, timeout5, ) def send_order_created(order_id: str, user_id: str, amount: int) - bool: payload { order_id: order_id, user_id: user_id, amount: amount, created_at: 1710000000, } try: producer.send(payload, keyorder_id) return True except Exception as exc: # 这里注意send 抛出异常并不代表消息一定没被代理接收 # 可能是代理已接收但回执在网络上丢了需要去重机制兜底 print(fsend failed: {exc}) return Falsekey参数值得单独说。它决定了消息在分区间的分配策略同一个 key 的消息会始终进入同一个分区。对订单类消息拿order_id当 key 可以保证同一个订单的所有事件天然有序如果对顺序性不敏感不传 key 也行代理会按轮询策略把消息打散到各分区吞吐更均衡。生产端有一个细节必须注意send 方法返回异常并不等于消息一定没被收到。字节流可能已经到了代理只是确认包在回程中丢失。因此业务侧在收到 send 异常时不能简单判断“发送失败”而要在消费端配合幂等去重处理否则重发会带来重复数据。这一点我在第 4 章展开讲。3.4 消费端接入回调函数与确认机制消费端接入更简单注册一个回调函数在收到消息时执行执行成功就返回确认代理收到确认后删除消息执行异常则触发重试。from hermes_agent import Consumer def handle_order_created(payload: dict): order_id payload[order_id] user_id payload[user_id] amount payload[amount] # 模拟业务处理扣减库存、发放积分、发送通知 print(fprocess order: {order_id}, user: {user_id}, amount: {amount}) # 处理成功后无需额外操作函数正常返回即视为确认 # 如果处理失败抛出异常代理会自动重试该消息 consumer Consumer( endpointhermes://10.10.10.10:7788, topicorder.created, groupinventory-group, callbackhandle_order_created, ) consumer.start()这套模型本质上是把消息消费做成一个“可重试的事务”。回调函数正常返回代表处理成功抛出异常代表处理失败这条消息会被代理重新投递。因此回调函数里不能有不可逆的操作。比如扣库存和发短信这种动作要保证即便重复执行也不会造成资损。一个稳妥的措施是引入去重表每条消息在第一次进入回调时把message_id写入一张去重表重复执行时先查表命中则直接返回成功。这样处理逻辑天然幂等无论代理重试多少遍业务数据都不会脏。这个处理方式虽然会给每笔请求增加一次存储查询但相比数据错乱带来的对账成本这点开销微不足道。4. 跑通之后才是真正开始线上环境踩过的坑菜单搭起来、代码写顺了不等于系统可以高枕无忧。这一章写的都是我实际运维时遇到并解决过的坑每个都花了不少时间才定位到根因。4.1 第一个坑消费端重复执行原以为是代理重复投递结果是自己业务没做幂等上线第二天运营反馈说某些用户的积分好像加了两次。我第一反应是查代理的重试日志结果发现消息只投递了一次消费者也成功确认了。随后翻数据库才发现问题出在消费回调里扣减积分的逻辑是先查询用户当前积分再加上本次积分然后写回库。而当两条不同的消息比如下单积分和签到积分几乎同时到达两个回调各自读到旧积分余额分别加完后写回后写的覆盖了先写的相当于丢了一次加分。这不是重复投递这叫并发覆盖写。排查链路是这样的先看消息日志确认消息被消费的时间点和数据库写入时间点是匹配的再看应用日志发现两条消息的处理日志几乎在同一秒打印最后看了数据库记录确认是覆盖写而不是追加写。到这里定位就清楚了——这不是消息代理的问题是消费端代码没有考虑并发串行化。解决方案是在数据库层面用原子更新把“读改写”三个动作合并为一个 UPDATE 语句保证同一用户的积分变更按顺序执行。如果用 Redis 做积分存储就用 Lua 脚本把读和写做成原子操作。这个教训让我养成了一个习惯所有消费回调里的写操作都得先问自己“这条 SQL 在并发下执行两次数据会不会错”。4.2 第二个坑连接风暴代理重启把消费端全部打挂有一次给代理节点做例行重启重启过程很顺利但随后监控显示消费端服务 CPU 一路飙升最后整批宕机。排查下来发现是这么一回事代理重启后所有消费端的长连接同时断开各消费端为了尽快恢复连接在极短时间内同时发起重连请求。服务器的进程数、文件描述符、握手线程全部被打满代理反而因为过载拒绝了新的连接形成一个恶性循环。这个坑在连接型中间件里非常经典英文叫 connection storm中文互联网一般叫“连接风暴”。根因是消费端重连时机过于集中没有加入抖动延迟。修复方案也很经典消费端在断线重连时不能立刻重连而要加上一个随机延迟延迟范围取 1 到 10 秒。这相当于让所有重连请求在时间轴上散开避免对齐式的集中冲击。同时我在代理侧也做了保护限制单 IP 的连接频率以及设置最大半开连接数。这类保护看着多余但在线上出问题的时候能救命尤其是当你通过域名接入多个消费端时DNS 刷新可能导致所有消费端同时换 IP 重连。4.3 第三个坑消息积压消费者扩容却没用问题出在分区数某个大促活动开始时订单量瞬间翻了好几倍代理端消息积压很快到了几十万条。我第一反应是给消费端扩容机器从 2 台加到 8 台结果积压数纹丝不动。后来仔细看文档才意识到模型里有个关键约束同一个消费组内一条消息在同一时间只会被一个消费者实例处理而消息是按分区分配给消费者的消费者数量超过分区数之后多出来的消费者实例只能闲着。换句话说分区数决定了消费并行度的天花板。Topic 创建时我设了 4 个分区哪怕后面往消费组里塞 100 个消费者实例同时最多也只有 4 个实例在拉消息。这便是积压无法缓解的真正原因。排查思路梳理如下表排查步骤操作结论1. 检查消费组内实例数8 个消费者都在线实例数不是瓶颈2. 检查代理分发状态只有 4 个消费者分到消息消息只按分区数并行分发3. 检查分区配置Topic 只有 4 个分区并行上限为 4扩容无效4. 修改分区数并滚动重启分区提升到 12 个积压开始快速下降这个坑暴露了一个容量规划原则Topic 的分区数在创建初期就要按业务峰值和消费可并行的最大规模来定而不是按日常流量定。分区数过小后续扩容消费端没用要修改分区数还需要新建 Topic 或做数据迁移分区数过大则会产生过多小文件影响代理本身的吞吐稳定性。实践中保守一点的做法是按当前日常流量的 4 倍峰值来预估所需分区再乘 1.5 的冗余系数。4.4 第四个坑ack 超时导致频繁重试消费者其实只是有点慢默认的ack_timeout是 30 秒平时运行没有任何问题。某天一个消费回调里多了一个外部接口调用响应偶尔会飘到 40 秒开外。代理端因为 30 秒没收到确认判定这条消息投递失败于是触发重试重试几秒后外部接口恢复消息成功被消费并确认。但原来的那条还没超时的消息继续处理最后也确认了。结果一条业务数据被执行了两次。这个问题的根子在于代理不知道消费者需要多久才能处理完超时时间只是一个拍脑袋的配置。处理方案很直观让代理的 ack 超时时间大于消费者的 p99 处理耗时即 99% 请求都能在该时间内完成处理而不是等于平均值。通过监控得出消费者的 p99 处理耗时是 50 秒我就把 ack_timeout 改成了 90 秒多留一点余量重试次数立刻降了下来。另外一个相关的点是客户端的心跳间隔。如果消费者长时间不发送心跳代理会误判消费者失联并释放连接导致消息被分配给其他实例。所以心跳间隔也要和 ack_timeout 配合设置我的习惯是把心跳间隔设为 ack_timeout 的一半比如 ack_timeout 90 秒心跳间隔 45 秒一次。5. 性能调优与容量规划用数据说话而不是拍脑袋5.1 吞吐量估算从业务目标反推技术参数接入阶段很多人的第一个问题是“我的业务需要多少个分区、多少台代理”。这个问题不该靠猜而应该从业务指标倒推。假设你的业务高峰期每秒产生 2000 条订单消息每条消息大小约 2 KB你期望消费端能在 10 秒内把消息全部消化掉且单条消息处理耗时约 50 毫秒那么每个消费实例的吞吐大约是 1 / 0.05 20 条/秒。消化 2000 条/秒的流量需要 2000 / 20 100 个消费实例并发而这又要求分区数至少等于 100。这个简单的算术题说明了一个直接结论分区数的上限不是消息量而是消费端的处理速度。如果你的消费端单条处理需要 500 毫秒那 1000 条/秒的流量就需要 500 个实例这显然不现实。更合理的方案不是无限扩容而是优化消费端的处理逻辑——比如把 500 毫秒的外部调用改为异步批量处理单条耗时降到 50 毫秒后实例数和分区数都降了一个数量级。5.2 磁盘和内存的平衡消息量大的时候磁盘 IO 和内存分配往往是瓶颈。本地存储模式下消息写入磁盘文件读取后删除。如果写入频率高、单条消息小会产生大量小文件碎片文件句柄数量上升操作系统层面的 IO 压力随之增大。我的调优动作主要有三个一是把segment_size调大减少分片合并的频率二是将消息批量刷盘一次写入合并多批数据减少磁盘寻道次数三是为数据目录单独挂一块 SSD同时设置独立的 IO 调度策略避免和其他业务的磁盘读写互相干扰。内存方面我建议给代理进程留足 page cache 空间。消息先写内存再刷盘系统会利用空闲内存缓存热点分片文件。如果服务器内存有 64 GB使用率不高的情况下不需要刻意限制代理的内存占用让它多用 page cache 反而是好事情。真正需要担心的是内存被打满导致 swap所以还是要给操作系统保留一定空闲内存。5.3 监控指标这几项不盯出事必手忙脚乱代理是否健康不能靠等 bug 出来再看我的监控面板里固定放着四个指标。第一是消息积压数量这个直接反映消费是否跟得上生产第二是重试消息占比这个指标如果突然上升说明消费端在报错或者 ack 超时第三是死信队列增长速率正常应该趋近于零一旦出现有规律的死信增长就是有某类固定消息被反复处理失败第四是代理进程的句柄数与连接数这个可以在连接风暴和设备文件耗尽之前提前预警。积压数量这个指标要额外加一个按 Topic 分组的统计因为全局积压看起来正常但不代表单个 Topic 没有异常。我曾经遇到过全局积压只有几百条实际上某一个订单 Topic 积压了 20 万条的情况因为其他 Topic 消费很快平均下来总量并不大结果掩盖了问题。6. 从消息代理走向任务调度hermes-agent 的两种进阶用法6.1 定时任务与延迟消息的落地实践消息代理用顺手之后我开始尝试扩展它的能力边界。第一个用法是把延迟消息做成定时任务调度器。场景是这样的下单后 15 分钟未支付需要自动关单用户申请退款后 48 小时未处理需要提醒客服。原生实现往往要起一个单独的任务扫描表定时轮询数据库效率低且存在扫描倾斜问题。把 hermes-agent 和延迟队列结合后整个流程变成一条流水线业务服务创建订单时投放一条延迟消息到代理代理在指定的延迟时间之后才将消息推给消费端消费端收到消息后再去数据库里核实订单的支付状态若未支付则执行关单若已支付则直接放弃处理。这个方案省掉了定时扫描任务也不用维护额外的任务表。这里有个很重要的细节延迟类的回调处理必须是“复核式”的不能收到延迟消息就无条件执行动作。因为延迟消息可能被重复投递而且从投放消息到触发回调的这段时间里原始业务状态可能已经变了。收到关单消息时先查一遍订单当前状态已经支付就什么都不做这样整个机制才是幂等可靠的。6.2 用消息代理做服务状态感知与自动补偿第二个进阶用法是状态感知与补偿说白了就是让消息代理充当业务的“神经末梢”。某个下游服务被调用失败时生产端不直接把错误吞掉而是投递一条失败事件到代理专门的补偿消费者会分析这类失败事件结合失败频率和业务规则决定是否启动重试、告警或者切换备用通道。本质上是把业务里的异常分支也变成一条消息流。以前异常处理散落在各个服务里有的记日志、有的打告警、有的重试三次就放弃行为完全不统一。接了消息代理之后我把所有业务的失败事件统一投递到同一个失败处理 Topic由同一个消费组专门负责分类处理。这套模式减少了每个服务的异常处理代码量也让系统里“哪里有坑”变得一目了然——失败事件累积最多的那个 Topic就是最需要重构的业务代码所在。6.3 消费能力扩展从内部解耦到对外开放回调最后一个扩展方向是对外开放回调能力。hermes-agent 本身支持 HTTP 推送消费者的形态所以可以不写服务端代码直接让代理把消息以 HTTP POST 的形式推给一个公网可访问的回调地址。这个特性在对接内部运维系统、对接企业微信机器人、或者做外部 Webhook 通知时特别方便不需要为了发一条通知而专门起一个 WEB 服务。对这类对外回调要注意两个安全问题一是回调地址必须配置白名单防止代理被恶意用来做内网请求转发二是回调地址必须是 HTTPS避免消息内容在传输过程中被截获。内部调试阶段用 HTTP 没太大问题一旦涉及真实用户数据证书配置这件事不能省。7. 写在最后消息代理不是一个终点而是一个起点从最初为了解决服务间同步耦合开始到后来逐渐承担了延迟任务、失败补偿、状态感知、外部回调这些职责hermes-agent 这个组件在我的技术体系里已经从“一个中间件”变成了“一套业务韧性基础设施”。回看整个过程几个教训我印象最深。分区数一定要按峰值和消费并行度来规划而不是按平均流量估否则线上扩容会撞上一个看不见的天花板。消费回调必须做幂等处理重复投递不是 bug不做去重才是。配置 ack 超时和重试策略之前先去监控面板看一眼消费端的 p99 耗时所有参数调整都应基于真实数据而不是猜。如果只让我留一条建议那就是消息代理类组件的核心价值不在“能把消息发出去”而在“发不出去、收不到确认、重复投递、消费积压这些异常场景下系统能不能稳定自愈”。把异常路径设计好代理才能真正成为那个名字里所说的信使——平时无感关键时刻绝不掉链子。
返回列表