Golang 实现高性能的消息推送系统架构设计
Golang消息推送系统通过三层架构实现高性能:连接层使用独立goroutine和定时心跳管理WebSocket,元数据存入sync.Map;分发层按topic和userID两级分片,批量推送并控制并发;投递层采用幂等ACK、指数退避重试及WAL日志异步落库,确保可靠性。
聊到 Golang 消息推送,很多人的第一反应就是开协程、堆并发。但真正在生产环境踩过坑的人都知道,并发只是入场券,真正决定系统能不能撑住 30 万活跃 PV 的,是三道闸门:连接怎么管、消息怎么分、失败怎么兜。用错一道,QPS 就卡在半山腰;用对了,4c8g 的机器也能扛得住。

直接说结论:连接层、分发层、存储层,每一层都有自己的门道。下面逐个拆解。
长连接管理:别让 net.Conn 泄漏或卡死
WebSocket 连接不是建完就完事的。如果不加管控,net.Conn 在高并发下会迅速吃光文件描述符,或者让 TCP TIME_WAIT 暴涨到让内核告警。真实场景里,90% 的连接异常都出在心跳和读写超时的设置上——要么没设,要么设错了。
- 每个连接必须配独立的
readLoop和writeLoopgoroutine,千万别共用一个循环来读写。为什么?因为conn.Write()可能阻塞,一旦堵住,心跳 ping 就发不出去了。 - 心跳间隔建议设为 30s,服务端用
time.AfterFunc定时发送websocket.PingMessage,客户端回复 pong 后重置定时器。如果超时未响应,比如连续 90s 没有动静,直接conn.Close(),别犹豫。 - 连接的元数据(
userID、deviceID、tags等)存入sync.Map,不要查 DB 或 Redis。热路径上的一次远程调用,就是延迟放大器。 - 连接对象本身建议用
sync.Pool复用,避免高频分配加剧 GC 压力。但注意:池中的对象用完后要显式清空字段,比如把userID置为 0,否则可能残留脏数据。
消息分发:按 channel 分片 + 批量 publish,别遍历全量连接
把所有在线用户塞进一个 map[userID]*Conn 然后 for-range 推送,这是典型的性能反模式。360 万用户里哪怕只有 10% 在线,单次广播就得遍历 36 万个连接,CPU 直接打满。
- 采用两级路由:先按业务 topic(比如
"order_update")哈希到 64 个 shard,再在 shard 内按userID % 16路由到子队列。每个子队列配一个独立的chan *Message和固定数量的 worker(比如 4 个),保证并发可控。 - Redis PubSub 别每条消息都单独
Publish。改用内存缓冲:同一 channel 的消息攒够 50 条,或者 100ms 触发一次批量redis.Client.Publish。这样一来,QPS 能降 90%,延迟反而更稳定。 - 客户端的订阅关系用
map[string]map[uintptr]struct{}来存——key 是 topic,value 是 conn 指针集合。增删操作使用sync.RWMutex,但锁内只做指针操作,绝不调用 DB 或 RPC。 - 如果需要支持动态标签,比如
tag=ios_vip,用布隆过滤器做预筛选,再加后置校验,避免每次推送都全量匹配字符串。
投递可靠性:ack 必须幂等,失败必须进 retry_queue
"已发送"不等于"已送达"。线上最常见的坑是:没有做 ack 去重,结果用户收到 3 条一模一样的订单通知;或者重试逻辑写在了主线程里,一条慢请求拖垮了整个分发链路。
- 每条消息带一个唯一的
messageID(用uuid.NewV7()或时间戳+序号),客户端收到后主动上报ACK {messageID, userID}。服务端用redis.SetNX("ack:"+messageID, "1", time.Hour)保证幂等。 - 未收到 ack 的消息,30s 后触发第一次重试,采用指数退避策略,最多重试 2 次。重试任务丢进内存延迟队列(用
timer.AfterFunc实现),如果仍然失败,则写入 Kafka 或 RocketMQ 的retry_topic,由独立的消费者兜底处理。 - 消息落库必须走异步路径:先写 WAL 日志(比如用 Badger),再发送到分发 channel。DB 写入操作要加
context.WithTimeout(ctx, 3*time.Second),超时直接丢弃并触发告警,绝不能阻塞主流程。 - 离线消息不能只靠 Redis List 缓存——容量不可控。正确的做法是用带 TTL 的
zset存储未读消息 ID,score 设为过期时间戳,再定时用ZRANGEBYSCORE清理过期数据。
说到底,消息推送真正难的从来不是"怎么推",而是"推错了之后怎么收场"。重试队列堆积时要不要自动降级?用户连续断连 3 次后该不该暂停推送?这些边界逻辑如果没写进代码,压测时永远看不出问题,一旦上了生产,就会给你"惊喜"。


































