
我们团队在日常业务里被消息流转折腾得够呛——Webhook 回调、内部邮件通知、企微/钉钉机器人、Kafka 事件、数据库变更订阅每个系统都要对接不同的协议每接一个新渠道就是一遍重复的适配代码出问题的时候更是七嘴八舌找不到根因。后来我抽时间把这块逻辑单独抽出来做成一个统一的消息代理服务代号就叫hermes-agent。这篇文章就把整个项目从设计、核心实现到部署踩坑的过程完整记录一遍希望对正在被同样问题折磨的开发者有实际帮助。hermes-agent 解决的核心问题并不复杂让业务系统只关心把一条消息交出去而不用关心它最终怎么到达目标——是走 HTTP 回调、邮件、消息队列还是聊天机器人。它适合那些内部系统多、渠道碎片化、又没有精力引入重型集成平台的中小型团队也适合想理解消息代理底层原理、打算自研一套轻量解决方案的开发者。1. 为什么我会动手写一个叫 hermes-agent 的消息代理故事的起点其实是三个具体的痛点每个经历过跨系统对接的开发者应该都不陌生。1.1 痛点复现渠道碎片化引发的维护地狱先看一个我们内部特别常见的场景。订单系统产生一个订单完成事件需要同时做四件事通知仓储系统开始拣货、给用户发一封订单确认邮件、在钉钉群同步一条运营消息、把事件写入 Kafka 供下游做数据分析。如果用最朴素的方式实现代码长这样def on_order_completed(order): # 通知仓储系统 try: requests.post(http://warehouse.internal/api/picking, json{...}) except Exception as e: logger.error(通知仓储失败, e) # 发邮件 try: email_client.send(order.user_email, subject, body) except Exception as e: logger.error(发送邮件失败, e) # 钉钉机器人 try: dingtalk_client.send_webhook(webhook_url, content) except Exception as e: logger.error(钉钉通知失败, e) # 写 Kafka try: kafka_producer.send(order-events, value{...}) except Exception as e: logger.error(写入Kafka失败, e)这段代码第一天写的时候很爽但两周之后就会变成噩梦。第一个问题每次新增一个通知渠道都要改动订单系统的核心代码第二个问题某个渠道超时了怎么办同步调用会让主流程变慢如果改异步又得自己管理任务队列和重试第三个问题完全没有统一的可观测性钉钉机器人挂了还是 Kafka 不可用只能靠人工翻日志。最要命的还不是代码而是这些逻辑散落在十几个服务里。支付系统也发邮件、权限系统也发钉钉消息每个服务的实现方式都不一样出了问题排查起来像破案一样困难。1.2 为什么不直接用现成的集成方案当时我也认真对比过几类现成方案。首先看了 Apache Camel这东西功能确实强大路由的抽象做得非常完善但引入它的代价是学习曲线陡峭而且 Camel 的 DSL 风格偏 Java部署起来也相对重对我们这种以 Python 为主的后端团队不太友好。也考虑过 Node-RED 这种流编排工具做原型演示很好用但生产环境跑高吞吐的消息流转它的性能和可维护性都有点让人心里没底。还专门调研过一些 SaaS 集成平台比如 Zapier 这类问题更直接——内部系统的数据要从内网出去走一遍第三方安全团队那一关就过不了。我当时的结论是我需要的是一个足够轻量、可以内网私有化部署、核心抽象足够简洁的消息代理内核。与其去学习一套重型框架的抽象不如针对自己的场景做一个刚好够用的轮子。1.3 hermes-agent 的设计目标在动手之前我给自己定了七条明确的设计目标后面的所有代码都是围绕这些目标展开的接入门槛要低业务方通过一个简单的 API 就能投递消息不感知下游细节。渠道接入要标准新增一个通知渠道只需要实现一个适配器不改核心逻辑。可靠投递消息必须至少成功投递一次不能因为某个下游抖动就静默丢失。可观测每条消息从收到、路由到投递成功/失败的全链路状态都要能追踪。支持动态路由规则的修改不需要重启服务运营人员也能改。隔离部署可以独立部署成一个中间件服务不强依赖某个具体业务系统。不要过度设计单机够用不上复杂的分布式协调组件。这些目标里至少一次投递这个语义是后来设计里最重要也最容易被忽略的后面我会单独讲。2. hermes-agent 的核心抽象把信使当作一等公民hermes-agent 有一个特别核心的设计决策——不从队列出发而是从信使messenger出发。这听起来有点咬文嚼字但其实是完全不同的建模思路。2.1 消息的统一信封模型传统队列模型里消息只是一段二进制数据路由信息存放在队列的 exchange/topic 中。但在一个多渠道代理的场景里消息本身要携带的东西远比数据内容多发给钉钉群的消息需要知道 webhook 地址和是否 某些人发邮件需要主题、收件人、优先级写 Kafka 需要指定 topic 和 key。所以我定义了这样一个统一的信封结构{ id: b3a2f5e1-8c4d-4a62-9f11-2d7c0a34e8b6, timestamp: 1712764800000, source: order-service, type: order.completed, priority: high, headers: {}, payload: { order_id: 202404101234, user_email: userexample.com, amount: 199.0 }, routes: [] }字段其实很少但每个都是经过考量的id是全链路追踪和幂等去重的基础要求全局唯一。type是路由匹配的核心字段类似事件的业务语义。source标记消息来源用于排障和审计。payload就是业务数据本体。routes是消息在代理内部被处理过程中逐步填充的路径记录每经过一个渠道就往里面追加一条投递记录。为什么不直接用 JSON Schema 定义一堆必填字段因为消息代理不是业务系统它不应该去理解消息内容的结构化含义只需要对type等信封字段做路由判断payload保持透明即可。这种信封透明原则让 hermes-agent 可以对接任意业务格式。2.2 通道适配器的统一接口有了统一信封接下来每一个渠道都是一个通道适配器Channel Adapter。我定义了一个极简的接口class BaseChannel(ABC): 所有通道适配器的基类 name: str abstractmethod async def deliver(self, message: MessengerMessage) - DeliveryResult: 把一条消息投递到目标渠道。 返回 DeliveryResult(successTrue/False, err_msg...) pass abstractmethod async def validate(self, config: dict) - bool: 校验通道配置是否合法在加载配置时调用 pass接口刻意保持得非常窄。这背后是一个重要的设计权衡很多人做这类系统时习惯在基础接口里塞满 init、health_check、retry_policy、timeout 等一堆方法结果每个适配器都要实现一堆自己根本用不到的函数。hermes-agent 把横切关注点重试、超时、限流放在代理内核里统一处理适配器只需要关心一件事——把消息送出去。用邮件通道举个具体例子实现一个EmailChannel无非就是填上 SMTP 配置然后在deliver里调用邮件库发送。钉钉通道则是构造一个 Markdown 消息体POST 到 webhook。整个核心逻辑都被收敛到deliver这一个函数里所有适配器的代码量都很小单个通道最长不超过 150 行。2.3 路由规则引擎的 DSL 设计路由引擎解决这样一个问题一条消息进来之后要发给哪些通道规则定义采用 YAML因为 YAML 对非程序员也友好rules: - name: order-completed-notify match: type: order.completed actions: - channel: dingtalk-ops template: templates/order_notify.md - channel: email-user template: templates/order_email.html - channel: kafka-events topic: order-events conditions: - field: payload.amount op: gte value: 100每条规则包含三部分match定义触发条件actions定义要投递的通道列表与使用的模板conditions在 match 的基础上做更精细的条件判断。条件判断引擎支持简单路径取值和比较操作符eq、neq、gt、gte、lt、lte、contains、startswith。不需要支持复杂的嵌套与或逻辑真需要复杂判断的业务方应该在自己的服务里完成而不是把业务规则塞给代理。这就是前面说的边界感。这里的模板系统也值得一提。实际项目中不同渠道对同一事件的展现方式完全不一样——钉钉希望一段简洁的 Markdown邮件希望完整的 HTML 排版。所以模板是为渠道定制的而不是为消息定制的。hermes-agent 在路由时根据template字段找到对应的 Jinja2 模板文件用消息的 payload 渲染出渠道特有的内容然后交给适配器发送。3. 关键实现细节worker 池、重试幂等与持久化框架搭好了接下来聊聊内核部分。这是 hermes-agent 最容易被写崩的地方也是跟玩具项目的分水岭。3.1 基于协程的 worker 池传入消息进入系统后会被投放到一个内存队列里由一组 worker 并发消费。这里我采用的是 asyncio 协程模型而不是多线程。原因很简单消息投递到钉钉、邮件、Webhook、Kafka 都是 IO 密集型操作协程在单个线程内就能很好地把等待时间重叠起来避免了线程上下文切换和 GIL 竞争的负担。class WorkerPool: def __init__(self, queue_size10000, workers16): self.queue asyncio.Queue(maxsizequeue_size) self.workers workers async def start(self): self.tasks [asyncio.create_task(self._worker_loop(i)) for i in range(self.workers)] async def _worker_loop(self, worker_id: int): while True: try: # 从队列取一条消息最长阻塞 1 秒以便检查停止状态 message await asyncio.wait_for(self.queue.get(), timeout1.0) except asyncio.TimeoutError: continue except asyncio.CancelledError: break await self._dispatch(message) self.queue.task_done()队列长度和 worker 数量是两个最关键的调参点。队列太长生产速度和消费速度严重失衡积压的消息会产生过高的延迟worker 太少消费能力不足队列堆积worker 太多又没必要毕竟瓶颈通常是下游通道的响应速度而不是本机 CPU。经过压测我最终把默认 worker 数设为 16队列长度设为 10000单机能扛住每秒数百条消息的正常流转这对绝大多数业务场景绰绰有余。3.2 重试、指数退避与幂等消息投递不可能永远成功所以必须设计重试机制。但简单的重试和可靠的重试是两回事。一开始我做的是固定间隔重试比如每次退 5 秒最多重试 3 次。这个方案在高峰期出过一个严重问题下游系统已经过载了所有 worker 还在争着重试失败的消息直接形成重试风暴把下游彻底打死。后来我把重试策略改成带抖动的指数退避第一次失败后等 2 秒第二次失败后等 4 秒第三次失败后等 8 秒以此类推最大间隔不超过 300 秒每次在基础退避时间上加入 ±30% 的随机抖动def backoff_delay(attempt: int) - float: base min(2 ** attempt, 300) jitter random.uniform(0.7, 1.3) return base * jitter抖动的价值经常被低估。如果没有抖动100 条失败消息会在同一时刻集体重试下游在重试时刻面对的并发压力跟故障时刻一模一样故障永远无法恢复。有了抖动重试请求就变成了平滑的随机分布下游才有喘息空间。对于重试次数上限我在内核里规定默认最大重试 5 次超过之后进入死信队列Dead Letter Queue。死信队列不是说消息就扔了而是把投递详情持久化下来等待人工介入或后续补偿任务。然后就是幂等。我们对外保证的语义是至少一次投递这就意味着消息有可能被投递多次。比如 Kafka 通道成功写入了但由于网络问题响应没回来适配器认为失败了触发重试然后又写入了一次。消费者可能因此重复处理消息。要解决这个问题关键是消息 ID 的透传。hermes-agent 的 envelopeid会在投递时写入各渠道支持的元数据字段——Kafka 消息的 key、HTTP Webhook 的X-Request-ID头、邮件的Message-ID头。下游消费者只要去重时以这个 ID 作为幂等键就可以忽略重复消息。3.3 内存队列之外的持久化保障如果你只把消息放在内存队列里服务一重启队列里还没消费的消息就全丢了。这直接违反可靠投递的目标。hermes-agent 的做法是把消息状态持久化到一个本地 SQLite 数据库。这条消息从进入代理开始就 INSERT 一条记录状态是received路由完成进入某个通道后状态变成dispatching投递成功变成delivered重试超过上限变成dead。CREATE TABLE messages ( id TEXT PRIMARY KEY, create_time TIMESTAMP, state TEXT, rule_name TEXT, payload TEXT, delivery_attempts INTEGER, next_retry_time TIMESTAMP, last_error TEXT );服务启动时代理会扫描状态为received和dispatching的记录重新投递。注意不会恢复已经delivered的记录否则会造成重复投递。这个策略带来的语义是崩溃恢复窗口内最多丢失已投递但尚未更新状态的那一小段记录结合各渠道的幂等机制可以接受。SQLite 够用吗对于单机、每秒几百条消息的规模完全够用。SQLite 写入速度本身就是万级 TPS 的瓶颈根本不在这里。不要一上来就上 PostgreSQL那是给分布式部署准备的后面有需要再说这是不要过度设计原则的体现。另外SQLite 文件的持久化注意放在 Docker volume 里不然容器一重建数据全没了这个坑后面展开讲。3.4 配置热更新与渠道健康状态管理业务方的通道配置比如钉钉 webhook 地址换了一个不应该要求重启服务。我用的是监听配置文件变更 定期重载的机制。每 30 秒检查一次配置文件有没有变化有变化就重新加载规则和通道配置但worker 池和队列不动这样可以避免运行状态的打断。每个通道还有健康状态管理。如果某个通道连续失败超过阈值比如连续 20 次失败代理会把这个通道标记为circuit_open状态后续路由到该通道的消息直接进入延迟重试队列不再实际发送。每 30 秒做一个探测请求探通后自动恢复正常。这个机制本质上就是一个简易熔断器防止故障通道拖垮整个 worker 池。class CircuitBreaker: def __init__(self, failure_threshold20, cooldown30): self.failure_count 0 self.failure_threshold failure_threshold self.cooldown cooldown self.state closed # closed - open - half_open def record_failure(self): self.failure_count 1 if self.failure_count self.failure_threshold: self.state open def allow_request(self) - bool: if self.state closed: return True return False4. 部署落地配置热更新与性能实测数据设计得再好最后都要落到部署和实测上。这一章给出一份可以照抄的部署方案以及一组真实的压测数据。4.1 Docker Compose 部署参考hermes-agent 部署非常简单官方镜像会同时包含代理内核、可视化控制台和一个内置 SQLite。最简化的docker-compose.yml是这个样子version: 3.8 services: hermes: image: hermes-agent:latest container_name: hermes restart: unless-stopped ports: - 8000:8000 # 投递 API - 8090:8090 # 控制台 Web UI volumes: - ./config:/etc/hermes # 配置文件目录 - ./templates:/etc/hermes/templates # 模板目录 - hermes-data:/var/lib/hermes # SQLite 数据文件 environment: - HERMES_LOG_LEVELinfo - HERMES_WORKERS16 - HERMES_QUEUE_SIZE10000 healthcheck: test: [CMD, curl, -f, http://localhost:8000/healthz] interval: 10s timeout: 3s retries: 3 volumes: hermes-data:这份配置里有两个细节值得注意。第一个是restart: unless-stopped消息代理属于基础中间件不能因为一次异常退出就彻底挂掉这个策略保证容器退出后自动拉起。第二个是数据卷hermes-data一定要挂否则 SQLite 文件存在容器可写层容器重建就丢了。我自己就在这个上面吃过亏改了配置重建容器结果积压的待重试消息全部清零后来才意识到是数据卷没挂。投递 API 的接入方式也很简单直接 POST 一条 JSON 即可curl -X POST http://127.0.0.1:8000/v1/messages \ -H Content-Type: application/json \ -d { source: order-service, type: order.completed, payload: {order_id: 12345, amount: 199.0} }响应会立即返回 202并附上消息 ID{id: b3a2f5e1-8c4d-4a62-9f11-2d7c0a34e8b6, status: accepted}业务方拿到这个 ID 就可以去控制台查询投递状态不用关心背后的路由逻辑。4.2 性能压测指标与调参参考部署完成后我专门用 Locust 做了一轮压测。压测场景是模拟订单完成事件通过 HTTP API 投递经过一条包含三个通道钉钉、邮件、Kafka的规则进行路由。硬件环境给的是 4 核 8G 的云主机消息 payload 平均 2KB。并发用户数投递速率条/秒P50 延迟P99 延迟失败率5012168 ms256 ms0.0%200396112 ms578 ms0.0%500385155 ms863 ms0.1%1000340410 ms2450 ms0.4%从数据里可以清楚看到性能瓶颈不在代理本身而在于下游通道的响应速度。500 并发之后投递速率不升反降是因为钉钉和邮件的 API 开始出现响应变慢的情况worker 池中的任务等待时间增加。P99 延迟飙到 2 秒以上就是熔断器开始生效部分请求进入退避重试的结果。根据这轮压测我给团队的建议是单实例的 hermes-agent 可以放心扛住大多数中小规模业务但如果下游通道响应普遍超过 500ms 且峰值投递需求超过 300 条/秒就要考虑把 Kafka 等高速通道拆成分离的 worker 组避免慢通道拖累快通道。另一个调参心得是关于 worker 数和队列长度的联动关系。queue size 如果设置得比 workers 少太多会直接导致生产端阻塞API 响应变慢。我后来采用的经验公式是queue_size workers * 60016 个 worker 对应 9600 左右。这个比例能保证消息在队列里最大等待时间约 1 秒不会造成明显的额外延迟。5. 踩坑记录围绕消息代理的三个真实事故做这类系统的过程中最贵的永远是事故。下面三个问题都是我在真实环境里遇到过、花了不少时间才定位清楚的写出来给大家省点时间。5.1 重试风暴固定重试策略的翻车现场第一个事故发生在灰度上线后的第三天。某天下午 Kafka 集群因为磁盘空间不足开始拒绝写入适配器报错。刚开始问题不大只有少数消息失败。但因为最开始我用的重试策略是固定 5 秒重试失败的消息每 5 秒就冲一次 Kafka失败得越多重试的任务积累得越多Kafka 的负载越高系统陷入了恶性循环。监控看到重试队列从几十条涨到几万条几乎是十几分钟内的事下游 Kafka 的 CPU 被打到 100%其他业务也受到了影响。这次事故的直接教训就是前面提到的重试一定要用带抖动的指数退避而且要加最大重试次数。另外一个隐性教训也很重要——重试风暴发生时最紧急的操作不是去修下游而是先暂停 worker 或拉高熔断阈值让系统冷却下来。我当时手忙脚乱地修 Kafka忽视了代理这边的压力控制导致影响面扩大。后来我把熔断器的触发阈值从连续 20 次失败调低到了 10 次失败信号一旦出现就能更早阻断。5.2 幂等键设计不当导致消息重复消费第二个事故更像一个隐藏的定时炸弹。Kafka 通道适配器把消息投递到order-events这个 topic 时使用了消息payload.order_id作为 Kafka 的 key。看起来挺合理对吧按订单维度保证分区顺序。但存在一类特殊场景订单支付失败后用户重新支付同一个订单会产生两条不同状态的事件order.payment_failed和order.payment_succeeded这两条消息的 key 相同但业务语义完全不同。下游消费者只拿order_id做幂等判断直接把第二条消息当成第一条的重复内容丢掉了。正确的做法是幂等键必须包含事件的唯一语义。我后来把设计改成了消息ID作为幂等键因为 hermes-agent 为每条消息生成的 UUID 已经覆盖了同一次事件的唯一一次投递这个语义。即使业务方要按订单去重也应该在 payload 里额外设计业务幂等键而不是偷懒直接用某个单一业务字段。这个教训也被我写进了使用文档的醒目位置代理层的幂等是投递幂等业务层的幂等是处理幂等两层不能混用。5.3 优雅关闭时的消息残留问题最后一个问题相对隐蔽是在发布新版本时陆续发现的。每次更新 hermes-agent 镜像重启容器总有几条消息莫名消失虽然量不大但对于要保证消息不丢的核心通道来说一条都不该少。排查过程费了不少劲。查日志发现消失的消息都有个共同的规律——它们在重启时刻正处于dispatching状态也就是说已经由 worker 从队列取出、正在投递中。由于容器Docker stop默认发的是 SIGTERM但 Python 进程如果没有注册信号处理函数可能直接被强杀导致 worker 线程里正在进行的 HTTP 调用和后续的状态落库来不及执行。修复方案包含两部分。第一部分是代码侧注册 SIGTERM 信号处理函数收到信号后worker 先等待当前正在执行的任务完成不再从队列取新任务等队列清空后再退出。第二部分是运维侧docker stop时给容器留足够的停止时间通过stop_grace_period配置默认 10 秒换成 30 秒。如果 30 秒后还有任务没完成也要保证 SQLite 里记录的状态是未完成的这样下一次启动时还能恢复重投。async def shutdown(signal, loop): logger.info(收到终止信号停止接收新任务等待正在投递的任务完成...) for worker_task in self.worker_tasks: worker_task.cancel() await asyncio.gather(*self.worker_tasks, return_exceptionsTrue) # 将未完成任务状态重置为 received以便重启后恢复 await self.message_store.reset_dispatching_to_received() logger.info(优雅关闭完成)6. 从 hermes-agent 延伸出的三个优化方向核心功能完成后我根据生产使用中的反馈又梳理了三个优化方向这里一并说下给打算基于类似思路做二次开发的团队一个参考。6.1 插件化通道市场让适配器变成热插拔现在的通道适配器虽然接口已经很薄了但每加一个新通道仍然要改代码重新构建镜像。下一步我计划引入插件机制通道以单独的 Python package 或 WASM 模块发布代理通过配置文件声明启用哪些插件。这样第三方贡献的通道比如 Telegram Bot、飞书、Slack、AWS SNS就可以像安装依赖一样pip install进来不用再碰核心代码。这也是社区类基础设施产品的标配能力。6.2 可观测性OpenTelemetry 与全链路 Trace目前控制台的查询还停留在按消息 ID 查状态的阶段运维体验不够好。我在计划引入 OpenTelemetry 标准把投递 API 收到消息、路由匹配、各通道投递尝试、重试事件都打点输出为 Span。这样一条消息从业务服务发出到所有渠道最终投递完成的完整链路就能在 Jaeger 或 Grafana Tempo 里可视化地展示出来。排障从翻日志变成看链路图效率会有一个质的提升。6.3 多租户与配额管理如果未来 hermes-agent 要作为一个内部消息中台向多个业务团队开放使用多租户隔离和配额管理就是必须的。每个业务团队是一个租户拥有独立的 API Token、独立的通道配置空间、独立的配额每分钟最多投递多少条消息。目前我的思路是在 envelope 里增加 tenant 字段路由和限流都基于 tenant 维度做隔离SQLite 层面增加 tenant_id 索引即可不用引入额外的组件尽量保持轻量。这些方向我目前只在设计文档里做了架构预留已经验证过其中的插件机制在 Python 的 entry point 机制下实现成本很低。有兴趣在这个方向深入的同学可以从这里切入。hermes-agent这个项目走到现在最大的体会就是消息代理这类基础设施设计时的克制比功能的堆砌更难。每一次往内核里加东西之前我都会先问一句——这个逻辑是不是真的应该由代理负责还是应该由适配器、业务方或者模板层负责边界越清晰系统就越容易维护踩的坑就越少。希望这份复盘能给你的自研中间件之路带来一些参考。