1. 项目背景与设计思路
1.1 为什么需要一套自研的推送系统
先交代一下背景。我做的这个项目代号叫“buzz”,取自蜂鸣声——就是一有事情发生,系统能立刻“嗡”一下通知到你的那个意思。当时接手的是一个中等规模的Web平台,后台管理、C端用户、定时任务、监控告警混在一起,业务方提了一堆“实时通知”的需求:用户下单后要给运营弹一条消息、服务器CPU飙高要给值班人员发提醒、用户A关注了用户B要让B的页面实时更新。
把需求盘了一圈,发现市面上的方案有几条路:直接买第三方推送服务,按量付费,一年下来账单不小,而且用户数据要过第三方,合规上还要多填一堆材料;自建一套完整IM系统,明显又杀鸡用牛刀,我们的场景根本没有会话历史、群聊权限这些复杂模型;上Kafka那一套流处理体系,为了一个“实时通知”搭上ZooKeeper、Broker集群,运维成本直接压垮小团队。
所以要的其实很简单:一个能扛住大量长连接、支持按主题订阅、延迟尽量低的轻量级推送网关。项目代号“buzz”,就这么定了。
1.2 核心设计目标与取舍
立项之初我把目标写在白板最顶上,只有三条:轻量、低延迟、好维护。剩下所有设计都围绕这三条展开,遇到拿不准的取舍,就回头对着这三条问一句。
- 轻量:单个服务能独立跑,不依赖外部存储组件,部署时一个二进制文件加一个配置文件就能起服务。
- 低延迟:从业务方调用推送接口到客户端收到消息,目标控制在500毫秒以内,日常要稳定在200毫秒左右。
- 好维护:目录结构简单,日志清晰,出问题时工程师能不看源码就定位到大概方向。
基于这三个目标,我砍掉了“消息持久化”这一块。推送服务不负责把消息存下来,消息存不存、存多久,那是业务方自己的事。buzz只做一件事:把消息从生产者手里接过来,然后尽快送到订阅者的长连接里。如果客户端当时不在线,这条消息就直接丢弃,等客户端重新连上来之后由业务方自行决定要不要补推——大多数场景下,客户端恢复在线时会主动拉一次最新状态,这个兜底逻辑能覆盖掉丢消息的影响。这个取舍在早期评审时被质疑过,但实际跑下来,丢消息的概率极低,业务侧又有补偿机制,反而是这套系统最省心的设计决策。
1.3 整体模块划分
buzz的代码结构按照一条消息从进来到离开的路径拆成四个模块:
- 接入层:负责处理客户端WebSocket连接,包括握手、鉴权、心跳、断线清理。
- 订阅管理:维护主题(topic)到连接(session)的映射关系,支持客户端动态订阅和退订。
- 推送引擎:接收业务侧通过HTTP接口投递的消息,解析目标主题,把消息投递到所有相关连接。
- 连接管理:核心中的核心,负责每个连接的状态记录、缓冲队列、慢消费者处理、优雅关闭。
这四个模块在进程里是完全独立的目录,彼此之间通过内部接口调用,不共享全局可变状态。这样后面做水平扩展时,把推送引擎抽出来单独部署成独立服务,改动成本都很低。
2. 技术选型:为什么是WebSocket加内存队列
2.1 传输协议不是越“高级”越好
先说传输层。实时推送的常用方案大概有四类:WebSocket、SSE(Server-Sent Events)、MQTT、轮询。我直接做了个对比表:
| 维度 | WebSocket | SSE | MQTT | 短轮询 |
|---|---|---|---|---|
| 连接方向 | 双向 | 服务端到客户端单向 | 双向 | 单向反复请求 |
| 浏览器原生支持 | 好 | 非常好 | 需额外库 | 好 |
| 消息格式 | 自定义,灵活 | 纯文本强制UTF-8 | 二进制+主题过滤 | 无 |
| 断线重连 | 需自己实现 | 内置,支持断点续传 | 有遗嘱机制 | 天然无状态 |
| 服务端复杂度 | 中 | 低 | 高(要搭Broker) | 最低 |
| 适合场景 | 通用实时双向通信 | 服务端单向推送、新闻流 | 物联网/弱网环境 | 低频定时拉取 |
最终选了WebSocket,理由很直接:我们有一部分消息确实需要客户端回执或上行指令,比如用户在前端点了“标为已读”,这个动作要回传服务端。SSE做不了上行,MQTT在Web端用起来又重,WebSocket是唯一一个“一条连接把上下行都办了”的方案。代价是服务端要高一点,但高出来的这部分复杂度集中在连接管理上——这本就是我们自己该扛的活,交给框架反而不放心。
2.2 服务端选型:Go还是Node.js
服务端的语言选型,我对比过Go和Node.js两个方向。Node.js的ws库很成熟,单机连接数也不错,写起来也快。但考虑到buzz要处理大量长连接,对内存占用和GC停顿更敏感,我最终选了Go。
Go在长连接场景有三个天然优势:第一,goroutine非常轻量,一个连接对应一两百KB内存,单机能轻松扛几万连接;第二,标准库和官方生态对网络编程支持扎实,WebSocket库虽然要第三方,但golang.org/x/net里维护的那个足够稳定;第三,部署产物是一个静态二进制,丢到服务器上就能跑,不用装运行时环境。
这里补一个踩过的坑:早期我用过github.com/gorilla/websocket,稳定可靠,但后来换成了gobwas/ws,因为后者把协议解析和IO操作做了更底层抽象,内存分配更少,高并发下GC压力明显更小。追求极致性能的可以看看nhooyr.io/websocket,在性能上有些更激进的优化,但社区成熟度稍弱,要自己评估。
2.3 订阅模型:一张哈希表就够用
订阅模型这一点,我见过很多团队一上来就上Redis Pub/Sub,问为什么,回答是“以后要扩展多节点”。但buzz初期就是单机,Redis Pub/Sub反而引入了一个中间件依赖,还多一层网络交互。
buzz的订阅模型设计得朴素但有效:进程里维护两个结构,一个map[string]map[*Session]struct{}存主题到会话的映射,一个map[*Session]map[string]struct{}存会话到主题的反向映射。订阅时加这两个map,退订时删,推送时遍历主题下的会话,复杂度都是O(1)或者O(连接数),简单到不会出错。
只有当单机支撑不住、需要横向扩展成多节点时,才值得引入Redis Pub/Sub或者内网消息总线来做跨节点广播。这个顺序不要反了——系统还没到瓶颈,先把架构复杂度加上去,这是新手团队最容易犯的错。
3. 核心实现:从接入到推送全链路
3.1 客户端接入:鉴权这步暗藏坑
客户端连上WebSocket,第一步是鉴权。buzz采用的是JWT方案:业务方在签发登录态时顺便签一个短期token,客户端连接时把它放到URL的query参数里。
这里有个很隐蔽的问题:浏览器WebSocket API不能自定义Header,很多人想把token放Authorization头里,结果发现原生浏览器根本不让你设。市面上有些库绕过了这个限制,但绕法的原理其实也是先发一个HTTP请求让服务端种Cookie,再升级到WebSocket,多走一跳。所以最省事的方式就是把token放query参数——缺点是token会留在访问日志里,所以buzz对query里的token做了脱敏处理,生产环境的访问日志里一律打码。
鉴权通过之后,连接就进入“可订阅”状态。客户端可以发一个订阅请求,指定要订阅的主题列表,比如{"type":"subscribe","topics":["user.10086.notify","order.center"]}。服务端处理这句话的逻辑很简单:把当前Session分别加入对应主题的映射里。同时维护反向映射,是为了后续某个连接退出时,能迅速从所有主题的Session集合中把自己摘干净,不需要遍历所有主题。
3.2 推送链路:一条消息的完整旅程
业务侧推一条消息,走的是HTTP接口,POST一个JSON体,大结构长这样:
{ "topic": "user.10086.notify", "payload": { "event": "new_order", "order_id": "20240508001", "message": "您有一笔新订单待处理" }, "ttl": 60 }服务端接到这个请求,做三件事。第一,校验topic格式,只能由字母、数字、点、下划线组成,防止注入诡异字符。第二,根据topic找到所有订阅会话。第三,把payload序列化成一段带消息ID的帧,写入每一条连接的写缓冲通道。
这个“写缓冲通道”是整个buzz设计里最值得说的地方。每一条连接上挂一个带缓冲的channel,容量默认128,超出这个值意味着客户端消费不过来——典型的慢消费者。推送时只往这个channel里投递,不直接调用socket写操作。为什么这样设计?因为如果直接同步写在某个连接上,一旦这个连接TCP窗口满了,write会阻塞住,整个推送线程就会被这个慢客户端拖住。用channel一隔,生产者永远不会阻塞,慢消费者只影响自己的channel堆积,不影响全局。
消费端是每个连接一个goroutine,专门从channel里取数据,往WebSocket连接上写。这个goroutine和读goroutine配合,形成了buzz的标准连接模型:一个读循环、一个写循环、一个带缓冲的写通道。
3.3 心跳保持与断线重连
长连接最怕的是什么?静默死亡。客户端网络断了,服务端可能很久都感知不到,TCP层面虽然超时时间能兜底,但往往要几分钟,这段窗口期里的推送全都会被塞进那个没人读的写缓冲里。
buzz的心跳策略是:服务端每30秒发一次Ping帧,客户端收到后自动回Pong帧;服务端如果连续3个周期没收到Pong(也就是90秒),就判定这个连接死透了,直接清理。同时客户端侧由前端SDK实现了指数退避重连:断线后等1秒重试,失败等2秒,再失败4秒,最大间隔60秒封顶,再加上一个0到500毫秒的随机抖动,防止大量客户端同时重连造成“重连风暴”。
这里有个数据值得分享:重连逻辑加上随机抖动之后,我们线上重启服务时,客户端重连的瞬间并发从原来的一下子几千个,被摊开到了几十个几十个的平稳批次,服务端CPU峰值下降了大概60%。
3.4 慢消费者处理策略
慢消费者是长连接系统里最常见的“隐形杀手”。一个连接订阅了热门主题,但客户端网络质量差,或者客户端在后台被系统挂起,导致消息一直堆积在缓冲通道里。等它恢复时,积压的消息瞬间涌入,用户端卡顿、服务端内存飙升。
buzz的兜底策略分两档:通道堆积超过阈值(比如80%)时,不再往这个连接里投递消息,直接标记为“慢消费者”并记录计数;如果连续触发两次,判定该客户端已失去实时接收能力,直接断开这个连接。客户端的SDK重连后会自动重新订阅,但此时只订阅当前需要的主题,不做全量重放,避免再次堆积。
这个策略上线之后,线上再也没有出现过“一个弱网客户端拖垮整个推送进程”的事故。我不止一次在技术群里看到同行讨论“WebSocket推送系统消息积压”的问题,其实根子就在于没有在连接层面做好背压管理。
4. 性能压测与稳定性优化实录
4.1 压测环境与实测数据
buzz做完第一个可运行版本之后,我搭了一套压测环境:一台4核8G的云主机,模拟客户端用Go写了个并发连接器,服务器上用wrk打HTTP推送接口。重点测了两个数据:一是单机能稳定维持多少条WebSocket长连接,二是消息吞吐延迟随连接数增长的曲线。
实测下来,单机稳定维持在6万条长连接,内存占用大概2.1G,CPU在30%上下浮动。推送一条消息到1500个订阅者的场景,P99延迟在180毫秒左右,P50在80毫秒,是在合格范围内的。如果把连接数推到8万,内存到2.8G,已经能感觉到GC频繁了一些,延迟偶发会跳到400毫秒。所以上线时我把单机连接数的容量阈值定在了6万,留出余量承接流量高峰。
4.2 内存暴涨问题:读写缓冲池化
压测中期遇到过一个典型问题:连接数涨到5万之后,内存增长突然变得不太线性,掉头向上的趋势很陡。打内存profile发现两个大头——每个连接上读缓冲和写缓冲各自预分配了4KB和8KB,看似不大,乘以5万就凭空吃掉了600MB。
解决的思路是池化。写缓冲池用sync.Pool管起来,每个goroutine用完归还,早期压测时这块缓存复用率能达到70%以上。读缓冲则是延迟分配,连接握手成功时不确定客户端会不会频繁发数据,先不预分配大块内存,按需增长,空闲连接不占内存。优化之后,同样5万连接的内存从2.3G降到了1.6G,效果好得明显。
4.3 心跳风暴与批量清理
另一个印象深刻的坑,是服务端做心跳超时扫描时出的问题。第一版实现是每个连接一个定时器,到时间就去查一下这个连接的上次Pong时间。5万连接就是5万个Timer,本身就占不少内存,而且每到整点附近的定时器触发会造成一大片goroutine被同时唤醒,CPU曲线跟心电图一样,一突一突的。
后来改成了时间轮算法,精度是1秒。把所有连接按照“上次活跃时间+超时阈值”归到对应的秒级槽位里,扫描线程每秒只处理当前槽位里的连接,摊平了整个负载。改完之后,心跳扫描几乎在CPU统计里看不出单独的波峰了。这个优化算是从Kafka的定时器实现里借鉴来的,通用的思路,换到任何长连接系统都适用。
4.4 优雅重启:先摘流量再断存量连接
上线后最怕的一件事是发布重启时客户端体验“咔哒一下断了”。我最初直接systemctl restart,结果是所有连接瞬间断开,客户端SDK一起重连,服务端短暂涌入几千个握手请求,CPU顶到95%。
后来总结出一套“优雅重启”流程。第一步,通过运维脚本通知注册中心摘掉这个节点的流量,新的WebSocket连接不再路由到本机;第二步,等待30秒或者观察活跃连接数下降到某个阈值;第三步,给所有存量连接发一个“服务端即将重启”的Ping帧,然后一次性关闭;最后才真正停进程。客户端SDK收到这个特殊帧之后,会主动平滑重连到其他节点,而不是等断开了再慌乱重试。
这套流程虽然看起来简单,但解决了线上发布时的两个核心矛盾:存量连接有感知的断开,和增量流量不再进来。实际用下来,发布期间业务上报的推送送达率下降控制在0.1%以内。
5. 常见问题速查与避坑清单
5.1 线上问题排查记录表
下面这张表是buzz上线半年内遇到的几个典型问题的汇总,按排查顺序记录了现象、定位手段和解决方案,以后遇到类似问题可以直接照着查:
| 现象 | 根因 | 排查手段 | 解决办法 |
|---|---|---|---|
| 内存缓慢增长直到OOM | sync.Pool未及时归还或连接未正常清理 | pprof heap对比两个时间点 | 检查连接关闭流程是否走到了资源清理分支 |
| 推送延迟突然飙升 | GC停顿或慢消费者阻塞 | 看GC曲线、检查P99延迟分布 | 池化缓冲、在推送路径上禁用finalizer |
| 某个主题的推送丢失 | 客户端退订和重连竞态 | 客户端日志对比订阅/退订时间戳 | 在重连成功后的订阅请求里加序号,服务端幂等处理 |
| 进程启动时CPU 100% | 大量客户端同时重连握手 | 看连接建立速率曲线 | 客户端SDK增加随机抖动,服务端握手限流 |
| 订阅关系错乱 | 反向映射和正向映射更新不是原子操作 | 代码review + 并发压测 | 订阅管理的写操作统一加锁 |
| 推送消息出现乱序 | 同一条连接有多个写goroutine并发写socket | 抓包看帧序列 | 严格要求每个连接只有一个写goroutine,缓冲通道串行消费 |
5.2 集成阶段最容易忽略的边界情况
新团队接入buzz时,出了问题几乎都出在边界处理上。最常见的一个是:客户端断线重连后,觉得“反正我重新订阅一遍就行”,但服务端可能还残留着旧连接的信息没有彻底清理,导致新连接推送正常后又收到一遍旧连接残留的消息。解决的土办法是,服务端在鉴权阶段如果发现同一用户ID已有活跃连接,强制踢掉旧连接,保证“一用户一连接”的语义。
另一个边界是消息体大小。WebSocket本身对帧大小限制其实比较宽,但过大的消息体(比如超过512KB)会撑爆写缓冲通道。buzz在接入层做了保护:大于64KB的消息体直接拒收并返回错误码,这是从成本和稳定性角度综合考虑出来的阈值。业务侧要传大文件走对象存储,不要把整个文件塞进推送通道。
还有一点要注意的是客户端SDK里的全局异常处理。WebSocket偶尔会出现一些线程栈上很难复现的异常,比如Java客户端里的CloseFrameNotSent,这类异常如果不catch住,会直接把整个调用线程打死,进而影响宿主业务。SDK里针对所有WebSocket回调都包了一层全局异常捕获,记录日志后吞掉异常,避免“推送没做死,先把业务线程搞崩了”。
5.3 踩了几次坑之后沉淀出的几条铁律
多条血的教训总结下来,我给这个项目立了几条铁律,基本不会再犯。
第一,连接清理必须是“全路径覆盖”的。不管连接是正常关闭、异常断线、服务端踢人、心跳超时,还是进程退出——每一条路径都要触发同一套Session清理逻辑。我最初漏了“服务端踢人”这条路径,导致被踢的旧连接残留在主题映射里,直到后来加了一个全量Session扫描任务才兜住。
第二,推送路径上禁止任何昂贵州操作。序列化、压缩、埋点、慢日志,这些统统不要放在单条消息的推送路径上。可以做成异步采样,比如每1000条消息采样一条埋点。你要相信,线上流量起来之后,推送路径上一个微小的开销都会被放大到肉眼可见。
第三,所有限流和熔断都要有“降级预案”。比如客户端重连的指数退避加上限之后,如果服务端仍在过载状态,SDK要能进入“降级模式”,也就是暂停重连期待下次业务事件触发时再拉取一次全量状态。这个降级逻辑平时不起眼,但真到故障演练时能救整个系统一命。
6. 上线前后的一些扩展思考
buzz上线跑稳之后,我开始琢磨它还能干点啥。目前做完了两件顺理成章的扩展:一是从纯HTTP推送扩展出了Webhook能力——某些主题可以挂在HTTP回调上,消息到达时向业务方的回调地址发一个POST请求,相当于“连接不好使的时候,用HTTP这条工业标准通道兜底”。二是支持了多租户隔离:不同业务线在topic命名空间上相互独立,推送接口也需要带上各自的密钥,互不可见。
而对于“离线消息补推”,我目前是建议业务方在客户端重连成功后主动调一次“状态拉取接口”,这个方案比在buzz里内置离线消息缓存简单得多,也让buzz保持纯粹——它就是一个管道,不是一个数据库。
最后说一点个人的体会。推送系统看起来实现简单,但真正的门槛全在“连接管理”这四个字里:连接的生命周期、心跳、重连、慢消费者、异常路径清理。这些细碎的工程问题堆在一起,才构成了生产级推送系统的真实面貌。如果你也要做类似的系统,不要被各种花哨的中间件吸引,先把Core做扎实——连接管理、缓冲通道、一个写goroutine模型、完整的心跳与重连策略。把这些基本功打牢,后面一切扩展都是水到渠成的事。