title: 分布式 WebSocket 集群中 Redis PubSub 实现方案
你想通过Redis存储WebSocket客户端的连接节点映射信息,实现集群环境下的跨节点消息推送(比如客户端连在节点A,服务端在节点B,能通过Redis找到节点A并完成推送),这是分布式WebSocket集群的核心实现方案。
Redis的核心优势是高性能K-V读写+丰富的数据结构+发布订阅(Pub/Sub),刚好适配这个场景:用K-V结构存储「客户端标识-连接节点」的映射关系,用Pub/Sub实现跨节点的消息广播/定向推送,整体方案轻量、高效,适配百万级客户端的分布式场景,也是业界的标准落地方案。
下面从核心设计思路、Redis存储结构设计、完整实现流程、Go代码落地(结合之前的WebSocket框架)、优化方案、生产级扩展六个维度讲解,方案适配Redis5.x+,兼容单机/Redis集群,可直接结合之前的Golang WebSocket代码落地。
前置场景回顾
假设你的WebSocket集群有多个节点(Node1/Node2/Node3),部署在不同服务器,每个节点独立管理自身的WebSocket连接:
客户端C1连接到Node1,客户端C2连接到Node2;
当Node2需要给C1推送消息时,无法直接访问Node1的连接池,需通过Redis找到C1的连接节点是Node1,再通过跨节点通信让Node1给C1推送消息。
一、核心设计思路
分布式WebSocket跨节点推送的核心是**「Redis做映射存储+消息路由,节点间做轻量通信」**,整体分2个核心模块:
连接元信息存储:客户端握手成功后,当前节点将**「客户端唯一标识→节点信息」**的映射存入Redis,连接关闭时删除该映射;
跨节点消息推送:当节点需要给某客户端推送消息时,先从Redis查询该客户端的连接节点:
若目标节点是自身:直接从本地连接池获取客户端,完成推送;
若目标节点是其他节点:通过Redis Pub/Sub或HTTP/gRPC将消息推送给目标节点,由目标节点从本地连接池获取客户端完成最终推送。
核心方案选择(2种主流方式)
| 方案 | 实现方式 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|---|
| Redis Pub/Sub | 跨节点消息通过Redis频道转发 | 轻量、无需额外通信服务、适配高并发 | 无消息持久化、丢包风险 | 实时推送、非核心消息(如通知) |
| HTTP/gRPC通信 | 跨节点通过HTTP/gRPC直接调用 | 可靠、有回执、支持重试 | 需开发节点通信接口、部署复杂 | 核心业务消息(如IM聊天、指令) |
| 推荐组合:日常实时推送用Redis Pub/Sub(轻量高效),核心消息用HTTP/gRPC做兜底(可靠),下文主要讲解Redis Pub/Sub方案(落地成本最低,适配绝大多数场景)。 |
二、Redis核心存储结构设计
需存储2类信息,均采用Redis String结构(K-V),支持高性能的单键读写(百万级QPS),避免用Hash/List等复杂结构增加开销,Key命名做统一规范(方便运维和排查)。
1. 客户端-节点映射(核心)
作用:通过客户端唯一标识,快速查询其连接的WebSocket节点信息。
Key:
ws:conn:{clientId}(clientId为客户端唯一标识,如用户ID、设备ID,需业务层保证唯一性)Value:
nodeId(节点唯一标识,如node-192.168.1.100:8080、node-01,需集群内唯一)过期策略:设置过期时间(如300s)+ 心跳续期,防止客户端异常断开(如断电、网络中断)导致Redis中残留无效映射,服务端心跳检测时自动续期。
2. 节点-客户端列表(可选)
作用:快速查询某节点上的所有客户端,用于节点下线时的客户端迁移、节点级广播。
Key:
ws:node:{nodeId}:clientsValue:
clientId1,clientId2,clientId3(客户端ID逗号分隔,若客户端数多,建议用Redis Set存储:SADD ws:node:{nodeId}:clients clientId1)过期策略:无需设置,节点下线时直接删除该Key。
补充:节点ID生成规范(集群唯一)
节点ID必须保证集群内唯一,推荐2种生成方式(按需选择):
静态配置:部署时给每个节点配置唯一ID(如
node-01、node-02),适合固定节点集群;动态生成:节点启动时用
IP+端口生成(如node-192.168.1.100:8080),适合动态扩缩容的K8s集群。
三、完整实现流程(Redis Pub/Sub方案)
结合WebSocket连接生命周期和Redis操作,分6个核心步骤,流程覆盖「客户端连接→跨节点推送→客户端断开」全链路,适配集群环境。
前置准备
每个WebSocket节点启动时,初始化Redis客户端(Go用
github.com/redis/go-redis/v9,推荐v9版本,支持Context);每个节点启动时,订阅Redis专属频道:
ws:channel:{nodeId}(仅当前节点消费该频道的消息,避免消息广播到所有节点);定义跨节点推送的消息格式(JSON),包含目标客户端、消息内容、推送类型等,保证节点间消息解析一致。
步骤1:客户端WebSocket握手成功
客户端向某节点(如Node1)发起WebSocket握手,携带业务层唯一标识clientId(如请求参数
?clientId=user123);Node1完成握手后,创建本地Client实例,加入本地连接池;
Node1向Redis执行2个操作:
SET ws:conn:{clientId} {nodeId} EX 300:存储客户端-节点映射,设置5分钟过期;SADD ws:node:{nodeId}:clients {clientId}:将客户端ID加入节点的客户端集合(可选);
Node1启动心跳续期协程:每2分钟向Redis执行
EXPIRE ws:conn:{clientId} 300,刷新过期时间。
步骤2:节点发起消息推送(本地/跨节点判断)
假设Node2需要给客户端user123推送消息,执行以下判断:
Node2从Redis执行
GET ws:conn:user123,查询到目标节点为node-192.168.1.100:8080(即Node1);Node2对比目标节点ID和自身节点ID:
若一致:直接从本地连接池获取
user123,执行本地推送;若不一致:进入跨节点推送流程(步骤3)。
步骤3:跨节点消息发布(发布到Redis专属频道)
Node2按统一消息格式封装跨节点推送消息(JSON),包含
clientId、msg、msgType(文本/二进制)等;Node2向Redis的目标节点专属频道
ws:channel:{nodeId}(如ws:channel:node-192.168.1.100:8080)发布该JSON消息;Redis将消息推送给唯一订阅者(即Node1),避免其他节点收到无效消息。
步骤4:目标节点消费Redis消息并完成推送
Node1持续监听Redis频道
ws:channel:node-192.168.1.100:8080,收到Node2发布的跨节点推送消息;Node1解析JSON消息,获取目标
clientId=user123和消息内容;Node1从本地连接池获取
user123的Client实例,执行Send方法完成消息推送;若推送失败(如客户端已断开),Node1向Redis删除该客户端的映射(
DEL ws:conn:user123),并返回失败回执(可选)。
步骤5:客户端心跳续期(防止Redis映射失效)
Node1的应用层心跳检测(Ping/Pong)每30s执行一次,若客户端正常回复Pong;
Node1立即向Redis执行
EXPIRE ws:conn:{clientId} 300,刷新映射的过期时间,保证Redis中始终是有效连接。
步骤6:客户端主动/被动断开连接
客户端主动关闭连接,或Node1检测到连接异常(超时、读失败);
Node1执行优雅关闭:关闭本地Client、从本地连接池删除;
Node1向Redis执行2个清理操作:
DEL ws:conn:{clientId}:删除客户端-节点映射;SREM ws:node:{nodeId}:clients {clientId}:从节点客户端集合中移除该客户端(可选);
若客户端异常断开(未触发步骤1),Redis中映射会在300s后自动过期,避免无效数据残留。
整体流程示意图
客户端C1(user123)→ 握手→ Node1 → 写入Redis(ws:conn:user123=Node1)
↑
Node2 → 推送消息给user123 → Redis查映射→ 发现目标是Node1 → 发布消息到Redis频道ws:channel:Node1
↓
Node1 → 监听Redis频道ws:channel:Node1 → 消费消息 → 本地推送C1
四、Go代码落地(结合之前的WebSocket框架)
基于之前的Golang WebSocket百万连接框架,集成Redis(go-redis/v9)实现连接映射存储和跨节点推送,代码为生产级可落地版本,包含所有关键细节(Redis初始化、心跳续期、跨节点发布/消费)。
1. 安装依赖
# Redis客户端(v9版本,支持Context)
go get github.com/redis/go-redis/v9
# JSON序列化(原生encoding/json即可,也可用easyjson优化)
2. 核心配置定义
定义Redis配置、节点配置、跨节点消息格式,抽离为全局配置,方便部署时修改:
package main
import (
"context"
"encoding/json"
"fmt"
"net/http"
"net"
"os"
"sync"
"sync/atomic"
"time"
"github.com/gorilla/websocket"
"github.com/redis/go-redis/v9"
)
// -------------------------- 核心配置 --------------------------
// Redis配置
var redisClient *redis.Client
var ctx = context.Background()
// 节点配置(动态生成IP+端口,集群唯一)
var NodeID string
const (
WsPort = "8080" // WebSocket服务端口
RedisExpire = 300 * time.Second // Redis映射过期时间
RenewInterval = 120 * time.Second // 心跳续期间隔(小于RedisExpire)
RedisChannelPrefix = "ws:channel:" // Redis频道前缀
RedisConnKeyPrefix = "ws:conn:" // 客户端-节点映射Key前缀
RedisNodeClientsKeyPrefix = "ws:node:" // 节点-客户端集合Key前缀
)
// 跨节点推送的消息格式(JSON)
type CrossNodeMsg struct {
ClientID string `json:"client_id"` // 目标客户端ID
Msg []byte `json:"msg"` // 推送消息内容
MsgType int `json:"msg_type"` // 消息类型:1=文本,2=二进制(对应websocket.TextMessage/BinaryMessage)
FromNode string `json:"from_node"` // 发送节点ID(用于排查问题)
}
// -------------------------- 原有WebSocket相关定义(略)--------------------------
// 保留之前的upgrader、Client、clientPool、connCount等定义,无需修改
var (
upgrader = websocket.Upgrader{
ReadBufferSize: 1024,
WriteBufferSize: 1024,
CheckOrigin: func(r *http.Request) bool { return true },
}
clientPool = sync.Map{}
connCount uint64 = 0
)
type Client struct {
conn *websocket.Conn
connID string
clientIP string
ctx context.Context
cancel context.CancelFunc
sendChan chan []byte
ClientID string // 新增:业务层唯一客户端ID(核心,用于跨节点推送)
}
// 保留NewClient、ReadLoop、WriteLoop、Send、Close等方法,后续仅修改关键处
3. 初始化Redis和节点ID
节点启动时初始化Redis客户端、动态生成NodeID、订阅Redis专属频道,作为程序启动的前置操作:
// 初始化Redis客户端
func initRedis() {
redisClient = redis.NewClient(&redis.Options{
Addr: "127.0.0.1:6379", // Redis地址,生产环境用Redis集群地址
Password: "", // Redis密码,生产环境必须设置
DB: 0, // 数据库编号
PoolSize: 100, // 连接池大小,高并发下调大
})
// 测试Redis连接
_, err := redisClient.Ping(ctx).Result()
if err != nil {
panic(fmt.Sprintf("Redis连接失败: %v", err))
}
fmt.Printf("Redis连接成功,地址:%s\n", "127.0.0.1:6379")
}
// 动态生成NodeID(IP+端口),保证集群唯一
func genNodeID() string {
// 获取本机内网IP(避免localhost/127.0.0.1)
addrs, err := net.InterfaceAddrs()
if err != nil {
panic(fmt.Sprintf("获取本机IP失败: %v", err))
}
var localIP string
for _, addr := range addrs {
if ipNet, ok := addr.(*net.IPNet); ok && !ipNet.IP.IsLoopback() && ipNet.IP.To4() != nil {
localIP = ipNet.IP.String()
break
}
}
if localIP == "" {
localIP = "127.0.0.1"
}
return fmt.Sprintf("node-%s:%s", localIP, WsPort)
}
// 订阅Redis专属频道,持续消费跨节点推送消息
func subscribeRedisChannel() {
// 订阅当前节点的专属频道:ws:channel:node-192.168.1.100:8080
channel := RedisChannelPrefix + NodeID
pubsub := redisClient.Subscribe(ctx, channel)
// 验证订阅
_, err := pubsub.Receive(ctx)
if err != nil {
panic(fmt.Sprintf("Redis订阅频道失败: %v, 频道:%s", err, channel))
}
fmt.Printf("Redis订阅频道成功,频道:%s\n", channel)
// 启动协程持续消费消息
go func() {
ch := pubsub.Channel()
for msg := range ch {
// 解析跨节点消息
var crossMsg CrossNodeMsg
if err := json.Unmarshal([]byte(msg.Payload), &crossMsg); err != nil {
fmt.Printf("解析跨节点消息失败: %v, 消息:%s\n", err, msg.Payload)
continue
}
// 执行本地推送
go pushLocal(crossMsg)
}
}()
}
// 程序启动前的全局初始化
func init() {
// 1. 初始化Redis
initRedis()
// 2. 生成NodeID
NodeID = genNodeID()
fmt.Printf("节点ID生成成功:%s\n", NodeID)
// 3. 订阅Redis频道
subscribeRedisChannel()
// 4. Go运行时调优
runtime.GOMAXPROCS(runtime.NumCPU())
}
4. 修改WebSocket核心方法(集成Redis操作)
基于之前的Client结构体和方法,新增ClientID字段,并在连接创建、心跳续期、连接关闭时添加Redis操作,实现连接映射的增删改查。
(1)修改NewClient:新增ClientID参数
// NewClient 创建新的客户端连接实例,新增ClientID(业务层唯一标识)
func NewClient(conn *websocket.Conn, clientIP string, clientID string) *Client {
ctx, cancel := context.WithCancel(context.Background())
return &Client{
conn: conn,
connID: genConnID(), // 保留原有内部连接ID
clientIP: clientIP,
ctx: ctx,
cancel: cancel,
sendChan: make(chan []byte, 100),
ClientID: clientID, // 新增:业务层客户端ID
}
}
(2)修改WsHandler:获取ClientID,写入Redis,启动续期协程
客户端握手时携带clientId请求参数,握手成功后写入Redis映射,并启动心跳续期协程:
// WsHandler WebSocket握手的HTTP处理器,新增解析ClientID
func WsHandler(w http.ResponseWriter, r *http.Request) {
// 1. 解析业务层ClientID(必须携带,否则拒绝握手)
clientID := r.URL.Query().Get("clientId")
if clientID == "" {
http.Error(w, "缺少clientId参数", http.StatusBadRequest)
return
}
// 2. HTTP升级为WebSocket
conn, err := upgrader.Upgrade(w, r, nil)
if err != nil {
fmt.Printf("WebSocket握手失败: %v, clientId: %s\n", err, clientID)
http.Error(w, "握手失败", http.StatusBadRequest)
return
}
// 关闭Nagle算法,优化IO
conn.NetConn().(*net.TCPConn).SetNoDelay(true)
// 3. 获取客户端IP
clientIP := r.RemoteAddr
// 处理反向代理真实IP(生产环境必加)
if realIP := r.Header.Get("X-Real-IP"); realIP != "" {
clientIP = realIP
} else if xff := r.Header.Get("X-Forwarded-For"); xff != "" {
clientIP = xff
}
// 4. 创建客户端实例(传入ClientID)
client := NewClient(conn, clientIP, clientID)
// 5. 本地连接池操作
clientPool.Store(client.ClientID, client) // 关键:改用ClientID作为Key,方便本地查询
atomic.AddUint64(&connCount, 1)
fmt.Printf("客户端[%s]握手成功,clientId: %s,IP: %s,节点:%s,当前活跃连接数:%d\n", client.connID, client.ClientID, clientIP, NodeID, atomic.LoadUint64(&connCount))
// 6. 写入Redis:存储客户端-节点映射
redisKey := RedisConnKeyPrefix + client.ClientID
err = redisClient.Set(ctx, redisKey, NodeID, RedisExpire).Err()
if err != nil {
fmt.Printf("写入Redis映射失败: %v, clientId: %s\n", err, client.ClientID)
client.Close()
return
}
// 可选:将客户端ID加入节点的Redis Set集合
redisNodeKey := RedisNodeClientsKeyPrefix + NodeID + ":clients"
redisClient.SAdd(ctx, redisNodeKey, client.ClientID)
// 7. 启动心跳续期协程(刷新Redis过期时间)
go client.renewRedisExpire()
// 8. 启动读/写协程
go client.ReadLoop()
go client.WriteLoop()
}
// renewRedisExpire 客户端心跳续期:定时刷新Redis映射的过期时间
func (c *Client) renewRedisExpire() {
ticker := time.NewTicker(RenewInterval)
defer ticker.Stop()
for {
select {
case <-c.ctx.Done():
// 连接关闭,退出续期协程
return
case <-ticker.C:
// 刷新Redis过期时间
redisKey := RedisConnKeyPrefix + c.ClientID
err := redisClient.Expire(ctx, redisKey, RedisExpire).Err()
if err != nil {
fmt.Printf("Redis续期失败: %v, clientId: %s\n", err, c.ClientID)
// 续期失败,主动关闭连接
c.Close()
return
}
}
}
}
(3)修改Close方法:清理Redis映射
连接关闭时,删除Redis中的客户端-节点映射和节点-客户端集合中的记录:
// Close 优雅关闭客户端连接,新增Redis清理操作
func (c *Client) Close() {
defer func() {
if err := recover(); err != nil {
fmt.Printf("客户端[%s]关闭时panic: %v, clientId: %s\n", c.connID, err, c.ClientID)
}
}()
// 1. 取消上下文,终止所有协程
c.cancel()
// 2. 关闭写通道
close(c.sendChan)
// 3. 关闭WebSocket连接
c.conn.WriteMessage(websocket.CloseMessage, websocket.FormatCloseMessage(websocket.CloseNormalClosure, ""))
c.conn.Close()
// 4. 从本地连接池删除
clientPool.Delete(c.ClientID)
atomic.AddUint64(&connCount, ^uint64(0))
// 5. 清理Redis映射(核心)
redisKey := RedisConnKeyPrefix + c.ClientID
redisClient.Del(ctx, redisKey)
// 6. 可选:从节点的Redis Set集合中移除
redisNodeKey := RedisNodeClientsKeyPrefix + NodeID + ":clients"
redisClient.SRem(ctx, redisNodeKey, c.ClientID)
fmt.Printf("客户端[%s]已优雅关闭,clientId: %s,节点:%s,当前活跃连接数:%d\n", c.connID, c.ClientID, NodeID, atomic.LoadUint64(&connCount))
}
5. 实现跨节点推送核心方法
封装本地推送和跨节点推送方法,对外提供统一的PushMsg接口,业务层只需调用该接口,无需关心目标节点是本地还是其他节点:
// pushLocal 本地推送:从本地连接池获取客户端,完成消息推送
func pushLocal(crossMsg CrossNodeMsg) {
// 从本地连接池获取客户端
v, ok := clientPool.Load(crossMsg.ClientID)
if !ok {
fmt.Printf("本地连接池无该客户端,clientId: %s,来自节点:%s\n", crossMsg.ClientID, crossMsg.FromNode)
// 清理Redis无效映射
redisClient.Del(ctx, RedisConnKeyPrefix + crossMsg.ClientID)
return
}
client := v.(*Client)
// 推送消息(根据消息类型,这里示例用文本,可扩展二进制)
if err := client.Send(crossMsg.Msg); err != nil {
fmt.Printf("本地推送失败: %v, clientId: %s\n", err, crossMsg.ClientID)
client.Close()
}
}
// PushMsg 对外提供的统一推送接口:自动判断本地/跨节点,业务层直接调用
func PushMsg(clientID string, msg []byte, msgType int) error {
// 1. 从Redis查询目标节点ID
redisKey := RedisConnKeyPrefix + clientID
targetNode, err := redisClient.Get(ctx, redisKey).Result()
if err != nil {
return fmt.Errorf("Redis查询目标节点失败: %v, clientId: %s", err, clientID)
}
// 2. 判断目标节点是否为自身
if targetNode == NodeID {
// 本地推送
pushLocal(CrossNodeMsg{
ClientID: clientID,
Msg: msg,
MsgType: msgType,
FromNode: NodeID,
})
return nil
}
// 3. 跨节点推送:发布消息到目标节点的Redis专属频道
crossMsg := CrossNodeMsg{
ClientID: clientID,
Msg: msg,
MsgType: msgType,
FromNode: NodeID,
}
// 序列化为JSON
msgJson, err := json.Marshal(crossMsg)
if err != nil {
return fmt.Errorf("跨节点消息序列化失败: %v", err)
}
// 发布到Redis频道
targetChannel := RedisChannelPrefix + targetNode
err = redisClient.Publish(ctx, targetChannel, string(msgJson)).Err()
if err != nil {
return fmt.Errorf("Redis发布跨节点消息失败: %v, 目标节点:%s", err, targetNode)
}
fmt.Printf("跨节点推送消息成功,clientId: %s,从节点:%s 到节点:%s\n", clientID, NodeID, targetNode)
return nil
}
// 封装广播方法(集群级,可选)
func ClusterBroadcast(msg []byte, msgType int) {
// 1. 本地广播
clientPool.Range(func(k, v interface{}) bool {
client := v.(*Client)
client.Send(msg)
return true
})
// 2. 跨节点广播:需遍历所有节点,发布到对应频道(生产环境可通过Redis获取所有节点ID)
// 此处省略,可通过Redis的keys ws:node:*:clients获取所有节点ID
}
6. 主函数:启动WebSocket服务
保留原有主函数,仅做少量修改,保证服务正常启动:
// genConnID 生成内部连接ID(无需修改)
func genConnID() string {
var id uint64
connMu.Lock()
connID++
id = connID
connMu.Unlock()
return fmt.Sprintf("conn-%d", id)
}
var connMu sync.Mutex
func main() {
// 注册WebSocket路由
http.HandleFunc("/ws", WsHandler)
// 启动HTTP服务,设置TCP参数(SO_REUSEPORT)
addr := ":" + WsPort
ln, err := net.Listen("tcp", addr)
if err != nil {
panic(fmt.Sprintf("监听端口失败: %v, 地址:%s", err, addr))
}
// 设置SO_REUSEPORT,提升多核接收效率
if tcpLn, ok := ln.(*net.TCPListener); ok {
tcpLn.SetSyscallConn(func(fd uintptr) {
syscall.SetsockoptInt(int(fd), syscall.SOL_SOCKET, syscall.SO_REUSEPORT, 1)
})
}
// 启动服务
fmt.Printf("WebSocket集群节点启动成功,节点ID:%s,监听地址:%s\n", NodeID, addr)
server := &http.Server{
ReadTimeout: 10 * time.Second,
}
if err := server.Serve(ln); err != nil && err != http.ErrServerClosed {
panic(fmt.Sprintf("服务启动失败: %v", err))
}
}
// 保留原有ReadLoop、WriteLoop、Send方法,无需修改
func (c *Client) ReadLoop() {
// 原有逻辑:读消息、重置读超时、处理业务
defer c.Close()
c.conn.SetReadDeadline(time.Now().Add(60 * time.Second))
c.conn.SetPongHandler(func(string) error {
c.conn.SetReadDeadline(time.Now().Add(60 * time.Second))
return nil
})
for {
select {
case <-c.ctx.Done():
return
default:
_, p, err := c.conn.ReadMessage()
if err != nil {
if !websocket.IsCloseError(err, websocket.CloseNormalClosure, websocket.CloseGoingAway) {
fmt.Printf("客户端读消息异常: %v, clientId: %s\n", err, c.ClientID)
}
return
}
fmt.Printf("收到客户端消息, clientId: %s, 内容:%s\n", c.ClientID, string(p))
// 业务层处理消息,可调用PushMsg实现跨节点回复
// 示例:回复客户端消息
_ = c.Send([]byte("收到消息:" + string(p)))
}
}
}
func (c *Client) WriteLoop() {
defer c.Close()
heartbeatTicker := time.NewTicker(30 * time.Second)
defer heartbeatTicker.Stop()
for {
select {
case <-c.ctx.Done():
return
case <-heartbeatTicker.C:
c.conn.SetWriteDeadline(time.Now().Add(10 * time.Second))
if err := c.conn.WriteMessage(websocket.PingMessage, nil); err != nil {
fmt.Printf("发送心跳失败: %v, clientId: %s\n", err, c.ClientID)
return
}
case msg, ok := <-c.sendChan:
if !ok {
c.conn.WriteMessage(websocket.CloseMessage, []byte{})
return
}
c.conn.SetWriteDeadline(time.Now().Add(10 * time.Second))
if err := c.conn.WriteMessage(websocket.TextMessage, msg); err != nil {
fmt.Printf("推送消息失败: %v, clientId: %s\n", err, c.ClientID)
return
}
}
}
}
func (c *Client) Send(msg []byte) error {
select {
case c.sendChan <- msg:
return nil
default:
return fmt.Errorf("写通道已满,clientId: %s", c.ClientID)
}
}
7. 测试跨节点推送
(1)启动集群节点
节点1:部署在192.168.1.100,启动服务,NodeID为
node-192.168.1.100:8080;节点2:部署在192.168.1.101,启动服务,NodeID为
node-192.168.1.101:8080;两个节点连接同一个Redis实例/集群。
(2)客户端连接
用WebSocket测试工具(如wscat、Postman)连接节点1,携带clientId:
# 安装wscat
npm install -g wscat
# 连接节点1的WebSocket服务
wscat -c ws://192.168.1.100:8080/ws?clientId=user123
(3)跨节点推送
在节点2的代码中调用PushMsg方法,给user123推送消息:
// 节点2的测试代码
func main() {
// 等待服务启动
time.Sleep(5 * time.Second)
// 跨节点推送消息给user123
err := PushMsg("user123", []byte("这是节点2的跨节点推送消息"), websocket.TextMessage)
if err != nil {
fmt.Printf("跨节点推送失败: %v\n", err)
} else {
fmt.Println("跨节点推送请求已发送")
}
// 阻塞程序
select {}
}
结果:连接在节点1的user123客户端会收到节点2推送的消息,实现跨节点通信。
五、核心优化方案(生产级必做)
上述基础方案可实现跨节点推送,要支撑百万级客户端的集群场景,需做以下优化,提升性能、可靠性和可维护性。
1. Redis层优化
使用Redis集群/哨兵:避免Redis单机单点故障,生产环境必须做高可用(Redis Cluster/主从+哨兵);
设置Redis连接池大小:根据节点并发量调大连接池(如200),避免Redis连接耗尽;
批量操作优化:节点下线时,用
SMEMBERS获取节点所有客户端,批量删除Redis映射,减少Redis调用次数;Redis键值过期监控:通过Redis的
KEYSPACE_NOTIFY监控过期键,及时清理无效的节点-客户端集合。
2. 跨节点推送优化
消息序列化优化:用
easyjson/msgpack替代原生encoding/json,提升序列化/反序列化效率(百万级推送时效果明显);非阻塞发布:Redis发布操作封装为非阻塞,避免Redis阻塞导致服务端卡死;
消息分片:超大消息(如>10KB)做分片发送,避免Redis频道传输过大消息导致的网络阻塞;
回执机制:核心消息推送后,要求目标节点返回回执(如通过另一个Redis频道),实现消息可靠性校验。
3. 节点管理优化
节点注册与发现:所有节点启动时,将NodeID和地址写入Redis Set(
ws:nodes),节点下线时删除,实现动态节点发现;节点健康检测:启动定时任务,检测Redis中
ws:nodes的节点是否存活(如TCP探活),标记异常节点并通知运维;客户端迁移:异常节点下线时,将其客户端列表推送给其他健康节点,实现客户端自动重连(需客户端配合实现重连逻辑)。
4. 客户端重连优化
客户端重连机制:客户端断开后,实现指数退避重连(1s→2s→4s),避免频繁重连导致集群压力;
重连节点选择:客户端重连时,从Redis的
ws:nodes中随机选择健康节点,实现负载均衡;连接状态缓存:客户端本地缓存最近连接的节点ID,重连时优先尝试该节点,提升重连成功率。
六、生产级扩展(解决核心痛点)
1. 替代Redis Pub/Sub:用Redis Stream实现可靠消息推送
Redis Pub/Sub的缺点是无消息持久化、消费者下线会丢消息,生产环境若需要可靠的跨节点推送,可用Redis Stream替代Pub/Sub,实现:
消息持久化:消息写入Stream后持久化,消费者下线后重启可继续消费;
消费组:每个节点作为一个消费组,避免消息重复消费;
消息确认:消费者消费消息后需手动ACK,确保消息被处理。
2. 节点间通信:HTTP/gRPC直接调用
对于核心业务消息(如IM聊天、金融指令),推荐用HTTP/gRPC实现节点间直接通信,替代Redis:
每个节点暴露推送API(如
POST /api/push,参数为clientId、msg、msgType);节点推送时,从Redis查询目标节点的IP+端口,直接调用该API完成推送;
实现重试机制(如3次重试)和超时控制,保证消息可靠送达。
3. 分布式限流与熔断
集群环境下需做分布式限流,避免单节点被压垮:
基于Redis实现令牌桶限流:限制每个客户端的消息发送频率(如10条/秒);
节点级熔断:当节点连接数达到阈值(如90万),拒绝新的握手请求,返回忙信号,引导客户端重连其他节点;
跨节点熔断:当某节点推送失败率过高(如>50%),暂时停止向该节点推送消息,避免无效调用。
4. 监控与告警
生产环境必须做全链路监控,及时发现问题:
节点监控:监控每个节点的连接数、CPU/内存/网络、Redis调用量;
Redis监控:监控Redis的QPS、内存、连接数、键值数量;
推送监控:监控跨节点推送的成功率、延迟、失败率;
告警机制:连接数骤增/骤减、推送失败率过高、节点下线时,通过邮件/钉钉/企业微信告警。
七、总结
通过Redis实现WebSocket集群跨节点推送的核心是「映射存储+消息路由」,方案轻量、高效,是业界落地成本最低的分布式WebSocket解决方案,核心要点总结:
Redis存储核心:用String结构存储
客户端ID→节点ID的映射,设置过期时间+心跳续期,避免无效数据;跨节点推送核心:用Redis Pub/Sub实现轻量消息路由,每个节点订阅专属频道,避免消息广播风暴;
统一推送接口:封装
PushMsg方法,自动判断本地/跨节点推送,业务层无需关心底层实现;节点管理核心:实现动态节点注册与发现,结合健康检测,保证集群稳定性;
可靠性优化:核心消息用Redis Stream/HTTP/gRPC替代Pub/Sub,实现消息持久化和回执机制。
该方案可直接结合之前的Golang WebSocket百万连接框架落地,支撑百万级客户端的分布式WebSocket集群,适配IM、实时推送、物联网设备连接等主流场景。
(注:文档部分内容可能由 AI 生成)
