资讯动态

hermes-agent:轻量级消息代理与数据流转管线实践指南

发布时间:2026/9/9 3:36:38 来源:尧图企业网站定制
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 里存的待处理消息和实际入库的数据老是对不上定位了很久才发现是确认时机不对。标准消息对象是接层入最关键的数据结构。它长这样消息 IDUUID、来源标识、到达时间、管线名称、原始 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 部分监听设备数据 Topiclisteners: - type: mqtt name: env_sensors broker: mqtt://10.0.8.10:1883 topics: - devices//telemetry qos: 1 queue_capacity: 2000设备会上报到类似devices/sensor_001/telemetry这样的 Topicpayload 可能是{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/(?Pdevice_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_001rename_field完成字段改名ts_to_string生成captured_at2023-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也许它不会让你的系统瞬间变得完美但至少当一条数据真的丢了的时候你能像个侦探一样顺着日志把它找回来而不是对着数据库发呆。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价