ARTICLE DETAIL

资讯详情

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

hermes-agent实战:轻量消息代理与异步任务调度系统拆解

hermes-agent实战:轻量消息代理与异步任务调度系统拆解 先交代一下背景hermes-agent这个名字玩技术的朋友一眼就能看出门道。Hermes 在希腊神话里是传递消息的信使神所以叫这个名字的项目基本都围绕“消息流转、任务分发、调度执行”在打转。再加上“agent”这个后缀说明它不是单纯的消息队列而是一套常驻后台、主动干活、按规则办事的代理服务。我最近把一个内部小项目用 hermes-agent 重构了一遍从原来散落各处的定时脚本加手工触发收敛成统一的消息驱动体系。这篇文章就把整个拆解和落地过程写出来包括为什么选它、哪些环节最容易踩坑、以及我是怎么一步步调通的。如果你手上也有一堆“定时任务回调通知多系统联动”之类的脏活或者团队里消息通路乱成一团、每个服务各写各的轮询那么这篇内容应该对你有直接帮助。我不打算写那种贴满源码的文档式教程更多是讲清楚设计思路、关键取舍以及实际跑起来之后才会遇到的真实问题。1. 内容整体设计与思路拆解1.1 为什么需要 hermes-agent 这类中间代理层先说一个很常见的场景你手上有三个内部系统A 系统产生订单数据B 系统负责风控审核C 系统做通知触达。最开始大家图省事直接让 A 系统在业务代码里调用 B、C 的接口。跑了一阵子就发现每次 B、C 升级接口或者临时下线维护A 系统就要跟着改代码甚至因为 C 系统超时把 A 系统的核心下单流程也拖慢了。这类问题的本质是“生产者”和“消费者”之间耦合太深。而 hermes-agent 这类代理层就是在中间插入一个独立服务生产者只把消息丢给它它负责把消息可靠地送到消费者手里。这样 A 系统不需要知道 B、C 现在是什么状态B、C 也可以随时扩容缩容不用跟 A 打招呼。我当时做方案选型的时候对比过几个方向。一是直接用 Redis 的列表结构自己写一个消息队列轻量但功能太原始没有重试、没有死信、没有消费者组管理到后期还是得自己造一堆轮子。二是部署一套完整消息中间件能力是强但运维成本高几十人的小团队有点扛不住而且项目初期并发量根本达不到需要上那套东西的程度。hermes-agent 刚好卡在中间它有一定程度的消息持久化和重试机制又有 Agent 端的调度能力还能挂回调通知属于“够用但不臃肿”那一档。注意我这里说的 hermes-agent 并不是某个特定商业产品或固定开源仓库的名字它更像一类“信使代理”模式的总称。很多团队内部都会自研类似组件也有大量开源实现基于这个思路。文章里的配置和目录结构是以一种比较通用的设计为例方便你迁移到自己的具体实现上。1.2 核心角色划分Producer、Agent、Dispatcher、Consumer一套完整的 hermes-agent 体系里通常有四个角色这个划分对我理解整个系统帮助很大也建议你接手任何类似项目的时候先把角色画清楚再动手。Producer生产者负责发送消息的一方。在重构前A 系统的下单模块就是 Producer订单创建成功后直接把业务数据打成消息发给 hermes-agent。生产者的核心任务是“发出即忘”它不关心消息最终被谁消费了、消费成功没有只负责保证自己把消息成功交给了代理层。Agent代理节点这是 hermes-agent 体系里最能体现“agent”含义的部分。它是一组常驻后台的进程负责接收 Producer 的消息、把消息写入存储、然后按策略把消息分发给下游。如果叫它“代理”不如叫它“管家”更贴切——它知道每条消息该往哪儿走也知道消息万一失败该找谁。Dispatcher调度分发器负责决定“消息现在该不该发、发给谁”。这是策略集中的地方比如最简单的直接投递按照消息里的 target 字段路由到指定消费接口再比如延迟消息先放到等待队列等到时间了再转给投递器还有限流如果下游接口每秒只能扛 50 个请求Dispatcher 就要做速率控制。Consumer消费者消息的最终处理方。对我们内部系统来说B 系统的风控审核接口和 C 系统的通知推送服务都是 Consumer。消费者的核心要求是接口幂等因为 hermes-agent 为了保证消息不丢一定会重试重试就会导致重复投递。消费端不支持幂等的话下游订单可能被处理两次、用户可能收到两条重复短信这就要出大事故。我用一个生活化的类比来帮团队新人理解Producer 就像你在电商平台下单付款你只需要知道付款成功即可后续商家是不是立刻看见了、仓库几点发货你是不会一直盯着看的。hermes-agent 是平台中间的系统它记录你的订单信息、判断要给哪个商家推送、如果商家一直不处理还要反复提醒。Consumer 是商家后台商家收到订单后要保证自己处理订单的动作不会因为系统卡顿而重复出库。1.3 状态机与消息生命周期设计 hermes-agent 的时候最核心的抽象是消息的状态机。我一开始没设计好导致后面各种奇奇怪怪的问题。后来收敛成六个状态基本覆盖了所有场景PENDING待处理Producer 刚把消息发给 AgentAgent 已经确认接收并写入存储但还没交给 Dispatcher。DISPATCHING分发中Dispatcher 正在把消息按路由策略送给对应 Consumer。RETRY_WAIT等待重试上一次投递失败了消息进入冷却期等待重试时间到。SUCCESS成功Consumer 明确返回了成功消息生命周期正常结束。DEAD_LETTER死信重试次数用尽仍然失败消息进入死信队列等待人工干预。CANCELLED已取消业务主动取消或者消息超过最大存活时间被系统作废。为什么要这么细致地设计状态因为代理层的核心责任就是“兜底”。Producer 发消息是异步的Consumer 处理也是异步的中间只要一个环节断掉消息就会被卡住。没有状态机的话出了故障你根本不知道消息是丢在发送路上、存在队列里还是送到消费者了。有了这些状态每一条消息都像快递物流一样有迹可循出问题可以按状态去检索。状态机代码里一般用一个整数字段0 到 5 对应上述状态。我当时加了一个 updated_at 时间戳为了做超时扫描——比如一条消息在一个状态停留太久说明它大概率卡住了。这个设计后来帮我们抓到了好几个因为下游服务假死导致的消息堆积问题。2. 核心细节解析与实操要点2.1 Agent 双线程模型的底层逻辑hermes-agent 能保证消息“至少一次”投递核心在于内部采用了一套双线程模型。这个模型值得你花点时间吃透因为它直接决定了整个 agent 的性能边界和故障处理方式。我把这条主流程称为“主线程 重试线程”的协作模式。主线程只干一件事把 Producer 发来的新消息接收下来写入本地存储并立即返回“接收成功”给 Producer。注意这里主线程绝不去调下游 Consumer 的接口因为下游接口一旦慢主线程就堵住了新的消息也收不进来。这就像餐厅前台只负责接单和记录绝对不会自己跑到后厨炒菜一样。那真正送消息的活谁干呢由独立的工作线程池去轮询待分发消息。工作线程休眠一个很短的时间片比如 200 毫秒醒过来就去存储里扫一批处于 PENDING 状态的消息逐个按照路由规则发出。如果发出成功就标记 SUCCESS如果失败就切入重试流程由重试线程负责后续的定时补偿。重试线程是整个模型里最容易出问题的部分因为你要控制“重试频率”和“系统压力”之间的平衡。如果失败后立刻重试下游接口本来就没恢复只会被反复打死如果重试间隔太长又会影响业务时效。我最终采用的是一个指数退避策略公式很简单delay (2 ^ retry_count) * base基础值 base 设了 10 秒。第一次重试等 10 秒第二次 20 秒第三次 40 秒依此类推。这样下游短时抖动能快速恢复长时间宕机也不会被高频率请求淹没。具体实现时一定要把“主线程接收消息”和“工作线程分发消息”的线程池分离开大小独立配置。接收线程的 Pool Size 可以设成 CPU 核数的一半分发线程则可以大一些因为分发过程涉及网络阻塞线程多了才能同时等多个下游响应。2.2 消息可靠存储先写库再分发提到状态持久化最容易踩的坑是用内存队列硬顶。项目初期消息量小的时候内存队列确实又快又方便但一旦 agent 进程重启所有内存里没发出去的消息直接蒸发。生产环境这是不可接受的。我当时在方案里明确了一条铁律所有消息必须先落存储再由分发线程取出投递。我选择的存储方案是 SQLite 加本地文件双写。为什么不是 MySQL因为 hermes-agent 部署时追求轻量独立每个节点自己管理消息状态如果依赖外部数据库就多了一个故障点。而 SQLite 作为嵌入式数据库几乎没有运维成本一条记录对应一条消息查询状态也很方便。至于本地文件我主要是做了 WAL预写日志防止 SQLite 本身在突然断电时出现文件损坏导致消息丢。具体消息表结构我设计成下面这样核心字段都有索引CREATE TABLE messages ( msg_id TEXT PRIMARY KEY, trace_id TEXT NOT NULL, producer TEXT NOT NULL, route_key TEXT NOT NULL, payload TEXT NOT NULL, status INTEGER NOT NULL DEFAULT 0, retry_count INTEGER NOT NULL DEFAULT 0, next_retry_time INTEGER NOT NULL DEFAULT 0, created_at INTEGER NOT NULL, updated_at INTEGER NOT NULL ); CREATE INDEX idx_status ON messages(status, next_retry_time); CREATE INDEX idx_route ON messages(route_key);这里有个细节值得注意为什么不直接加一列叫create_time用 DATETIME 类型而要用INTEGER存毫秒时间戳因为 agent 内部大量逻辑要做时间比较比如“扫描所有next_retry_time小于当前时间的消息”用整数比较比字符串日期快一个数量级索引命中也更高效。这个习惯在我后来做数据处理的项目里一直延续着算是一个通用的优化思路。“先写库再分发”还有一个非常重要的附加效果它让 agent 天然具备了“断点续传”能力。即便 agent 进程突然被 kill -9 杀掉重启后只需要扫一遍数据库里所有非终态非 SUCCESS 和非 DEAD_LETTER的消息重新调度它们的分发工作即可。这本质上是一个可靠的 oversized envelope 模式消息内容永驻存储分发过程只传递引用。2.3 路由规则与通配符匹配的设计路由是 hermes-agent 里最灵活也最容易失控的部分。路由规则决定了“这条消息该去哪几个消费者”。一开始我打算做完全配置化像 API 网关那样每条消息都指定目标 URL但这很快暴露了问题下游服务调整接口地址的时候Producer 端也要跟着改配置又变成隐式耦合了。所以我最终还是落回了route_key 加通配符的设计。Producer 发送消息时不直接指定 URL而是指定一个抽象的业务标识比如order.created、risk.audit.passed、notify.sms.send。Agent 这边的路由表再把 route_key 映射成具体的下游地址而且支持通配符比如routes: - pattern: order.* targets: - service: risk-audit url: http://risk-service.internal/audit weight: 100 - pattern: notify.* targets: - service: sms-sender url: http://notify-service.internal/send-sms weight: 100这套设计的精妙之处在于Producer 端稳定只认业务语义。下游系统的变化全部收敛在 agent 的路由表里。而且通配符让同一个消息可以同时路由给多个消费者例如order.*同时匹配order.created、order.updated、order.cancelled下游如果只需要处理订单创建订阅order.created这个更精确的 key 就行了。实践下来我发现通配符规则还有个附加优势可以轻松实现灰度发布。比如新上线的风控 V2 接口想先接 10% 的流量就在路由表里给risk.audit.passed配置两个 target一个权重 90 指向老接口一个权重 10 指向新接口。这比在业务代码里写if (rate 0.1)这种硬编码要干净得多。2.4 回调通知机制Producer 如何拿到处理结果消息代理只负责把消息送出去是不够的生产者很多时候还关心消息下游处理成功没有。比如用户提交了一个申诉工单管理员审核通过后系统要发邮件通知用户。如果审核这个动作是异步消息就得想办法把“审核成功”这个信息传递回最开始提交工单的上游业务系统。hermes-agent 的回调机制设计上分两种同步回调和异步回调。同步回调是指 Producer 在发送消息时传入一个 callback_urlagent 在下游消费完成后往这个 URL 发起一次 HTTP POST携带消息 ID 和处理结果。异步回调则是 producer 自己订阅结果主题agent 把处理结果作为新消息反写回某个 topic 里生产者异步消费即可。我实际项目里两种都用了。面向接口的调用也就是请求响应模型几乎都是同步回调面向流程引擎的事件驱动模型比如订单状态流转把每一步都广播出去就用异步回调。这里要提醒一个很容易犯的错误回调 URL 一定是 agent 主动去调用 Producer 提供的接口而不是 agent 告诉 Producer “你等一会自己来查”。因为后者很容易做成轮询又退化回高耦合的老路了。回调数据里最重要的不是结果本身而是 trace_id。整个链路里从 Producer 发消息开始就生成一个唯一的 trace_id它跟着消息进入 agent跟着投递请求打到 Consumer再跟着回调返回 Producer。有了 trace_id排查请求链路就等于有了主线之前那种“消费者报错但不知道对应哪条业务记录”的窘境就消失了。3. 实操过程与核心环节实现3.1 自建 hermes-agent 最小初始化目录结构这里我按自建一个最小可用的 hermes-agent 项目来讲。下面这个目录结构是我实际动手搭建时用的不依赖特定框架清晰且适合二次扩展hermes-agent/ ├── bin/ │ ├── agent-server # 主服务入口负责消息接收与分发 │ └── agent-ctl # 命令行管理工具含手动重发、死信查询等 ├── config/ │ └── config.yaml # 基础配置包括监听端口、路由表、线程池大小 ├── internal/ │ ├── store/ # 存储层封装 SQLite 的读写接口 │ ├── dispatcher/ # 分发器核心轮询调度逻辑 │ ├── worker/ # 工作线程池与重试线程池 │ └── callback/ # 回调通知逻辑 └── data/ └── messages.db # SQLite 数据文件运行时自动生成可以看到这套结构把一个 agent 最核心的几件事拆得清清楚楚。store 只负责数据的持久化和状态更新dispatcher 决定什么时候该分发哪些消息worker 只管真正发 HTTP 请求callback 负责把结果通知出去。每个模块之间的依赖方向是单向的调试和单元测试都很方便。配置方面我用 YAML 格式。为了好上手我把配置项的注释直接写进文件里团队新人看到就知道每个参数是干什么的。示例配置server: listen_addr: 0.0.0.0:9090 receive_timeout_sec: 10 storage: dsn: ./data/messages.db wal_enabled: true worker: receive_pool_size: 4 dispatch_pool_size: 16 poll_interval_ms: 200 dispatcher: max_retry_count: 5 base_retry_delay_sec: 10 dead_letter_threshold: 20 routes: - pattern: order.* targets: - service: order-sync url: http://internal-order-sync:8080/api/sync - pattern: notify.* targets: - service: notify-center url: http://internal-notify:8081/api/send3.2 启动服务并手动模拟一条消息的全流程服务代码收工后我通常会把启动和验证的完整链路写成命令。先初始化存储目录并启动 agentcd hermes-agent mkdir -p data ./bin/agent-server --config config/config.yaml启动日志里看到“listening on 0.0.0.0:9090”就说明 Agent 已经进入工作状态。这时候我用另一个终端窗口手动构造一条测试消息投递进去curl -X POST http://127.0.0.1:9090/message \ -H Content-Type: application/json \ -d { trace_id: trace-20240101-001, producer: order-service, route_key: order.created, payload: {\order_id\: \A100001\, \amount\: 2999}, callback_url: http://order-service/internal/callback }agent 返回202 Accepted后我可以立即去 SQLite 库里看看消息是不是已经落库。这一步很关键它验证的是“接收与持久化”链路是否正常SELECT msg_id, route_key, status, retry_count FROM messages;正常情况下能看到一条 status 为 0PENDING的记录。等待 poll_interval_ms 配置的时间后分发线程会把它取走并通过路由表找到order-sync服务完成 HTTP POST 投递。此时状态更新为 SUCCESS如果下游返回 2xx。我再查一次库确认状态流转SELECT status, retry_count FROM messages WHERE route_key order.created;这一步走通说明 agent 的最小可用链路已经建立。后面的开发重点就是充实现有模块比如加上前端可视化界面、遥测指标导出等。3.3 手动触发一次死信流程并找回消息死信流程是我在测试阶段反复验证的一个环节它保证消息不会无限重试导致系统持续产生无效流量。我设计的是当重试次数超过max_retry_countagent 会把消息标记为 DEAD_LETTER 并存进一张独立表。下面是一个死信流程的现场记录。首先我把下游order-sync服务临时停掉然后发送 5 次相同的测试消息。agent 每次投递都会失败重试次数逐渐递增。因为 base_retry_delay_sec 设成 10 秒实际等待时间会很可观。为了快点看到效果调试时我会把基础重试延迟临时调成 1 秒测完再改回去。结果如下表投递次数重试次数状态变化触发逻辑第 1 次0PENDING - 失败下游连接超时第 2 次1RETRY_WAIT10 秒后再次尝试第 3 次2RETRY_WAIT20 秒后再次尝试第 4 次3RETRY_WAIT40 秒后再次尝试第 5 次4RETRY_WAIT80 秒后最后一次尝试第 6 次5DEAD_LETTER触发最大重试次数进入死信表查死信SELECT * FROM dead_letter_messages ORDER BY created_at DESC LIMIT 5;看到死信记录后我快速恢复下游服务再用管理工具把死信重新投递./bin/agent-ctl --resend --msg-id 20240101-0001这条命令会让 agent 把死信记录重新赋回 PENDING 状态分发线程下一轮就会正常投递。投递成功后死信记录自动从死信表移除同时保留一份归档副本方便审计。3.4 5 分钟搭建一个带 Hermes 风格通知的告警实战项目上线后我给它配了一个运维告警场景用来验证 realtime 通知能力。这个场景是“服务存活检测 消息通知 工单补齐”其实就是把 hermes-agent 当作一个事件中枢。我在 Agent 上注册两条路由一是monitor.host.down发给告警中心二是monitor.host.down同时发给工单系统。因为该场景希望两个下游都收到消息所以路由规则里直接配了多 target。测试时我用一个为监控脚本单独写的 Producer 工具模拟主机宕机事件curl -X POST http://hermes-agent:9090/message \ -H Content-Type: application/json \ -d { trace_id: trace-monitor-1001, producer: node-monitor, route_key: monitor.host.down, payload: {\host\:\web-01\,\ip\:\192.168.1.20\,\message\:\ping timeout 30s\} }消息发出后约几百毫秒告警中心就收到了推送工单系统也自动创建了一个 P1 优先级事件单。全程没有人工干预没有硬编码的 HTTP 调用。后来同事开玩笑说这是“Hermes 信使跑腿两边都不用等”。这个例子虽然小但它完整呈现了 hermes-agent 的定位从事件生产者到多个事件消费者之间搭建了一条可靠、可重试、可追溯的消息通路。4. 常见问题与排查技巧实录4.1 问题速查表长期维护 hermes-agent 的过程里我遇到过的典型问题可以说五花八门这里挑出最常见的一批整理成速查表方便大家直接对照。现象可能原因排查命令 / 方法解决方案Producer 发消息返回超时Agent 接收线程阻塞或 DB 写入慢查看 agent 日志中“receive”耗时调大 receive_pool_size优化 DB 写入参数消息一直停在 PENDING 状态分发线程未触发或路由未命中检查 dispather 日志和路由表确认 poll_interval_ms 已启动pattern 匹配正确消息状态频繁进入 RETRY_WAIT 但业务正常Consumer 响应码不符合 2xx 预期查看 Consumer 完整响应体调整 agent 的健康判断规则允许 3xx 重定向延迟过高消息从发出到消费超过 5 秒分发线程打满下游响应慢pstack或线程池监控增加 dispatch_pool_size对下游做超时熔断死信表堆积大量记录下游接口代码 bug 或参数格式变更抽样 dead letter payload人工恢复后使用 agent-ctl 批量重发重启后部分消息丢失上次运行时未来得及落库检查 WAL 开启情况确保 wal_enabled 设为 true落库逻辑在所有接收路径最前面4.2 消息重复如何兜底消费端幂等设计hermes-agent 用的是“至少一次”投递语义所以重复消息是必然会发生的事件不是小概率故障。比如消费者处理完成后在返回成功前网络闪断了agent 判定投递失败并重试消费者就会收到两条相同 ID 的消息。这个坑在我做订单同步业务时真实踩过重复投递导致下游把同一个订单状态翻转了两次最后数据对不上。解决思路不在 agent而在消费者。我当时在业务表里加了一张 processed_message 记录表把每次收到的 trace_id 和 msg_id 当作唯一键存储。消费者处理前先查一遍INSERT INTO processed_messages (msg_id, handled_at) VALUES (?, ?) ON CONFLICT(msg_id) DO NOTHING;如果影响行数为 0说明这条消息之前处理过直接返回成功即可。如果影响行数为 1说明首次处理才进入真正的业务逻辑。这种做法不需要分布式锁简单可靠而且能很自然地兼容并发重复投放的情况。强烈建议任何接入 hermes-agent 的消费者都把这一步做上这是兜住全局的最后一道安全网。4.3 下游超时与线程池耗尽问题有一次我把 dispatch_pool_size 设成 16结果下游接口突然慢到 30 秒才返回线程池里 16 个线程全部被占住新消息全部堵在 PENDING 状态消息堆积肉眼可见地往上飙。这个问题的本质是HTTP 调用的阻塞特性决定了线程池数量得跟着下游的最大响应时间走。我做的调整有两步。第一步在 agent 配置里把下游 HTTP 调用的超时时间缩短到 3 秒如果 3 秒没响应就立即标记为失败进入重试流程而不是傻等。第二步把线程池容量从固定大小改成可监控、可动态扩容的模型最低 8 个最高 32 个每次新增线程都打印一条日志。这样即使下游抖动也只会占用少量线程不至于全面瘫痪。4.4 消息积压与死信快速清理经验消息积压后最忌讳的是人工一条一条去点。我的做法是给 agent-ctl 增加一个批量重发通道每次从死信表取最新 500 条按创建时间升序重新入队。这样处理后下游不至于被同一个瞬间的并发打爆值得参考./bin/agent-ctl --resend-batch --limit500 --statusdead_letter重发过程中还要盯住 agent 的指标dispatch_active、dead_letter_count、retry_wait_count。如果重发 100 条后死信表数量不降反升说明下游的问题还没解决这时候要立刻停止批量重发先把下游接口修好再慢慢回放积压消息。我这次踩过的坑就是着急回放旧数据结果下游服务的 CPU 直接被顶到 100%问题更严重了。后来我总结了一个原则“先恢复消费能力再恢复消息数据。”5. 关于 hermes-agent 设计的最终思考到这里整个 hermes-agent 的设计和落地就梳理得差不多了。从头到尾我其实没有引入什么高深技术核心思路就四个字异步解耦。把生产者和消费者的关系从“你直接调用我”变成“你告诉我我来安排”这个思维的转换是让系统稳定性上台阶的关键节点。我个人在实际操作中的体会是hermes-agent 这类代理层最适合的场景往往不在高并发大流量的核心链路而在那些“链路不长但分支很多”的业务里。比如订单状态变化后要同时通知库存、通知财务、更新日志、推送短信用同步调用的方式写出来代码会非常痛苦而一旦换成消息代理整个逻辑就变得清爽多了。另外我想多提一句很多人会把 agent 和消息队列混为一谈其实它们解决的问题在粒度上不一样。消息队列解决的是“多个消费者同时读同一个数据副本”的问题而 agent 解决的是“一个事件需要按路由规则到达不同目的地”的问题。如果你手上的场景只是多个消费者各取所需那么 agent 可能大材小用了但如果你有复杂的路由规则、重试策略、回执回调那么 agent 这种轻量代理的价值就体现得非常明显。最后一点经验接手任何 hermes-agent 系统前先把 config 里的路由表读一遍把 dispatcher 的重试参数读一遍这比看业务代码更重要因为大多数故障都出在“消息找不到路由”或者“重试策略不合理”上而不是业务本身。希望这篇文章能给你一个清晰的认知框架不是推荐你马上换某个具体组件而是提供一套思考路径。如果你的项目里也遇到系统之间通知混乱、任务没有兜底机制这些问题试试画一张生产者、代理、消费者关系图也许下一行代码该怎么写你自己就有答案了。
返回列表