尧图精选

轻量级消息代理与任务执行智能体Hermes-Agent的设计与实践

🕒 发布时间:2026/9/10 7:21:38 📁 来源:尧图网络
凌晨三点被一条“支付回调积压超时”的告警吵醒起来查日志发现order-service自己的回调重试逻辑崩了任务表里压了几千条数据消息队列里的事件却还在无脑堆积。这种场景我相信做过微服务的兄弟都不陌生——业务一多服务间的通知关系就是一团乱麻。后来我一咬牙写了这个项目就是今天要聊的“hermes-agent”一个轻量级的消息代理与自动化任务执行智能体把消息路由、插件执行、回执追踪这几件事统一收敛到一个中枢里。这个项目解决的核心问题很直接服务之间不再需要互相约定“你调我哪个接口、重试几次、失败怎么补偿”所有事件统一投递给hermes-agent由它按照规则表决定分发给哪个处理器、执行什么动作、以什么节奏重试。对那些受够了回调脚本满天飞、任务状态不可查、告警定位全靠肉眼翻日志的团队来说这套东西能省下大量心力。我自己压测跑过稳定态下单节点吞吐能做到几千条消息每秒规则热更新不需要重启进程插件可以独立注册、独立升级整体足够轻不依赖重型框架。如果你正在做微服务改造或者被异步任务、事件通知、定时动作这一类“脏活儿”折磨这篇分享值得读一读。我会把架构设计的取舍、核心模块的实现思路、上线前踩过的三个大坑以及最终调优的参数都展开聊。1. 一个老问题服务多了之后消息分发为什么变得这么难受1.1 从一次凌晨的告警说起那次事故的根子不复杂下单服务支付成功后发了一个事件回调服务需要拿到这个事件去调用上游发货接口。但代码是一个月前写的当时只有两三个服务大家直接在代码里写了HTTP调用失败就Sleep三秒重试。等业务涨到十几个服务互相调用的链路开始变成蜘蛛网每个服务维护一套自己的重试表状态用Redis加过期时间硬顶日志格式谁也不统一。后来我梳理过一次服务间的调用关系画出来的图基本是“全互联”状态任何一个服务发布都可能导致连锁超时。真正让人崩溃的是排查问题一条消息从产生到最终成功中间经过了哪个服务、哪一步耗时最长、哪一步被静默吞掉异常完全不可见。告警只有“积压数量超过阈值”这种模糊信号连哪个环节堵了都不知道。这种状态下做扩容也好、加补偿也好都是在猜。1.2 市面上的轮子为什么不顺手当时我先把市面上的成熟方案过了一遍发现它们不是不好而是和我的诉求错位了。方案优点和我的需求错在哪RabbitMQ / RocketMQ功能全、可靠投递强部署和运维成本高规则路由要额外开发对一个小团队来说太重Redis Stream轻量、性能好消费者组和消息处理要自己写重试和死信策略不完整Celery任务语义清晰和Python绑定太深消息路由能力弱是任务队列而不是消息中枢NATS / MQTT轻快偏向实时消息推送缺少“执行动作 回执确认”这种任务语义我要的其实是一个既能当消息路由器、又能执行动作的中枢上游只负责发一个事件至于这个事件要不要分发、分发给谁、分发之后由哪个插件处理、处理失败按什么节奏重试都应该在中枢里配置化解决。同时它要足够轻不需要专门养一个中间件团队。这么盘下来自己写一个就变成了最合理的选项于是“hermes-agent”立项。2. hermes-agent的架构分层路由、执行、回执三件事分开做2.1 三个核心角色对应三类实际职责我在设计架构时一直坚持一个原则路由、执行、回执这三件事绝不能耦合在一起。很多项目写着写着就变成“一个大函数处理所有消息”表面看省事实际上后面加规则、加插件、加监控全得动核心代码。hermes-agent把这三块拆成独立角色Router路由层只干一件事——根据规则表把入站消息匹配到对应的插件和队列。它不关心插件内部怎么运行业务也不关心消息最终是否成功只负责“分流”。Worker执行层从队列里拉取消息包裹按包裹里的元信息调用插件。它不关心规则长什么样只负责“干活”。Registry回执与注册中心记录每个消息的完整生命周期包括入站时间、路由目标、执行开始时间、执行结果、重试次数。它负责“证据留存”。这三层之间通过消息包裹Envelope通信。Envelope是一个带元信息的标准结构里面包含消息ID、TraceID、Topic、Headers、Payload、重试计数、截止时间等字段。上游看到的是一个简单的投递接口不需要知道底层有多少队列和插件。2.2 目录结构与工程组织工程组织上我建议按模块边界划分而不是按技术类型划分。一个看起来比较舒服的目录结构长这样hermes-agent/ ├── hermes/ │ ├── api/ # HTTP/gRPC入口接收外部事件 │ ├── router/ # 规则引擎相关代码 │ ├── queue/ # 队列抽象层内存/Redis Stream/Kafka实现 │ ├── worker/ # 消费者、执行器、限流器 │ ├── plugin/ # 插件基类、加载器、内置插件 │ ├── registry/ # 回执存储、进度查询 │ └── observability/ # 日志、指标、链路追踪 ├── plugins/ # 外部插件目录运行时扫描加载 ├── configs/ │ ├── routes.yaml # 路由规则支持热加载 │ └── hermes.yaml # 服务配置 └── tests/这个结构的核心意图是让“加一个插件”和“加一条规则”不碰主进程代码。实际运行时主进程只负责把消息送进正确的队列插件代码以独立文件或模块方式被加载器扫描进registry。后续业务扩展的时候别人只需要往plugins目录放一个实现类然后在routes.yaml里注册不需要重新发版主程序这个体验非常关键。3. 核心实现路由规则、插件注册与任务执行的落地代码3.1 路由规则设计谁来解决“这条消息应该给谁”路由规则是整个系统的“大脑”它必须足够灵活又必须足够简单。一开始我想用简单topic匹配就完事但很快发现业务里经常出现“同一个topic来源不同处理方式不同”的情况。比如同一个order.event来自下单服务的要触发发货插件来自售后服务的要触发退款插件光靠topic区分不了。所以hermes-agent使用“header匹配 优先级”的组合。规则配置长这样rules: - name: payment_callback_router match: topic: payment.* headers: source: order-service event_type: pay_success priority: 100 action: plugin: payment_callback_plugin queue: high_priority_queue retry: 3 retry_backoff: 10s, 30s, 5m timeout: 30000 - name: refund_router match: topic: payment.* headers: source: refund-service event_type: refund_request priority: 90 action: plugin: refund_handler_plugin queue: default_queue retry: 5 retry_backoff: 30s, 1m, 5m, 15m, 30m timeout: 60000规则引擎的处理逻辑分成两步第一步用topic前缀和headers字段做快速过滤筛出候选规则第二步按priority降序排列取第一条完全匹配的规则执行。为什么执行逻辑要放进action而不是放在匹配条件里因为匹配条件关心的是“这条消息是谁发来的”而action关心的是“这条消息应该怎么被处理”职责清晰以后路由规则本身变成了一张可以人为阅读和审计的表。规则文件的加载有一个细节值得提不要每次消息来了都重新读YAML那样性能太差。我用了一个独立的规则管理器启动时加载规则并构建内存中的匹配树之后通过文件系统的watch机制监听变化检测到变更后重新编译匹配树并原子替换。这样规则热更新不需要重启进程代价只是毫秒级的切换时间生产环境实测没有丢消息。3.2 插件机制让agent的能力可持续扩展插件机制是整个系统最有价值的部分。它决定了一个agent项目能否脱离“一次性脚本”的宿命变成真正可持续增长的平台。在hermes-agent中每个插件只需要实现一个固定接口from abc import ABC, abstractmethod from hermes.envelope import Envelope from hermes.registry import PluginResult class BasePlugin(ABC): 所有插件的基类。 property abstractmethod def name(self) - str: 插件唯一名称用于路由规则引用。 property abstractmethod def version(self) - str: 插件版本号升级检查用。 abstractmethod async def handle(self, envelope: Envelope) - PluginResult: 处理入口可以访问Envelope的完整元数据。插件发现机制我用了简单的目录扫描 装饰器注册避免引入重量级依赖框架。插件目录下的每个Python文件可以放置任意多个插件类加载器扫描时检查是否继承自BasePlugin并完成注册。PluginResult包含状态成功/失败/重试/拒绝和描述信息执行层根据这个结果决定是否重试、是否写入回执。为什么要用异步接口而不是同步接口因为worker并发处理时如果插件里存在任何IO操作调上游接口、读写数据库、访问外部存储阻塞都会直接卡住worker的并发能力。我用asyncio让所有插件默认跑在事件循环里IO密集场景下并发能力提升非常明显。实测同一个无IO插件同步模型QPS在300左右异步模型能到1500以上这个差距后面还会细说。3.3 任务执行与回执消息生命周期管理执行层的核心是一个固定流程从队列拉取Envelope - 检查重试次数和截止时间 - 调用插件handle方法 - 根据返回结果更新回执 - 成功则确认消费失败则按策略重新入队或者进入死信队列。为了不让这个流程退化成“一团面条”我把每一步的状态变更显式记录到回执存储中。回执不是一个简单的“成功/失败”布尔值而是一个完整的事件流RECEIVED - ROUTED - EXECUTING - SUCCEEDED或者FAILED_RETRY。这样查询一条消息的经历时可以直接看到它在哪个环节停住了。回执存储我同时支持内存模式和Redis模式。内存模式适合本地开发进程重启后就丢Redis模式用于生产环境key设计为msg:{message_id}:eventsvalue用列表结构存储事件流方便后续追踪和分析。每条消息入站时会生成一个独立的MessageID同时通过TraceID把上游调用链、下游插件执行、回执事件串到一起。4. 上线前必须解决的三个问题并发安全、可靠投递、任务丢失4.1 高并发下的插件状态共享问题这个问题是后续压测时才暴露出来的。最开始几个插件用了一个全局字典存状态单个worker单线程跑没问题但压测场景下worker数一多全局状态出现脏读和覆盖。比如一个统计插件统计“当前处理中的任务数”两个worker同时读到100各自加1后写回最终结果还是101而不是102。解决方式是把插件状态分成两类一类是“无状态逻辑”也就是输入Envelope、输出PluginResult不依赖任何本地存储另一类是“有状态逻辑”必须使用时不要用进程内全局变量而是存到Redis或者数据库里并且带上插件名、业务ID这样的唯一键。我的经验是能做成无状态的就别碰本地状态实在需要状态宁可多一次网络IO走Redis也不要为了省事在worker进程里共享可变对象。4.2 投递失败的重试与幂等重试听上去简单但真正做好很难。最初版本我用了最朴素的“失败就立即重试”结果上游接口短暂抖动时几千条消息在几秒内反复打在同一个接口上老接口直接被打挂。后来老老实实做指数退避并规定一个消息最多重试N次超过就进死信队列。更关键的其实是幂等。因为重试机制天然意味着“同一条消息可能被执行多次”插件如果没做幂等就会出现重复发货、重复退款这种恶性问题。我在Envelope里给每个插件提供了execution_id每次执行生成一个唯一值插件处理前先到外部存储检查这个execution_id是否已经处理过。虽然这增加了一点复杂度但在处理支付类、通知类等对一致性敏感的场景时这笔投入完全值得。4.3 重启场景下的任务恢复第三个坑是重启导致任务丢失。内存模式开发时关掉进程再启动所有队列消息和回执状态全部清空这个其实可以接受。但即使是Redis Stream模式也有一条消息正在执行、进程突然被杀掉、执行结果还没写回执的窗口期。恢复策略如果不处理好这条消息就被永久卡在半路。我的处理方式是给“执行中的消息”加一个持久化的执行标记worker在开始处理前向Redis写入executing:{message_id}设置一个合适的过期时间和插件timeout挂钩处理完成后删除。进程恢复时启动一个重扫任务把所有过期但仍存在的执行标记找出来判定为“可能中断”重新放回队列。这里要注意重放又得益于幂等机制兜底所以最终整体语义是“至少一次投递插件侧保证幂等从而接近精确一次”。5. 实测数据与调优参数不同负载下的表现和配置建议5.1 压测环境与结果目前版本我用的是3节点的容器集群单节点规格2C4GPython 3.11 uvloop队列后端分别测了内存模式和Redis Stream模式。压测工具向入口API持续发送不同QPS的消息每条消息路由后执行一个模拟IO的插件sleep 1ms。模式输入QPS成功处理QPSP99延迟(ms)失败率内存模式1000980800.2%内存模式500046201803.5%Redis Stream10009502100.5%Redis Stream500041004608.2%从数据能看出内存模式在低并发下延迟很低但到高并发阶段失败率开始飙升因为内存队列的消费者处理不过来导致消息溢出。Redis Stream模式多了一层网络IO延迟明显高一些但好处是消费者组管理、消息持久化全都省心。生产环境我建议用Redis Stream模式不差这几毫秒延迟换来的是崩溃恢复不掉数据的确定性。5.2 几个关键配置的调优方向如果只看一个配置我会选worker并发数。它决定了单个节点同时处理多少条消息。太小时消费能力跟不上太大时线程切换开销上升、外部依赖压力陡增。我的经验值是“worker数 CPU核心数 x 5”起步然后根据插件类型调整IO密集插件可以再往上调CPU密集插件必须下调。第二个值得调的是队列长度阈值。Redis Stream的队列可以无限增长但积压太多意味着消息时效性在下降。我设置了两个阈值正常阈值和告警阈值超过后不是简单拒绝而是先通知到运维平台让负责人判断是扩容还是上游降速。第三是批量拉取参数。Redis Stream的消费支持批量拉取一次取几十条再逐条分发能有效减少网络往返。我实测将批量值从1调到20P99延迟下降了约35%代价是单批消息的最大等待时间略微变大对实时性要求不高的场景收益明显。6. 可观测性agent跑起来之后我怎么知道它在干什么6.1 三条链路调用链、消息轨迹、指标监控一个消息中枢如果自己不可观测那就是在给团队挖坑。hermes-agent上线前我就设计了三层可观测体系。第一层是调用链入口API中间件在收到消息时生成TraceID这个ID通过Envelope的headers透传给插件插件发出的任何HTTP调用都自动带上这个ID配合OpenTelemetry一条消息从进来到最终结果的全链路都能串起来。第二层是消息轨迹也就是前面说的回执事件流。每条消息的每个关键节点都会打一个带时间戳的事件存到Redis或者直接发到日志平台。排查问题时输入一个MessageID就能看到每个节点耗时多少、卡在哪里不需要再到各个服务翻日志。第三层是指标监控。我暴露了一个/metrics端点输出Prometheus格式的指标包括处理总数、成功数、失败数、各插件执行耗时直方图、队列积压数等。Grafana面板上挂了几个核心图吞吐量、P99延迟、积压趋势。凌晨再有告警不再是“有积压”这句废话而是直接能定位到是哪个插件变慢了、哪个topic冲量了。6.2 常见的日志乱象与解决的格式化方案日志这块我踩过的坑很多。最开始各个插件直接print输出格式五花八门排查问题时能把人看疯。后来我在入口和出口统一做了结构化日志改用JSON一行输出并且强制加入trace_id、message_id、plugin_name、execution_status这四个字段。{time:2024-06-15T10:30:01.123Z,level:INFO,trace_id:a3f2e1c9,message_id:m_8f92ab1d,plugin_name:payment_callback_plugin,execution_status:SUCCEEDED,latency_ms:15,queue:high_priority_queue}统一之后日志检索效率提升了一大截。以前要grep关键字再对时间戳现在直接按trace_id或message_id查一条消息的完整旅程就出来了。这也引出一个经验任何让团队里所有人遵守的规范都要靠框架强制不能靠个人自觉。日志格式一样就能少吵很多架。在实战中还有一个技巧想分享规则文件一定要纳入版本管理。routes.yaml就是架构的一部分每次调整规则都走代码评审流程而不是登录服务器随手改。有一次我图省事直接在服务器上改了规则没同步仓库结果同事发布时把旧规则覆盖回来了线上消息路由直接错乱那个教训到现在都记得。规则配置即代码这个习惯值得刻进团队规范里。hermes-agent从凌晨那次告警到现在跑了大半年期间加过插件、调过规则、压过性能也是边用边改的状态。回看整个过程它最大的价值不是替代了哪个中间件而是让团队对“消息究竟是怎么流动的”这件事重新有了掌控感。如果你也被服务间回调、通知、异步任务折腾得够呛不妨按这个思路搭一套轻量中枢从一条最简单的路由规则开始大概率你不会想再退回原来的模式。
上一篇/下一篇内容由系统自动关联 返回资讯列表 →