ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

自研轻量级发布订阅中间件buzz:从设计到生产实践

自研轻量级发布订阅中间件buzz:从设计到生产实践 buzz 这个名字最初是顺手起的——服务之间传消息像蜜蜂嗡嗡嗡叫 buzz 贴切。三年后再看它已经成了我们团队主力用的轻量级发布订阅中间件从“顺手”变成了“离不开”。上个月有个做 IoT 的朋友问我设备状态同步想从 HTTP 轮询改成消息推送Kafka 和 RabbitMQ 选哪个。我说先停一下把你真实的量级、运维预算和丢消息容忍度列出来再选。他列了三条两千台设备每 30 秒上报一次状态峰值是平时的十倍但丢一两条状态数据无所谓。这几乎就是 buzz 的主场。buzz 是 Go 写的单一静态二进制零外部依赖消息模型只有 topic JSON部署加跑通第一条链路不超过十分钟。本文不贴完整源码但把需求判断、核心设计、部署接入、压测调优和生产踩坑的完整过程摊开讲适合正在评估轻量消息方案的团队也想给想自己写中间件的朋友当参考。1. 为什么放着现成方案不用非要自己写一个 buzz1.1 HTTP 轮询的笨只有被它坑过才懂最开始的系统其实是 HTTP 依赖轮询A 服务每 30 秒调一次 B 服务的状态接口。设备量少的时候没问题但接入设备从几百台涨到两千台后轮询的问题开始扎手。实时性上限就是轮询周期本身你设置 30 秒那消息在管道里最多睡 30 秒浪费也明显大部分轮询返回的都是没变化的数据更难受的是消费方变多时每个新消费方都要让生产方改配置“上线一个新模块要改三个服务配置”这件事干一次就够够的了。那段时间日志里全是 200 OK 的空转我一度怀疑自己在跑的不是业务系统是心跳监控系统。轮询是个老实的方案但老实的代价是所有的业务都按最坏延迟来设计。很多团队一开始觉得“反正数据几秒后才用得上”可一旦业务量上来这“几秒”会在日志里变成几千个无效请求、几台白跑的机器以及若干个“为什么状态总是慢半拍”的工单。1.2 Redis Pub/Sub 很香但它是个火坑我们第一次迁移到 Redis Pub/Sub上手确实快但问题接踵而至。Redis 的 Pub/Sub 是 fire-and-forget 模型消息发出去订阅者掉线就是丢了没有任何补救更危险的是 Redis 不提供背压如果一个消费方处理不过来消息会积压在客户端的输出缓冲区里直接把内存顶上去严重时甚至因为慢客户端占用缓冲而影响其他 key 的读写。其次Redis 的订阅者断开后必须手动重订阅“重连后重新 SUBSCRIBE”这件事每个接入方都要实现一遍而且实现得还都不一样。最让我接受不了的是 Redis 通配符匹配是对所有订阅者扫一遍的模式主题一多性能就下滑也拿不出订阅者维度的指标。它本质是个缓存不是消息总线。硬拿它在业务里扛等于拿菜刀改锥——能拧几个螺丝但迟早把刀口崩了。1.3 Kafka 很强但背包客不需要一辆火车Kafka 那套能力我们到现在都认可持久化、分区、消费者组、消息回放全是好东西。但付出的代价是一套三节点的 Kafka 集群光是操作系统调优、磁盘规划、分区平衡、消费者组 rebalance 这些概念就已经超过我们整个团队的运维负担。我们的平均消息量也就每秒六七十条为了这点量养一套需要专职关注的集群怎么看都不划算。后来我们的结论变成了Kafka 解决的问题是“海量消息 可靠投递 回放”我们实际的问题是“量不大但很碎、希望轻量、允许小概率丢失”。中间隔着一个量级和一层运维复杂度。很多人选型时容易犯的错是按“某某大厂在用”而不是按“我的真实负载和失败容忍度”来选。1.4 收拾心情我们到底需要什么把需求落到纸上其实只有六条多级主题加通配符订阅订阅端持久连接断线自动重连并重订阅at-most-once 投递允许丢单进程撑得住 2000 以上并发连接协议可读方便抓包排查。至于持久化、死信、事务、镜像集群全部明确不在范围里。写这六条的过程就是 buzz 立项的过程。中间件的格局往往不是你想做什么而是你敢不做什么这句话后来成了我们的口头禅。你砍掉的功能会在文档里变成一行清晰的“不支持”而一个没有边界的系统最后一定会变成所有人的垃圾场。2. 核心设计拆解topic JSON 和一个长度前缀帧协议2.1 消息模型只有三个概念buzz 全系统只有三个概念连接、主题、消息。Topic 是 UTF-8 字符串用斜杠分层级发布和订阅都围绕它转。订阅表达式能收到的消息iot/device/001/temp精确匹配该主题iot//temp匹配所有设备的温度消息只通配一层iot/#匹配iot及以下所有层级#通配剩余所有层payload 就是一个 JSON 对象默认上限 1MB。服务端不解析 payload只把它当成字节串原样转发——这是后面压测性能和内存的关键决策之一。没有 Exchange、没有 Binding Key、没有 Headers这一堆 RabbitMQ 的概念一个都不引入。我们踩过不少系统发现只要引入超过两个核心概念学习成本和误用的概率就会指数级上升。2.2 帧协议四字节长度前缀自定义协议是拍脑袋拍出来的但后来证明这个决策值得。我们需要一个能跑在 TCP 上的、可读的、支持“连接后随时发订阅命令”的协议。HTTP/2 和 gRPC 也能做到但对调试不友好protobuf 序列化性能虽好但对“人肉抓包”是灾难。所以最终是四字节长度前缀帧加 JSON{op:sub,topic:iot/#,sid:device-sync-01} {op:pub,topic:iot/device/001/temp,id:f47ac10b-58cc-4372-a567-0e02b2c3d479,ts:1735000000.123,payload:{value:23.5,unit:c}} {op:msg,topic:iot/device/001/temp,id:f47ac10b-58cc-4372-a567-0e02b2c3d479,ts:1735000000.123,payload:{value:23.5,unit:c}}sub/unsub订阅、取消订阅sid是订阅会话标识。pub发布消息id是消息 UUIDts是毫秒时间戳。msg服务端向订阅者投递的消息。ping/pong心跳默认 30 秒对端 75 秒没收到就判死。err错误反馈。协议可读带来的好处是排障时可以直接上一把 ncprintf {op:sub,topic:iot/#,sid:test} | nc -v 127.0.0.1 7600中间件的调试体验就是生产力任何看不懂的协议最后都会变成一堆解释不清楚的 bug。2.3 服务端架构主题树 有界缓冲通道服务端只有四个组件主题树trie按层级匹配#和作为特殊节点。订阅变更走写锁消息投递走读锁把“读多写少”的特点榨干。订阅表sid到订阅信息的映射记录主题表达式、缓冲通道和统计指标。分派循环收到pub后主题树找出所有匹配订阅把消息写进每个订阅的有界通道通道满了就按策略处理。会话管理持有 TCP 连接、心跳计时、重连后的重订阅握手。慢消费者策略有三种drop_newest丢最新保旧消息、drop_oldest丢最旧留新、disconnect踢掉该订阅。默认用drop_newest因为对状态类数据旧的比新的更没价值。2.4 我们主动砍掉的功能持久化不落盘进程重启订阅全没、消息全没。要在意这件事的人第一站就应该去看 Kafka。ack 确认不提供 at-least-once不维护待确认消息队列。消费者组buzz 是 fan-out 模型每个匹配订阅者都会收到全量消息不做负载均衡。集群多实例各自独立、无共享状态。横向扩展就是多起几个进程上游按业务分片。砍掉这些不是为了省事是为了让行为可预期。“不保证什么”写得清清楚楚比“好像支持什么”害人少得多。任何把“不确定”伪装成“可能”的设计最后都会变成深夜的 oncall。3. 从部署到第一个消息全程不到十分钟3.1 编译与启动代码库很小拉下来之后一条 make 就能编译出两个二进制buzz-server服务端和 buzz-cli命令行客户端。前置就一个 Go 1.21 以上版本连 go mod vendor 都不需要。cd buzz make ./buzz-server -listen :7600 -slow-consumer-policy drop_newest启动参数速查参数默认值说明-listen:7600服务监听地址-max-payload-size1048576单条消息最大字节数-sub-channel-size1024每个订阅的缓冲通道容量-slow-consumer-policydrop_newest慢消费者处理策略-metrics-addr:19090metrics 端点裸 Prometheus 格式整个服务就一个进程没有配置文件没有外部 Redis没有后台线程养它基本等于养一个静态文件服务。启动日志会打印监听地址和几个默认参数没有任何玄学输出。3.2 CLI 验证链路终端 A./buzz-cli sub iot/# --verbose终端 B./buzz-cli pub iot/device/001/temp {value:23.5,unit:c}终端 A 会输出[msg] iot/device/001/temp - {value:23.5,unit:c} (latency0.8ms)这个动作我建议每次部署后都做一遍。它同时验证了连接、订阅、通配符、发布和延迟五个环节比任何健康检查接口都实在。3.3 Python SDK 接入业务消费端from buzz import BuzzClient client BuzzClient(127.0.0.1:7600) client.subscribe(iot/#) def on_device_state(topic, payload): device_id topic.split(/)[2] update_device_state(device_id, payload) client.connect() client.run_forever()生产端from buzz import BuzzClient client BuzzClient(127.0.0.1:7600) client.publish(iot/device/001/temp, {value: 23.5})SDK 内部已经把心跳、断线重连、重订阅、时间戳这些脏活封装好了应用层回调只负责业务。有一点必须强调run_forever里的回调如果阻塞超过 1 秒SDK 会打一条 warning 日志。这不是噪音它是在提醒你下游处理速度已经开始跟不上上游了。3.4 参数怎么调先看这条三角关系max-payload-size乘sub-channel-size决定了单个订阅在最坏情况下会占用多少内存。举个极端例子payload 1MB、channel 1024一个订阅者就能吃掉 1GB 内存。生产里我们普遍用-max-payload-size 65536配-sub-channel-size 256配合消费端的分批落库数据型业务从没有因为缓冲满拖垮过服务。这条三角关系的经验值是批量数据、图片、日志片段都不该走这条管道buzz 只运输小而高频的状态和事件。工具的自由度越大越需要使用者自己想清楚边界。4. 压测与调优从 45k 到 100k 的那几步4.1 压测环境和口径压测机是 8C16G 的虚拟机服务端内存限制 512MB消息体固定 512 字节 JSON。压测工具是内部写的 load_gen输出发布速率和投递速率要分开记录——发布速率是生产端每秒发多少条投递速率是服务端向所有匹配订阅者累计下发多少条两者之间就是扇出倍数。4.2 三组基线结果场景发布速率p99 延迟表现30 生产者 / 300 订阅者 / 1000 个 topic45k 条/秒9ms无丢弃50 生产者 / 500 订阅者 / 1000 个 topic82k 条/秒41ms无丢弃延迟上升1 个热门 topic / 200 个订阅者20k 条/秒210ms触发 drop_newest丢弃 1.1%第三行是真正的瓶颈热门主题把一条消息复制 200 份内存和 GC 开销立刻显现。扇出 200 倍的情况下20k 条/秒的发布速率意味着投递速率是 4M 条/秒这是单实例内存型中间件的现实天花板。4.3 三个有效的优化payload 零拷贝转发服务端不解析 payload 字段原文抄送减少了约 30% 的 CPU 开销。批量下发订阅通道不是逐条 flush而是每 2ms 攒一批一次写 socket。热点场景吞吐提升约 40%代价是延迟最多增加 2ms。读写锁分离主题树的订阅变更用写锁消息匹配用读锁。订阅变更本来就极低频这个改动把高并发下的锁竞争干掉了大半。全套优化后常规业务场景30 生产者 / 300 订阅者 / 1000 个 topic从 45k 提升到了 103k 条/秒p99 稳定在 11ms热点场景从 20k 提升到 28k但扇出放大效应依然是无解的物理瓶颈。4.4 没用的优化和教训我们还试过对消息对象做 sync.Pool 池化、用自旋锁替代互斥锁最终都因为代码复杂度上升、压测收益几乎为零而回滚。教训是不要为了优化而优化先跑压测让数据告诉你瓶颈在哪。buzz 的性能提升有大半不是来自更聪明的算法而是来自“少做事”少解析、少拷贝、少加锁。5. 生产踩坑实录四条完整的排查链路5.1 坑一慢消费端把服务端内存顶上去现象某天监控发现 server RSS 在半小时内从 600MB 涨到 2.1GB同时 drop_total 持续增长整个服务 p99 从 10ms 涨到 200ms。排查链路先看全局连接数、主题数没异常说明不是连接风暴再看 metrics 里的 per-subscription drop 计数发现其中一个订阅者的 drop 曲线把其他所有订阅者加起来都碾过去最后确认这个订阅者是数据接入模块回调里做了同步的第三方 API 调用单次耗时最多 500ms。20k 发布速率下它的消费速度只有几百条每秒通道持续打满。drop_newest 策略让它“看起来”没死实际上丢掉的都是最新状态。修复消费端改成异步线程池加批量写回调时间从 500ms 降到 5ms服务端把sub-channel-size降到 256并给 backlog 超过 60% 的订阅者加了一条告警。从此这一类事故再没出现。教训pub/sub 系统的瓶颈永远是下游消费速度。服务端能做的是别撒谎把丢了多少、谁丢了如实记到指标里。5.2 坑二断网重连引发的订阅风暴现象网络设备升级导致所有客户端断线 3 秒随后 5 分钟里服务端 CPU 飙到 100%连接数在 600 到 1200 之间来回震荡。排查链路看日志发现所有客户端几乎在同一秒完成重连形成“惊群”。SDK 第一版的重连策略是固定 1 秒间隔600 个客户端齐齐整整每秒重建一次期间每次重连还伴随全量重订阅主题树的写锁频率瞬间翻了上百倍。修复客户端重连改成指数退避加随机抖动第一次 500ms之后乘以 2 再加 0 到 2 秒随机量最大 30 秒。服务端同时把订阅清理策略从“断线立即删”改成了“断线后保留 60 秒”重连的客户端在重新完成订阅握手之前不会丢消息。之后又模拟拔网线压了一次连接重建曲线平滑p99 只抖了一下就恢复。教训断线重连是消息中间件最容易翻车的高频操作一定要设计成“错峰加抖开”别让重连变成重放攻击式的问题。5.3 坑三#通配符到底匹不匹配父层级现象有团队报故障说iot/#订阅收不到iot这条消息另一个团队却抱怨iot/#把iot这条消息也推过来了。两边都不满意而且都在骂服务端。排查链路两个团队的订阅表达式一样行为却相反直接去看代码——第一版 trie 实现里#匹配的是“当前层级以下的所有层级”不包含父层级本身。后来某次重构为了对齐 MQTT 语义改成了#包含父层级。改了代码却没同步所有人的认知于是出现了“同一版本上线后两边都骂”的局面。修复把语义固定下来并写进文档第一屏iot/#匹配iot和iot/anything/...iot//temp只匹配iot/一层/temp。同时给 topic 命名规范加了建议消息至少两级主题避免订阅者用split(/)[1]取业务字段时越界。教训中间件的任何语义更改都是契约变更要有变更公告而团队里一半人不读文档所以文档里要配树状示例图文字描述覆盖不了所有误读。5.4 坑四多生产者下的乱序与“没收到”现象运维平台发设备重启指令同一条命令被两个生产者各发了一次消费端按到达顺序执行结果设备挂了两个不同版本的参数直接表现成“指令执行错乱”。排查链路buzz 只保证“同一连接内按发送顺序投递”两个生产者从不同连接发到同一 topic服务端按到达时间分发先后顺序天然无法保证。这不是 bug是模型边界。乱序问题暴露后回查代码运维平台的工具包确实没做幂等同一个指令生成了两条不同 msg_id。修复对命令类渠道约定只能由单一生产者实例发消费端对设备指令做(device_id, seq)排序加幂等去重对消息实时性敏感但不要求完整性的模块统一改成“buzz 喊一嗓子状态以主动回查为准”。教训把“不保证全局有序、不保证送达”写进 README 前五行比任何功能都重要。中间件的价值一半是功能另一半是让业务方明白边界。6. 选型边界什么时候别用 buzz6.1 一张表讲清楚业务诉求推荐选型高频小消息、状态同步、事件通知允许小概率丢失buzz / 同类轻量总线订单、支付、库存任何需要可靠签收和回放的场景Kafka / Pulsar复杂路由、死信、定时重试需要灵活队列语义RabbitMQ已有 Redis 且量极小只要最简单广播Redis Pub/Sub接受丢核心判断标准就一句话消息丢了你的业务能不能自己把自己修好。不能就别省那点运维成本。buzz 是喊话工具不是快递公司喊话的内容丢了可以再喊合同原件必须签收留存。6.2 维护三年我们最值钱的经验README 第一屏放“我们故意不做什么”比放功能清单重要。监控只需要四个指标drop_total、backlog_per_sub、connection_count、pub_rate90% 的故障都写在里面。不要在早期做可插拔的存储抽象、插件机制这些设计最后只会变成没人删得掉的累赘。语言选型最该看的是团队能不能改得动而不是谁的性能参数最好看。最后分享一个我很私人的习惯每次新环境部署完 buzz我都会故意拔掉一个客户端的网线盯五分钟监控确认重连曲线和 backlog 走势符合预期。这个动作三分钟成本已经为我提前发现过两次配置错误。中间件的故障总是挑你最忙的时候上门前置演练永远比事后复盘值钱。
返回列表