hermes-agent:轻量级消息代理与数据流转管线实践指南
2026/9/9 3:36:36 网站建设 项目流程

1. 项目定位:hermes-agent 到底想解决什么问题

先说结论:hermes-agent 是一个面向系统间数据搬运场景的轻量级消息代理与任务执行组件。它做的事情不复杂——把各种来源的数据收进来,按照你预先定义的规则做过滤、转换、路由,然后投递到指定的目标系统里,并且全程留痕、可重试、可观测。

我最早动手做这个项目,是因为手头同时维护了六个内部系统,它们之间互相要数据的方式五花八门:有的走 HTTP 回调,有的往 Redis 队列里塞消息,有的则是直接连数据库轮询。每个系统都自己写了一套消息接收和推送逻辑,重复代码一大堆,而且一旦某个下游系统短暂不可用,消息就丢得无声无息。排查问题时,你根本不知道一条数据到底是在哪个环节丢的。

hermes-agent 的定位就是把这一层“消息接入、处理、分发”的逻辑统一收口到一个独立服务里。你可以把它理解成一个带智能路由功能的快递中转站:上游只管把包裹交给中转站,至于包裹怎么分拣、走哪条线路、送到哪个目的地、没送到怎么重试,都由中转站统一负责。取这个名字也是借用 Hermes 作为神使的意象——它不生产信息,只负责把信息准确、可靠地传到该去的地方。

适合读这篇文章的人,我猜大概是这几类:一是后端开发同学,正在为系统间数据同步、事件通知、任务分发这些事头疼;二是架构师,想找一个轻量方案替代那种“每个微服务都自己连 MQ”的混乱局面;三是运维或全栈出身、需要在有限的资源里快速搭一个消息中转服务的实践者。读完你至少能搞清楚 hermes-agent 的模块是怎么拆的、消息流经各个节点时发生了什么事、以及哪些坑是我已经替你踩平的。

2. 整体设计与核心思路拆解

2.1 为什么做成“代理”而不是直接上消息队列

很多人在设计这类系统时,第一反应是“直接部署一套 RabbitMQ / Kafka 不就行了”。这个思路没错,但要注意它解决的是“传输”问题,而不是“接入和适配”问题。实际工作中我们的大部分痛点恰恰集中在接入侧:上游系统可能是老旧的单体应用,只支持 HTTP POST;某个设备网关只认 MQTT 协议;还有一部分数据源根本没有主动推送能力,只能靠定时去拉文件。

hermes-agent 作为一个代理层,它的核心价值在于把“接收方式”和“处理逻辑”解耦。agent 本身不维护持久化队列,也不具备消息堆积能力,它就是一层无状态的转发与编排层。如果你已经有 MQ 或数据库作为可靠存储,agent 的角色更像是一个聪明的消费者和生产者——从源头拿到数据,做完加工,再交给下一个环节。这样设计的好处是:agent 挂了,上游的 MQ 和下游的业务库都不受影响;agent 扩容,只需要多部署几个实例,配合路由规则做负载分摊就行。

我见过不少团队在项目初期直接上 Kafka,结果整个团队要为分区策略、消费位点、消息顺序这些概念背上包袱。而 hermes-agent 反而适合那种“消息量一天几十万、对顺序不敏感、但对接方特别多”的场景,它用一个统一配置就能把几十个接入端点管起来,不用为每条链路单独写代码。

2.2 整体模块划分:接入层、管线层、执行层、管理面

hermes-agent 的代码结构我拆成了四个相对独立的模块。第一个是接入层(Listener),负责监听各种协议来源,包括 HTTP Webhook、MQTT Topic、Redis 队列、本地目录文件等。接入层只做一件事:把不同协议的数据统一封装成内部的标准消息对象,然后丢给管线层。

管线层(Pipeline)是 agent 的核心,也是它区别于普通消息转发组件的关键。每条消息进入管线后,会依次经过过滤器(Filter)、转换器(Transformer)、路由器(Router)三个环节。过滤器用来丢弃无效数据或重复数据;转换器负责字段映射、格式转换、数据补全;路由器则根据预设规则决定这条消息去往哪个目标系统。这套设计借鉴了 Apache NiFi 的处理器流概念,但做了一个轻量级实现——配置只需要写一个 YAML 文件,不需要拖拽画布。

执行层(Dispatcher)负责真正的数据投递。它管理着所有目标端的连接池、超时参数、重试机制和限流策略。我把执行层单独抽出来是因为它最容易出问题:下游系统不稳定、网络超时、返回格式诡异、认证失效,这些都需要在这一层做兜底处理。

管理面(Admin)则是一组 REST API 加一个极简的前端页面,用来查看消息流转状态、配置管线和手动触发重投递。不要小看这个管理面,线上排查问题、运营同学偶尔手动补偿数据,都靠它。没这个界面的版本我曾经被业务方追着问了三天“那条订单数据到底有没有到你们这边”。

2.3 技术选型的基本考量

技术栈上,hermes-agent 主进程用的 Python 3.11 + FastAPI,管线执行框架是自己写的一套基于 asyncio 的任务调度器,消息持久化用的是 SQLite(生产环境建议切 PostgreSQL),缓存和分布式锁用的 Redis。选 Python 而不是 Go 或 Java,不是因为性能不重要,而是因为这类代理组件的瓶颈几乎全在 IO 和外部系统上,Python 的 async 模型完全跑得动,但开发效率和对各种协议库的兼容度(paho-mqtt、aiohttp、redis-py)明显更高。

FastAPI 承担了两个职责:对外提供 Webhook 接收端点,以及为管理面提供 API。管线调度部分我特意没有使用 Celery。原因很简单:Celery 的重心在分布式任务异步执行,它的任务队列模式并不适合做“流式消息逐个处理”的场景,而且引入 broker 和 worker 的概念会让部署变复杂。我用 asyncio.Queue 加上一组 worker 协程实现了轻量级的并发处理模型,每一路管线有独立的队列和 worker 数量配置,互不干扰。单机实测下来,接入 3 路 HTTP、2 路 MQTT、总共 1000 条/秒的消息流转,CPU 占用不到 20%,内存稳定在 300MB 左右,对绝大多数内部系统来说完全够用。

3. 核心模块解析与配置说明

3.1 接入层:统一协议与消息封装

接入层是整个组件里适配代码最多的地方,因为每个协议的行为模式完全不同。HTTP 接入最直观:agent 暴露一个/webhook/{pipeline_name}端点,上游系统用 POST 把 JSON 丢过来,agent 立刻把 body 封装成标准消息对象,进入对应管线的队列。这里有一个容易踩坑的细节:你是先回复上游“我收到了”,还是先处理再回复?hermes-agent 的设计是收到消息后立即落库并返回 200,异步再进入管线处理。这样即使后续处理逻辑崩溃,消息也已经持久化,不会因为 HTTP 超时导致上游以为失败而重复推送。

MQTT 接入稍微啰嗦一点。你需要配置 broker 地址、端口、Topic 列表和 QoS 等级。QoS 0 在网络抖动时会丢消息,QoS 2 虽然绝对可靠但会引入大量确认报文,实测在物联网设备数据采集场景下 QoS 1 是最平衡的。代码里我用paho-mqtt的线程接口,配合一个内部 buffer 把收到消息塞进 asyncio 队列,注意这里必须处理背压问题——如果 MQTT 消息速率超过管线处理能力,队列长度会持续增长,所以我在监听器里增加了队列最大容量的配置,超限时暂停 MQTT 消费。

Redis 队列接入和文件监听接入逻辑上都是“定时或阻塞拉取”的模式,实现相对简单,但有一个共通点需要留意:消费完一条消息后,务必要在业务逻辑全部成功之后再确认(ACK)或将文件移动到 processed 目录。如果一半成功就 ACK,消息就彻底丢了。刚开始我把 ACK 放在入队之后、处理之前,结果 Redis 里存的待处理消息和实际入库的数据老是对不上,定位了很久才发现是确认时机不对。

标准消息对象是接层入最关键的数据结构。它长这样:消息 ID(UUID)、来源标识、到达时间、管线名称、原始 payload、解析后的 JSON 数据、重试次数、自定义标签。所有下游模块只认这个对象,不直接操作原始协议字节。定义清楚这个对象后,新接一个数据源只需要写一个 Listener 类做协议适配,管线层和执行层完全不用动。

3.2 管线编排:过滤、转换、路由的配置式实现

管线层是 hermes-agent 的灵魂。我见过很多人做数据接入时,把过滤和转换逻辑散落在业务代码里,每条链路一套逻辑,最后根本没有办法统一管理。管线的思路就是把这些横切逻辑收拢成配置。

管线配置用的是 YAML。一个最小的管线长这样:

pipeline: name: order_sync queue_size: 1000 workers: 4 filter: - type: field_exists field: order_id - type: deduplicate key: order_id ttl: 3600 transform: - type: rename_field mapping: orderNo: order_id customerName: user_name - type: enrich url: http://user-service/api/user/info cache_ttl: 300 field: user_info router: - type: by_field field: channel rules: app: pipeline_app_export web: pipeline_web_export fallback: - type: log_and_drop

过滤器里我用的最勤的是deduplicate这个类型。它底层就是开一个带 TTL 的 Redis Set,消息来了先看 key 在不在集合里,在就跳过,不在就写入并设置过期时间。TTL 的取值需要根据业务去调,我用 3600 秒是应对订单系统那个“上游偶尔重推同一单”的场景。如果你把 TTL 设得太长,Redis 内存会被无意义的消息 ID 占满;设得太短又起不到去重效果,这个需要自己权衡。

转换器里的enrich类型很有意思,它允许你在消息流转过程中调用外部接口补全信息。比如订单消息里只带了 user_id,但下游系统需要完整的用户手机号,就可以在转换阶段调用用户服务接口把信息塞进消息对象。这里我建议一定要加 cache_ttl,否则高峰期每来一条消息就打一次用户服务,压力全转到下游了。cache_ttl 设成 300 秒,配合 LRU 本地缓存,实测能把 80% 的外部调用挡掉。

路由器规则支持按字段值、按消息来源、甚至按正则表达式匹配。如果所有规则都不匹配,消息会走fallback段。我强烈建议不要把 fallback 设为丢弃,至少在调试阶段设为log_and_retry或者转发到一个专门的 dead-letter 管线。线上的脏数据形态千奇百怪,留一条活路比直接丢掉安全得多。

3.3 执行层与重试机制

执行层(Dispatcher)负责把处理完的消息投递到目标端。目标端可以是 HTTP 接口、Kafka Topic、RabbitMQ 队列等。每个目标端有一套独立的重试策略配置:

dispatcher: targets: - name: order_kafka type: kafka bootstrap_servers: "kafka1:9092,kafka2:9092" topic: order_sync_topic retry: max_attempts: 5 backoff: exponential initial_interval: 2 max_interval: 60 concurrency: 8

重试策略是我在这个项目里投入最多精力的一块。早期版本用的是固定间隔重试,比如每 5 秒重试一次,结果下游一旦故障,所有 worker 都在反复打同一个已经挂了的目标,形成重试风暴,还把自己的 Redis 连接池打满了。后来改成指数退避加抖动(exponential backoff with jitter),第一轮等 2 秒,然后 4 秒、8 秒、16 秒、32 秒,每次加上一个随机扰动,把重试请求在时间轴上打散,问题才缓解。

还有一个关键设计是“重试次数超限后怎么办”。max_attempts用光之后,消息会被标记为failed,进入一张独立的重试记录表。管理面可以查这张表并手动触发批量重投。实际操作中这个功能救过我很多次:有一次下游某系统的建表脚本执行错了导致接口一直 500,问题是 DBA 凌晨才修复,而消息早就到执行层了。修复后我在管理页面上选了一个时间段,把这几千条消息一键重新投递,业务方完全无感知。

执行层还需要处理目标端响应格式不统一的问题。HTTP 目标端有人返回 200 就算成功,有人返回 201,有人返回 200 但 body 里带一个"status": "error"。我在目标端配置里加了success_criteria字段,允许自定义什么响应算成功:

- name: biz_api type: http url: http://biz-service/api/v1/orders success_criteria: http_code: [200, 201] body_field: status body_value: ok

这样就不用为每个目标端写胶水代码了。

4. 实操一条完整消息链路:从设备上报到业务入库

4.1 场景设定与初始配置

光说模块没意思,我拉一条真实存在的链路完整走一遍。场景是:现场有一批环境监测设备,通过 MQTT 上报温湿度数据,最后这些数据要被清洗后写入业务系统的 HTTP 接口。这个场景在物联网项目里太常见了,而且它同时覆盖了 agent 的三种接入能力(MQTT、转换、HTTP 分发)。

先看 agent 的主配置文件。MQTT 部分监听设备数据 Topic:

listeners: - type: mqtt name: env_sensors broker: mqtt://10.0.8.10:1883 topics: - "devices/+/telemetry" qos: 1 queue_capacity: 2000

设备会上报到类似devices/sensor_001/telemetry这样的 Topic,payload 可能是:

{"temp": 23.5, "hum": 60.2, "ts": 1698835200}

业务系统需要的格式却是:

{"device_id": "sensor_001", "temperature": 23.5, "humidity": 60.2, "captured_at": "2023-11-01 18:40:00"}

明眼人一下就能看出来,这里至少要做三件事:提取 Topic 里的设备 ID、字段改名(temp 转 temperature)、时间戳转标准时间格式。这正是管线层该干的。

4.2 管线配置与参数选择

针对上述需求,管线配置如下:

pipeline: name: iot_telemetry_ingest queue_size: 3000 workers: 6 filter: - type: field_exists field: temp - type: range_check field: temp min: -40 max: 80 transform: - type: extract_from_topic pattern: "devices/(?P<device_id>[^/]+)/telemetry" target_field: device_id - type: rename_field mapping: temp: temperature hum: humidity ts: captured_at - type: ts_to_string field: captured_at format: "%Y-%m-%d %H:%M:%S"

这里解释几个选择的原因。field_exists判断 temp 字段是否存在,防止设备上报了空 payload 时后面直接异常。range_check是对温度做合理性校验,-40 到 80 摄氏度是这类传感器的合理边界,超过这个范围基本可以判定是设备故障或恶意数据,没必要继续流转了。有人可能会问,为什么不在设备端就直接把 Topic 里的 device_id 塞进 payload?因为很多设备固件是写死的,升级固件的成本远高于在 agent 里做一层转换,这个场景下用extract_from_topic属于最务实的解法。

时间戳转换我单独写了ts_to_string转换器,是因为不同设备的上报时间字段可能是秒级时间戳,也可能是毫秒级。配置里我还留了一个unit参数可以指定,默认是秒。这里有一个隐藏很深的坑:如果设备时间戳是字符串形式的"1698835200",某些语言解析整数没问题,但 JSON 允许的数字精度有限,超过 2^53 就会失真。用字符串作为中间格式流转,可以避免这种精度问题。

4.3 目标端配置与执行效果

数据处理完了,接下来就要投递到业务系统的 HTTP 接口。目标端配置:

dispatcher: targets: - name: biz_api type: http url: http://10.0.20.5:8080/api/v1/telemetry headers: Authorization: "Bearer ${BIZ_API_TOKEN}" http_method: POST retry: max_attempts: 5 backoff: exponential initial_interval: 2 max_interval: 60 concurrency: 4 success_criteria: http_code: [200]

注意这里我用了环境变量${BIZ_API_TOKEN}而不是明文写 token,配置文件会走模板渲染,密钥不落盘,这在多环境部署时是必需的安全习惯。

启动 agent 后,我通过 MQTT 客户端模拟一条设备数据:

mosquitto_pub -h 10.0.8.10 -t "devices/sensor_001/telemetry" -m '{"temp": 23.5, "hum": 60.2, "ts": 1698835200}'

在管理面的实时日志里能看到这条消息的完整时间线:

  • 收到 MQTT 消息,生成消息 IDmsg_8f3a2c,入队;
  • 管线开始处理,field_exists通过,range_check通过;
  • extract_from_topic提取 device_id 为sensor_001
  • rename_field完成字段改名;
  • ts_to_string生成captured_at=2023-11-01 18:40:00
  • 进入分发阶段,POST 到业务接口,返回 200,消息状态置为success

整个过程在我这台机器上耗时约 15 毫秒。如果业务接口返回非 200,消息会进入重试流程,前三次重试的日志会按指数退避的间隔出现,到第五次仍失败时消息会转为failed,等待人工介入。

5. 实战中踩过的坑与问题排查

5.1 重试风暴:从 5 秒一次到指数退避

我前面提到过重试风暴这个事,这里展开说具体现象。第一版的重试逻辑把所有失败消息塞进一个全局队列,固定 5 秒重试一次。某个周一早上,下游订单系统因为数据库连接数耗尽挂掉了,agent 里积压了大概 3 万条消息。这 3 万条消息每 5 秒就对下游发起一次集体冲击,直接导致下游系统在被压垮的边缘反复横跳——稍微恢复一点就被重试流量打垮。

排查时我先在监控面板上看到下游接口的 P99 延迟从 800ms 涨到了 15 秒,然后发现 agent 自身的 CPU 和内存也在快速攀升。当时我第一反应是代码里是不是有死循环,但后来打了一台实例的线程快照才发现,所有 worker 都阻塞在 HTTP 调用上等着超时,任务队列里积压的消息越来越多。

修这个问题的过程让我明白了两件事。第一,重试一定要带退避,指数退避是所有分布式系统教科书里都会讲但很多人懒得做的细节,它确实是保命的设计。第二,重试要按目标端隔离,不能全局一个重试队列。某个目标端挂了,不应该影响其他正常目标端的消息投递。现在的实现里每个 target 都有一个独立的重试队列和独立的退避状态,这类问题基本绝迹。

5.2 消息丢失排查:日志链路追踪的建立

有过一次比较吓人的事故。某个早上业务方反馈说前一天晚上的报表数据少了大概 2%。我第一反应是哪里丢了消息,但查了半天数据库里success状态的消息量确实和上游发送量对不上。

好在接入层开始就给每条消息生成了唯一 ID,并且每经过一个处理节点都会在日志里记录当前状态。通过日志检索,我发现丢失的消息都有一个共同特征:它们的来源都是某个老的 Python 服务,而这个服务用的 HTTP 库在连接被重置后不会正常抛异常,而是静默返回一个空响应——上游代码误以为发送失败,也没有做重试,消息就这么悄无声息地没了。

严格来说这其实不是 agent 的问题,但暴露了一个接入层设计缺陷:代理层只负责收消息,无法得知上游是否真的把消息发出来了。后来我在 HTTP 接入端点加了一个简单的校验逻辑,要求上游请求体里必须带上业务层生成的唯一 ID,如果同一个 ID 重复到达,在过滤阶段就能识别并记录duplicate日志。这样一来,即使上游静默丢了请求,agent 侧至少能通过消息量对比发现异常。在管理面的统计页面上,我现在能看到每个来源的接收量、成功量、失败量、去重量,数据对不上时一眼就能定位到是哪个环节的问题。

5.3 性能调优与资源控制

管线 worker 数量和队列长度的配置最影响资源占用。刚开始我图省事,把每条管线的 worker 都设为 10,队列容量设成 10000。结果在以文件监听方式接入的场景里,因为文件批量导入,瞬间涌入大量消息,10 个 worker 全部忙于处理第一波数据,队列被迅速填满,等到文件监听器再读到新文件时根本进不了队列,消息直接丢弃。

现在的配置策略是:队列容量要大于文件批量导入的最大批次行数,worker 数量则根据下游接口的响应耗时来决定。如果下游接口平均 50ms 响应,单个 worker 的处理能力是每秒 20 条,那么要支撑 100 条/秒的数据量,至少需要 5 个 worker。公式很简单:worker 数 = 目标吞吐量 × 单次处理耗时 / 1000,再乘一个 1.5 的冗余系数。

还有一个容易被忽略的参数是 HTTP 客户端的连接池大小。默认连接池只有 10,在高并发投递时会出现大量 TCP 连接排队等待。我在配置里把max_connections调到了 50,配合 keep-alive,单目标端的吞吐提升非常明显。

5.4 常见问题速查表

现象可能原因排查方法
消息一直处于 pending管线 worker 数过少,队列积压查看队列深度指标,按公式调大 worker
消息反复重试但仍失败目标端成功判定条件配置错误检查 success_criteria 的 http_code 和 body_field 是否和目标端实际返回一致
管理面图表数据为 0agent 实例时钟不同步,统计时间窗口错位检查所有实例的 NTP 状态
同一消息被处理多次上游重复推送,且去重 TTL 设置过短调大 deduplicate 的 TTL,或改用持久化去重存储
MQTT 消息有时收不到QoS 等级设置过低或 broker 断连未重连确认 QoS 至少为 1;检查 MQTT 客户端的 reconnect 机制

6. 个人经验与后续扩展

如果让我说这个项目做得最值的一个决定,那就是把“中间态可视化”做进了核心功能里。一个消息代理组件,最怕的就是被当成一个黑盒——数据进去了就不知道去哪了。管理面上能看到每一条消息从接入到分发的完整流转记录,这个功能在系统出问题时节省的时间,远超开发它所花的时间。

另一个经验是:接入层的协议适配一定要做得薄。我刚开发时想在一个 Listener 里把所有协议细节都处理完,最后代码里堆满了各种 if/else。后来下决心把“协议解析”和“业务处理”彻底分开,每个 Listener 只负责把原始数据变成标准消息对象,其他什么都别干。这个设计决策让后来新增接入源(比如再加一个 WebSocket 数据源)变成了纯增量工作,风险非常可控。

关于后续的扩展方向,我在规划里有几个想法。一是把路由规则从静态配置升级为可热加载,运行时通过管理面 API 修改路由规则,不用重启 agent。二是增加动态扩缩容能力,用 Kubernetes 部署时可以根据队列长度指标自动调整 worker 数或实例数。三是做更细粒度的权限控制,现在管理面的接口是没有鉴权的,在内部环境问题不大,但如果未来要暴露给外部团队使用,必须加上基于角色和管线的访问控制。

其实还有一个小改进特别想做,就是消息内容的字段级血缘追踪。现在能看到消息经过哪些环节,但如果想追溯某个业务字段在转换器里到底被哪一步修改成了什么样子,还做不到。真要实现这个,需要在每个转换器里记录字段级别的变更日志,对存储的占用会明显上升,所以一直没动手。等哪天数据治理的需求足够强烈了,这也许就是 hermes-agent 的下一个亮点功能。

无论如何,这个项目的核心价值始终没变:让消息流转的每个环节都可控、可观测、可配置。如果你也在为各种系统之间剪不断理还乱的数据同步发愁,不妨试试用这一套思路搭一个自己的 agent,也许它不会让你的系统瞬间变得完美,但至少,当一条数据真的丢了的时候,你能像个侦探一样顺着日志把它找回来,而不是对着数据库发呆。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询