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

资讯详情

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

统一消息代理设计与实践:从路由引擎到可靠性保障

统一消息代理设计与实践:从路由引擎到可靠性保障 如果你的系统里订单回调、支付通知、监控告警、定时任务各走各的入口代码里到处都是“收到某某消息后要干什么”的散装逻辑你大概率也会像我一样最终忍不住写一个统一的东西来收口。这个收口的东西就是我这次要分享的 hermes-agent。它本质上是把“消息从哪来、要去哪、由谁来处理”这三件事从业务代码里抽出来的轻量级代理。项目取名借用了希腊神话里信使神赫尔墨斯的名字因为它在整个系统里干的活就是传递信息、裁决路由、保障送达。如果你正在维护一个中后台系统天天被各种回调、事件、任务通知搞得焦头烂额这篇文章会比较对你胃口。我会从设计动因讲到架构拆分再讲路由引擎、可靠性设计最后把上线半年踩过的几个坑一并倒出来。不吹不黑都是实际生产环境里验证过的方案。1. hermes-agent要解决的核心问题消息入口太多、逻辑太散1.1 一个让我决定动手改造的具体场景先说当时的真实处境。订单系统创建订单后需要通知库存服务扣减、通知积分服务累加、给用户发送站内信支付服务回调进来要验签、改订单状态、推送前端监控系统会往一个Webhook地址推送告警告警进来要判断级别、决定是否发短信还有一个定时任务每天凌晨扫描超时未支付的订单。如果你只是把这几条链路画出来会觉得也不复杂。但代码写久了就知道真正可怕的是“新消息来了我该找谁”这个问题。新来的同事想接入一种新的通知类型可能要在四五个服务里翻半天最后发现类似逻辑已经有了但是参数格式不一样、重试策略也不一样。有的地方用Redis Pub/Sub有的地方直接HTTP回调有的地方干脆就是一个函数在另一个函数里被直接调用。那段时间我特别头疼一件事排查问题要反复对时间线。一个订单从创建到支付成功中间跨了好几个服务、好几条消息通道一旦某个环节静默失败恢复现场的成本非常高。于是我给自己提了一个需求能不能做一个统一入口让所有消息都从同一个地方进来由这个入口决定该给谁处理、怎么处理、处理失败怎么办。1.2 为什么选了“代理”而不是再造一个消息队列这个决定当时纠结过一阵子。第一反应是引入Kafka或者RabbitMQ但仔细一盘算团队运维成本不允许而且核心业务体量也远没到需要专业消息队列的程度。我们用的是Redis本身已经承担了缓存、分布式锁等职责再用它的Stream做队列支撑完全够用。另外我需要的不仅仅是“把消息从A传到B”。我希望这个组件具备三件事统一接入、智能路由、执行保障。这比一个纯消息队列多了一层“代理人”的角色。它要理解消息里的内容根据topic和meta判断归谁处理还要负责重试、死信、ACK这些脏活累活。所以我没造一个新的MQ而是做了一个面向业务事件的代理层——这就是hermes-agent最初的定位。最终选型是Python 3.10加asyncio接入层支持HTTP、Redis Stream、本地函数调用三种方式。Redis Stream作为核心的持久化队列确保进程重启后消息不丢。整个代理不依赖额外的存储组件业务服务只需要往HTTP接口或者Redis Stream里扔一条JSON消息剩下的事情交给agent。1.3 能力边界与不适用的场景这里必须说清楚它不适合干什么。hermes-agent不是高吞吐消息管道如果你期待它扛住每秒几万条的消息量那应该去看Kafka或者Pulsar而不是这个组件。它也没有分布式事务编排能力做不到跨服务的一致性事务回滚。它适合的是几百到几千QPS的中后台场景比如订单事件、支付回调、告警通知、定时任务派发这种量级。它更大的价值在于让团队对“消息的处理方式”有统一的认知所有业务方只关注注册一个handler、声明自己关心什么消息触发后干自己那点事。消息的接收、投递、路由、重试、死信全部收敛到agent内部。上线一段时间之后新业务接事件的处理时间从按天变成了按小时。2. 架构与核心概念一条消息的完整旅程2.1 三层模型接入层、路由层、执行层hermes-agent在逻辑上分成清晰的三层我用一个非常简单的链路来描述消息流转消息源 - 接入层(Entrypoint) - 校验与封装 - 路由层(Router) - 执行层(Executor) - 回执(ACK)接入层是整个系统的入口目前实现了三种。HTTP入口适合外部服务回调Redis Stream入口适合内部服务投递本地函数调用入口适合同进程内的代码直接触发。无论消息从哪个入口进来都会被统一封装成标准消息结构然后进入路由层。路由层是hermes-agent最核心的部分。它拿到消息后不会马上丢给某个handler而是根据配置的规则表达式判断这条消息应该由哪个执行器、哪个handler来处理。匹配逻辑在内存中完成规则在启动时被编译成语法树运行时只做一次求值操作性能和实时性都足够。执行层负责真正跑业务代码。它维护了一个线程池每个handler在独立的worker中执行避免某个handler里的耗时操作阻塞整个事件循环。handler执行完成后agent才会对Redis Stream里的原始消息做ACK确认。如果handler抛异常消息不会被确认而是进入重试流程。2.2 消息结构设计为什么把meta和payload分开这是我在设计初期就确定的一个关键决策。每个进入hermes-agent的消息都遵循下面的JSON结构{ msg_id: 20250214-8f3a2b7c, topic: order.created, meta: { region: cn, amount: 99.9, priority: normal }, payload: { order_id: A10001, user_id: U2333 }, created_at: 1739494800 }topic是消息类型比如order.created、payment.success、alert.critical它决定了消息的“大方向”。meta是路由决策用的附加信息比如区域、金额、优先级路由规则能直接引用meta字段做条件判断。payload是业务真正关心的数据只传给最终匹配到的handler。这样设计最重要的原因是让路由和执行解耦。路由只需要扫描topic和meta不需要关心payload里面的业务字段所以规则引擎可以非常轻量。而业务handler拿到的payload已经是完整数据不需要再去其他地方拉一次。msg_id从消息一进来就固定下来用于幂等去重和全链路追踪。调试的时候只要拿着msg_id就能在日志里捞出这个消息完整走过了哪些节点。2.3 Handler的编写方式业务侧只需要做一件事使用hermes-agent业务方要写的东西很少。我尽量把样板代码降到最低最终一个handler长这样from hermes_agent import MessageContext, handler handler(topicorder.created, whenmeta.region in (cn, sg)) def on_order_created(ctx: MessageContext): order ctx.payload # 这里直接写业务逻辑扣库存、加积分、发通知都行 send_notification(order[order_id], order[user_id]) return okhandler装饰器会在进程启动时扫描自动把函数注册到路由表。when参数是可选的不写就表示这个topic下的所有消息都归它处理。handler的返回值只用来表示“我处理完了”最终ACK由agent控制。我见过不少人第一次看到这个API时觉得太简单担心复杂场景不够用。真实情况是规范业务的复杂度主要落在两个地方一是路由规则怎么定二是handler内部出错了怎么办。这两块都有完善的机制兜底后面我会详细讲。3. 路由引擎设计消息准确到达目标执行器3.1 自研简单规则引擎的原因路由引擎我可以直接用一个现成的规则引擎但对比之后放弃了。像Drools、json-rules-engine这类组件功能很强大可以支持复杂的if-else嵌套、循环、甚至引入外部数据。但它的学习成本和配置复杂度也相应上来了。对于一个消息代理来说90%的路由需求其实就是几个基础判断topic前缀匹配、字段等于、字段在集合里、大小比较。所以我选择写一个极简的表达式解析器只支持有限的语法但足够覆盖绝大多数场景。规则表达式长这样topic order.createdtopic starts_with order.meta.region in (cn, sg)meta.amount 100 and meta.priority ! silentmeta.type not_equals blocked解析器会把表达式拆成词法单元再构建成抽象语法树最终在内存中执行求值。整个过程不到几百行代码但稳定性和可控性都很好。如果未来真的出现特别复杂的路由需求还留了一个扩展口允许handler注册自定义的match函数。这套轻量表达式的优势很直接写配置的人不需要学一门新语言看着就像自然语言。团队里负责接入消息的同事第一次看到这个语法基本都能直接上手。3.2 匹配顺序与优先级路由规则是按顺序执行的但又不像普通的if-else那么简单。我设计了一套“精确优先于通配”的策略。规则顺序表达式作用1topic order.created精确匹配订单创建事件2topic starts_with order.匹配所有订单域事件3topic starts_with payment.匹配所有支付域事件4fallback兜底处理所有未匹配消息在实际匹配时agent会先检测规则中是否存在精确匹配。如果存在优先执行精确匹配的规则如果精确匹配不到再按顺序走通配规则。这样设计是为了防止出现“通配规则先命中导致精确规则永远等不到”的尴尬。每条消息默认只能命中一个handler。这是为了避免同一个事件被多个地方重复处理造成订单状态被改多次或者重复发通知。如果你确实希望一个事件触发多个执行器可以显式配置broadcast: true它会跳过“已匹配”标记让消息继续走后面的规则。3.3 路由不匹配的消息去了哪里总会有一些消息因为规则写错、topic命名不规范等原因没有匹配到任何handler。这些消息不能被丢弃否则出了问题连查的机会都没有。因此我设计了一个内置的fallback处理器命名叫rest_handler所有无法路由的消息都会进入这里。rest_handler做的事情是把原始消息完整写入一个名为hermes.rest的Redis Stream同时发送一条低优先级的告警通知。这样既不影响主流程运行又保留了线索。上线初期这个兜底帮了大忙至少有三四个因为topic拼写错误导致的消息问题都是靠它发现的。如果你不想收到告警可以在配置里把告警通道关掉但Redis存档我强烈建议保留。4. 可靠性设计不丢消息、不乱重试的工程实现4.1 at-least-once 与 ACK 的关系做消息投递首先要直面一个现实没有人能在分布式环境下做到严格意义的exactly-once。网络会超时、进程会重启、Redis连接会闪断这些都会导致消息重复或不确定。我最终选择的是at-least-once语义配合业务侧幂等来保证最终效果。ACK机制是实现at-least-once的关键。消息进入Redis Stream之后会处于pending状态。handler执行成功agent才向Redis发送XACK确认这条消息处理完成。如果handler抛异常或者进程在handler执行中崩溃消息会一直停留在pending列表下一次重新投递时仍然会被拉到。这里有一个容易被忽略的设计ACK的时机一定是handler成功返回之后而不是消息刚被读取的时候。如果你在handler执行前就ACK一旦handler后续崩溃消息就永久丢失了。反过来如果handler确实成功处理了但因为网络抖动导致ACK没送达那么消息会被重复投递一次这个时候业务幂等就能兜住。4.2 重试策略指数退避 抖动 最大次数消息处理失败后不能无限重试也不能立刻原样重试。无限重试会把下游服务压垮立刻重试通常也会继续失败因为问题大概率不是瞬时抖动。我采用的方案是指数退避加抖动retry_delay base * (2 ** retry_count) random(0, jitter)假设base是1秒jitter是0.5秒第一次重试会延迟1秒到1.5秒第二次是2秒到2.5秒第三次是4秒到4.5秒以此类推。加抖动的原因是如果有大量消息同时失败不加抖动的退避会在某一刻形成整齐的重试波峰直接打爆下游。每条消息有一个最大重试次数默认是5次。超过这个次数后消息不会继续重试而是进入死信Streamtopic固定为hermes.deadletter。人工介入处理时只要看这条Stream里的原始消息和重试历史即可。生产环境里死信消息大多是下游依赖暂时不可用或者业务数据异常处理起来很快。4.3 一个生产可用的完整配置一个经过生产验证的配置骨架大概长这样app: name: hermes-agent-demo debug: false entrypoints: http: enabled: true port: 8000 path: /events redis: enabled: true stream: hermes:inbox group: hermes-workers executors: default: type: thread workers: 32 order_worker: type: thread workers: 8 router: fallback: rest_handler rules: - id: order_flow when: topic starts_with order. to: order_worker - id: payment_notify when: topic payment.success and meta.action notify to: notify_worker reliability: max_retries: 5 retry_base_seconds: 1 retry_jitter_seconds: 0.5 dead_letter_topic: hermes.deadletter每个字段都有明确作用。executors里的workers控制handler并发度不建议无脑调大32个线程在4核8G的机器上已经很高调大了反而会增加线程切换开销。router.rules里的每条规则都指向一个执行器同一执行器里的多个handler会排队执行不同执行器之间可以并行。reliability里的重试参数是全局的单个handler也可以通过装饰器参数覆盖。5. 上线后的实测数据与优化5.1 压测环境与结果这个组件写完第一版后我并没有马上部署到生产而是先压了一轮。压测环境是4核8G的容器Redis放在另一台同网段的机器上模拟器的并发请求从100到400递增。压测场景并发RPSP99延迟HTTP入站 - 纯路由不触发业务handler200318045msHTTP入站 - 触发1个轻量handler200214078msHTTP入站 - 触发3个handler含广播2001200125msRedis Stream入站 - 触发1个handler200260060ms从数据能看出两个结论。第一agent本身的性能瓶颈几乎不在路由和解析路由求值就是一次很轻的表达式计算。第二真正吃性能的是handler数量每多一个handler耗时就会明显增加。所以如果你在生产环境遇到性能问题优先优化业务handler而不是折腾代理本身的参数。5.2 生产实测带来的三个优化压测只是第一步上线后的真实流量暴露的问题更有价值。我根据生产数据做了三轮优化。第一轮是线程池参数调整。默认的32个线程在高峰时会占用大量CPU但吞吐反而没有明显提升。后来把大多数执行器改成8个线程CPU占用下降了约30%P99延迟反而降低了。线程多了以后锁竞争和上下文切换会让性能变差并不是线程越多越好。第二轮是批量消费。消息量大的时候逐条从Redis读取的效率很低。我改成在Stream读消息时一次性读取多条然后再逐条分发给handler。这个改动让Redis的往返次数大幅下降整体吞吐提升了大约20%。第三轮是规则编译缓存。启动时把所有的规则表达式编译成语法树并缓存不用每次消息进来都重新解析字符串。这个改动对性能的提升没有前两个明显但代码路径上更干净了也让启动时的规则校验变成可能。现在每次发布前都会在启动阶段直接校验所有规则表达式是否合法不合法直接拒绝启动。6. 踩坑复盘几个真实问题与解决办法6.1 handler抛异常被“静默吞掉”的教训上线不久后的一天我接到反馈说订单通知没发出去但死信队列是空的日志里也没有异常。排查了很久最后发现是executor的异常处理逻辑有bughandler抛出的Exception被捕获后只是打印了日志却没有把消息标记为重试也没有进死信。表面上看消息是被正常消费了实际上什么都没干。这个问题暴露了一个设计缺陷我最初把“所有异常”都当成“可重试异常”来处理但某些异常根本不该重试比如参数缺失、业务校验失败。修复方式是引入两个异常类一个是RetriableError表示这次失败可以重试另外一批被认为是永久性错误的异常直接进死信避免无意义的重试增加下游压力。现在的代码很明确handler里推荐用raise RetriableError(下游超时)表达“请重试”其他非预期异常默认也走重试但连续重试5次仍然失败自动进死信。这样就再也不会出现“消息悄悄消失”的情况。6.2 重试风暴差点把下游告警打爆还有一次印象特别深。某个下游服务因为发布问题短暂不可用大量handler同时抛异常触发了整个agent的重试机制。因为重试策略里有抖动短时间内没有出现整齐的波峰但还是造成了短时间的重试流量放大下游服务的告警系统差点被打爆。这次之后我加了两道保险。第一道是全局重试速率限制agent每秒最多发出的重试消息数不能超过一个阈值超过的部分继续在pending里等待。第二道是给执行器加熔断开关当某个执行器的失败率在连续时间窗口内超过50%这个执行器会暂停接收新消息直到失败率恢复正常。熔断期间消息不走重试而是进入短暂延迟的保留队列等下游恢复后再继续。这个机制上线以后再遇到下游故障agent的表现稳定了很多。它不会盲目地去“抢救”每一条消息而是有节奏地试探。值得强调一点重试不是越快越好更不是越多越好有序比有力重要得多。6.3 多实例部署导致的消息重复消费严格来说这不完全是bug而是at-least-once语义的必然结果。部署两个agent实例之后同一个Redis Stream会被两个实例同时消费。正常情况下消息会分散到两个实例上但一旦某个实例在handler执行完、ACK发送前崩溃这条消息就会重新被另一个实例消费于是同一个handler被调用两次。如果是发通知这种操作重复一次影响不大但如果是扣库存、加积分这种涉及金额的操作重复消费就是事故。最终靠的是业务侧的幂等每个handler都必须基于msg_id实现幂等。我提供了一个内置的小工具handler可以调用is_duplicated(msg_id)来判断这条消息是否已经处理过。它利用Redis SET NX命令实现如果msg_id已经存在说明处理过了直接返回成功即可。这里我想提醒一下幂等不能只在agent这一层做。任何异步系统只要涉及资金、库存、积分这类敏感操作业务自身也必须做好幂等保护。组件可以帮你减少重复但最终防线一定是业务代码。6.4 配置大小写问题引发消息“漂移”另一个很隐蔽的问题是大小写。某次区域参数从上游传过来是CN但路由规则写的是meta.region cn结果所有大区消息全部命中不了精确规则顺着通配规则跑到了另一个handler产生了错误行为。这个问题暴露出规则解析时缺少对大小写的统一约束。修复方案是在启动阶段做规则校验时自动把已知的枚举字段统一大小写。同时我增加了一个dry-run工具可以拿一条真实消息在本地跑一遍路由把匹配结果打印出来写规则的人能立刻看到自己的逻辑是否会按预期走。现在所有的规则变更都会先在测试环境跑一轮dry-run把线上捞出的真实流量回放一遍确认匹配结果没有变化后再发布。这个习惯帮我避免了好几次因为规则误写而引发的事故。6.5 规则更新后线上消息的“漂移”问题规则变更不只是在测试环境验证就万事大吉。有一次我在线上调整了某个topic的匹配模式把原来精确匹配的规则改成了前缀匹配。结果原来被打入死信的一批历史消息在新规则下变成可路由状态全部被重新投递了一次导致下游收到了一堆重复的旧订单通知。这次经历让我养成了一个习惯改变路由规则时必须评估存量消息的影响。hermes-agent现在支持在规则上增加生效时间配置了effective_from的规则在时间未到达前不会参与匹配。发布变更时我会错开时间窗口确保存量消息不会再被新规则影响到。另外生产环境建议保留“流量回放”机制。每天定时拉取一小部分真实流量在预发布环境的agent上跑一遍对比当前规则和历史规则的路由结果。如果差异超过预期就说明规则变更有可能影响线上行为需要进一步检查。这套机制虽然费力但对消息中台类组件来说非常值得。7. 给hermes-agent的后续扩展留的问题组件到现在运行了半年多整体稳定。我一直在思考几个可以继续扩展的方向也分享出来供参考。第一个方向是把HTTP接入层做成标准Webhook网关支持签名校验、IP白名单、限流。很多外部系统的回调需要这些能力目前是前置了一层Nginx在做如果能收敛到agent内部部署和配置会更统一。第二个方向是可观测性。现在agent暴露了基础的Prometheus指标包括消息进入量、路由命中量、ACK数量、死信数量。我下一步想增加按topic维度的聚合指标这样一眼就能看出哪种消息的失败率最高而不是出了问题以后再翻日志。第三个方向是可视化规则管理。目前配置还是YAML文件改动需要重新发布。对于中小团队来说这够用但如果agent未来要开放给更多业务方使用图形化编辑和实时生效会是刚需。这些扩展我都已经列入计划但说实话并不急着全部做完。做组件类的东西最重要的不是功能多而是稳。现在这套设计在中小规模场景下已经被验证是可靠的如果业务量真的涨到现有架构撑不住下一步不是给agent打补丁而是考虑引入更专业的基础设施了。最后分享一个小经验我习惯在生产环境多占一条Redis Stream专门记录所有消息的路由结果每天花几分钟过一眼流量日志。这个习惯帮我在问题真正影响用户之前就发现过好几次苗头。消息代理这种基础设施平时感受不到它的存在但你对它的运行规律越熟悉系统出问题时你能越快找到答案。
返回列表