基础知识
TCP 三次握手
客户端 - 服务器
三次握手目的:客户端的收发没问题,服务器的收发也没问题。
第一次握手:客户端发 SYN(同步请求)
- 客户端:我啥也不知道,我发个 SYN 消息出去,看看有没有人接收到。
- 如果没人回复,我再发几次(超时重传)。
- 这个时刻,客户端知道的事实是:
- 我发出去了,但不知道有没有人收到 ❌
- 我能不能收消息?不知道 ❌
- 对方存不存在?不知道 ❌
第二次握手:服务器回 SYN+ACK(同步+确认)
-
服务器收到了客户端的 SYN。
-
服务器回复 SYN+ACK:“我收到你的消息了,我也发一条试试,看看你能不能收到。”
-
这个时刻,服务器知道的事实是:
- 客户端发消息是没有问题的 ✅(因为我收到了)
- 服务器收消息是没有问题的 ✅(因为我收到了)
- 服务器不知道服务器发消息有没有问题 ❌(我发出去了,但不知道对方收到没)
- 服务器不知道客户端收消息有没有问题 ❌(我发出去了,但不知道对方收到没)
- 服务器需要发送一条 SYN+ACK,等对方回复 ACK 来确认“我能发、对方能收”
-
当客户端收到 SYN+ACK 后,客户端可以确定以下事实:
- 我发消息是没有问题的 ✅(我发的 SYN 被收到了)
- 我收消息也是没有问题的 ✅(我收到了对方的回复)
- 服务器收消息是没有问题的 ✅(对方收到了我的 SYN)
- 服务器发消息也是没有问题的 ✅(我收到了对方的 SYN+ACK)
- ✅ 客户端这边已经确认全双工通信没问题了!
- 但是服务器他还有顾虑(不知道自己能不能发、不知道我能不能收),我再发个 ACK 给他吧
第三次握手:客户端发 ACK(确认)
-
客户端收到 SYN+ACK 后,回复一个 ACK:“我收到你的 SYN+ACK 了,我这边收发都没问题,你也可以放心了。”
-
当服务器收到这个 ACK 后,服务器可以确定以下事实:
- 客户端发消息是没有问题的 ✅(第一次握手已验证)
- 服务器收消息是没有问题的 ✅(第一次握手已验证)
- 服务器发消息也是没有问题的 ✅(我发的 SYN+ACK 被客户端收到了)
- 客户端收消息也是没有问题的 ✅(我发的 SYN+ACK 被客户端收到了,他能收到就说明他能收)
- ✅ 服务器这边也确认全双工通信没问题了!
-
随后双方进入
ESTABLISHED状态,可以开始传输正式数据。
⚠️ 如果第三次握手的 ACK 丢了怎么办?
- 客户端视角:我发了 ACK,也直接发了正式数据(因为客户端认为连接已建立)。
- 服务器视角:我没收到 ACK,我还处于
SYN-RCVD状态,我不确定自己能不能发、客户端能不能收。 - 服务器行为:超时后重传 SYN+ACK(有限次数)。
- 客户端收到重传的 SYN+ACK 后:
- 客户端意识到“哦,我上次的 ACK 可能丢了”
- 客户端再次发送 ACK,并且捎带重传之前发的正式数据(如果有的话)
- 如果正式数据先到达服务器:
- 服务器的重传计时器还没超时,收到了客户端发来的正式数据包
- 这个数据包的 TCP 头部携带了 ACK 标志,正好确认了第二次握手
- 服务器立刻将连接状态转为
ESTABLISHED,收下数据
- 如果双方都超时了:
- 服务器有限次数重传 SYN+ACK 后无回应 → 放弃连接,释放资源
- 客户端有限次数重传正式数据后无回应 → 放弃连接,通知上层应用
- 双方都回到
CLOSED状态
🤔 为什么不能只握手两次?
如果只有两次握手(客户端发 SYN,服务器回 SYN+ACK 就结束):
| 角色 | 知道的事实 | 不知道的事实 |
|---|---|---|
| 客户端 | 我发没问题 ✅ 我收没问题 ✅ 服务器发没问题 ✅ 服务器收没问题 ✅ |
全部确认 ✅ |
| 服务器 | 客户端发没问题 ✅ 服务器收没问题 ✅ |
服务器发没问题 ❌ 客户端收没问题 ❌ |
后果:服务器在不确定自己能不能发、客户端能不能收的情况下就开始发送正式数据。如果服务器的 SYN+ACK 实际上客户端没收到,服务器还在傻傻地发数据,客户端根本收不到,导致死锁。
三次握手的本质:第三次 ACK 是给服务器吃的定心丸——让服务器确认“我能发,对方能收”,然后才敢放心发数据。
🤔 三次握手能携带数据吗?
-
第一次握手(SYN):不能携带数据。
- 因为此时连接还没建立,客户端还不知道服务器存不存在
- 如果携带数据,万一 SYN 丢了,数据白发了,且容易引发 SYN 攻击
-
第二次握手(SYN+ACK):不能携带数据。
- 服务器还没收到客户端的最终确认,不确定自己能不能发、客户端能不能收
- 如果带数据,万一 ACK 丢了,数据白发了
-
第三次握手(ACK):可以携带数据(RFC 793 允许)。
- 客户端已经确认了双方收发能力,可以放心捎带正式数据
- 如果 ACK 丢了,服务器会重传 SYN+ACK,客户端收到后会重传 ACK 并捎带重传数据
- 如果 ACK 没丢,服务器收到 ACK 后直接收下数据,一举两得
这就是 TCP 的“捎带确认”(Piggybacking)机制:利用第三次握手顺便传输数据,提高效率。
📊 三次握手状态变迁图
| 步骤 | 操作 | 客户端状态 | 服务器状态 |
|---|---|---|---|
| 初始 | 无 | CLOSED |
LISTEN |
| ① | 客户端发 SYN | SYN_SENT |
LISTEN |
| ② | 服务器收到 SYN,回 SYN+ACK | SYN_SENT |
SYN_RCVD |
| ②→③ | 客户端收到 SYN+ACK | ESTABLISHED |
SYN_RCVD |
| ③ | 客户端发 ACK | ESTABLISHED |
SYN_RCVD |
| ③→ | 服务器收到 ACK | ESTABLISHED |
ESTABLISHED |
💡 一句话总结三次握手
客户端先探路(SYN)→ 服务器回应并反问(SYN+ACK)→ 客户端最后确认(ACK),双方都确认“我能发、我能收、对方能发、对方能收”后,才开始正式通信。第三次 ACK 是给服务器的定心丸,确保服务器不会在“心里没底”的情况下发送数据。
TCP 四次挥手
客户端 - 服务器
四次挥手目的:客户端和服务器各自宣布“我没话说了”,并确认对方也“没话说了”,确保双方都完成数据收发后再断开连接,避免数据丢失。
第一次挥手:客户端发 FIN(结束请求)
- 客户端:我的数据都发送完了,我不想再发数据了,但我还能收数据。
- 我发一个 FIN 消息给服务器,告诉服务器:“我这边没话说了,我要准备关闭我的发送通道了。”
- 如果没人回复,我再发几次(超时重传)。
- 这个时刻,客户端知道的事实是:
- 我发了 FIN,但不知道服务器收到没有 ❌
- 服务器还有没有数据要发?不知道 ❌
- 我还能收数据,但不知道还能收多久 ❌
第二次挥手:服务器回 ACK(确认收到)
-
服务器收到了客户端的 FIN,知道客户端不再发数据了。
-
服务器回复一个 ACK:“收到你的 FIN,我知道你没话说了。”
-
这个时刻,服务器知道的事实是:
- 客户端不再发新数据了 ✅
- 服务器收消息是没问题的 ✅(收到了 FIN)
- 服务器还能继续发消息(因为服务器可能还有数据没发完)📦
- 客户端还能继续收消息 ✅(收到了 ACK 就说明客户端还在听)
- 服务器还没有发 FIN,因为数据没发完
-
当客户端收到 ACK 后,客户端可以确定以下事实:
- 我发的 FIN 服务器收到了 ✅(我发消息没问题)
- 服务器知道我不再发数据了 ✅
- 客户端不确定服务器还有没有数据要发 ❌
- 客户端不确定服务器什么时候会关闭 ❌
- 客户端还开着接收窗口,等着可能到来的数据 👂
- 进入**半关闭(Half-close)**状态:客户端不再发数据,但还能收数据
这个 ACK 的作用:只是告诉客户端“我知道你关了”,但服务器不会立即关闭自己的发送通道,因为可能还有数据要发给客户端。
第三次挥手:服务器发 FIN(结束请求)
-
服务器:我的数据也全部发送完毕了,我也没话说了。
-
服务器发一个 FIN 给客户端:“我的数据也发完了,我也要关闭我的发送通道了。”
-
这个时刻,服务器知道的事实是:
- 我发了 FIN,但不知道客户端收到没有 ❌
- 我不知道客户端会不会给我回 ACK ❌
- 如果客户端不回 ACK,我不知道该不该彻底关闭 ❌
- 我需要等待客户端的 ACK 来确认我的 FIN 已送达 ✅
-
当客户端收到 FIN 后,客户端可以确定以下事实:
- 服务器发 FIN 了,说明服务器确实没话说了 ✅
- 客户端之前担心“服务器还有没有数据”的顾虑消除了 ✅
- 客户端需要回复一个 ACK,让服务器安心关闭 ✅
- 但客户端还不能立即关闭,因为要等 ACK 送达确认
第四次挥手:客户端回 ACK(最终确认)
-
客户端收到服务器的 FIN,回复一个 ACK:“收到你的 FIN,确认关闭。”
-
当服务器收到 ACK 后,服务器可以确定以下事实:
- 我发的 FIN 客户端收到了 ✅
- 客户端知道我要关闭了 ✅
- 服务器可以安心关闭连接了 ✅
- 服务器释放资源,进入
CLOSED状态 🗑️
-
但是客户端收到 FIN 后,不会立即关闭:
- 客户端发完 ACK 后,进入 TIME_WAIT 状态 ⏳
- 客户端需要等待 2MSL(最大报文生存时间,通常是 2 分钟)
- 为什么要等?
- 如果最后一个 ACK 丢了,服务器会超时重传 FIN
- 客户端在 TIME_WAIT 期间如果收到服务器重传的 FIN,可以再发一次 ACK
- 确保服务器能收到最终的 ACK,避免服务器一直傻等 ❌
- 保证网络中所有旧数据包都消失,不会干扰后续的新连接(同一个端口复用)
- 等待结束后,客户端才真正关闭连接,进入
CLOSED状态
⚠️ 如果第四次挥手的 ACK 丢了怎么办?
- 客户端视角:我发了 ACK,进入 TIME_WAIT,等着服务器可能重传的 FIN。
- 服务器视角:我发了 FIN,怎么没收到 ACK?是不是丢了?
- 服务器行为:超时后重传 FIN(有限次数)。
- 客户端在 TIME_WAIT 期间:
- 如果收到了服务器重传的 FIN → 我上次的 ACK 丢了
- 客户端再次发送 ACK,并重置 TIME_WAIT 计时器(重新开始 2MSL 倒计时)
- 如果服务器多次重传 FIN 都收不到 ACK:
- 服务器放弃等待,强制关闭连接(进入
CLOSED),释放资源
- 服务器放弃等待,强制关闭连接(进入
- 如果客户端 TIME_WAIT 结束还没收到重传的 FIN:
- 客户端认为服务器已经收到了 ACK,正常关闭,进入
CLOSED
- 客户端认为服务器已经收到了 ACK,正常关闭,进入
这就是 TIME_WAIT 的意义:客户端在发完最后一个 ACK 后不能立即关闭,必须等一段时间,确保这个 ACK 万一丢了,还能重发,让服务器安心关闭。
🤔 为什么挥手是四次,握手是三次?
| 三次握手 | 四次挥手 | |
|---|---|---|
| 能否合并 | 服务器的 SYN 和 ACK 可以合并成 SYN+ACK | 服务器的 ACK 和 FIN 不能合并 |
| 为什么 | 服务器收到 SYN 时,不需要等待任何东西,立刻就知道“我也要建立连接”,所以可以把 SYN 和 ACK 合成一条发 | 服务器收到 FIN 时,可能还有数据没发完,不能立即关闭。必须先回 ACK 说“我知道了”,等数据全部发完后再单独发 FIN |
| 本质区别 | 建立连接是双方同步,没有先后顺序 | 关闭连接是各自独立,谁先关谁后关取决于谁先发完数据 |
实际场景比喻:
- 三次握手:两个人同时伸出手准备握手(SYN),同时回应对方(SYN+ACK),最后确认(ACK)。同步进行。
- 四次挥手:A 说“我说完了”(FIN),B 说“知道了”(ACK),但 B 还在继续说,等 B 说完了(FIN),A 说“知道了”(ACK)。各自独立。
📊 四次挥手状态变迁图
| 步骤 | 操作 | 客户端状态 | 服务器状态 |
|---|---|---|---|
| 初始 | 数据传输中 | ESTABLISHED |
ESTABLISHED |
| ① | 客户端发 FIN | FIN_WAIT_1 |
ESTABLISHED |
| ② | 服务器收到 FIN,回 ACK | FIN_WAIT_2 |
CLOSE_WAIT |
| ②→③ | 客户端收到 ACK | FIN_WAIT_2 |
CLOSE_WAIT |
| ③ | 服务器发 FIN(数据发完后) | FIN_WAIT_2 |
LAST_ACK |
| ③→④ | 客户端收到 FIN | TIME_WAIT |
LAST_ACK |
| ④ | 客户端发 ACK | TIME_WAIT |
LAST_ACK |
| ④→ | 服务器收到 ACK | TIME_WAIT |
CLOSED |
| 2MSL 后 | 客户端等待结束 | CLOSED |
CLOSED |
💡 一句话总结四次挥手
客户端先关(发 FIN)→ 服务器回 ACK(同意,但还在发数据)→ 服务器后关(发 FIN)→ 客户端回 ACK(最终确认,并傻等一会儿)
核心目的:
- 确保双方都没话说了(数据全部送达)
- 确保最后一个 ACK 不会丢,避免服务器“悬着一颗心”不知道对方收没收到
- 确保网络中旧数据包消失,不影响后续新连接
- 半关闭(Half-close)机制:允许一方关闭发送通道后,继续接收另一方数据,直到另一方也关闭。
WebSocket 交互过程
1. 基础 WebSocket 交互流程(含握手)
sequenceDiagram
participant Client as 客户端
participant Server as WebSocket服务端
Note over Client,Server: 1. HTTP握手升级
Client->>Server: HTTP GET /ws
Upgrade: websocket
Sec-WebSocket-Key: xxx
Server-->>Client: HTTP 101 Switching Protocols
Upgrade: websocket
Sec-WebSocket-Accept: yyy
Note over Client,Server: WebSocket连接已建立(全双工)
Client->>Server: 发送文本消息
Server-->>Client: 回复消息
Server->>Client: 主动推送消息
Client-->>Server: 确认收到(业务层)
Note over Client,Server: 2. 关闭连接
Client->>Server: Close Frame (状态码: 1000)
Server-->>Client: Close Frame (确认关闭)
Note over Client,Server: 连接已关闭
2. 详细交互流程(包含心跳 Ping/Pong)
sequenceDiagram
participant Client as 客户端
participant Server as WebSocket服务端
participant App as 业务应用层
Note over Client,App: 阶段1: 握手建立
Client->>Server: HTTP请求(Upgrade: websocket)
Server-->>Client: 101 Switching Protocols
Server->>App: 触发 @OnOpen 事件
Note over App: 保存Session到内存
Note over Client,App: 阶段2: 双向通信
Client->>Server: 发送业务消息(如:查询订单)
Server->>App: 触发 @OnMessage
App-->>Server: 处理业务逻辑
Server-->>Client: 推送响应数据
Server->>Client: 主动推送(如:系统通知)
Client-->>Server: 收到确认(业务层ACK)
Note over Client,App: 阶段3: 心跳保活(WebSocket协议层)
loop 每30秒
Server->>Client: Ping Frame (协议层)
Client-->>Server: Pong Frame (协议层)
Note over Server: 更新最后心跳时间
end
Note over Client,App: 阶段4: 异常断开或主动关闭
alt 主动关闭
Client->>Server: Close Frame (1000 Normal)
Server-->>Client: Close Frame (确认)
Server->>App: 触发 @OnClose
else 异常断开(网络超时)
Note over Server: 心跳超时检测到
Server->>App: 触发 @OnError
Server->>App: 触发 @OnClose
end
Note over App: 清理Session引用
3. 分布式场景下的 WebSocket(多节点 + Redis)
sequenceDiagram
participant Client1 as 客户端1
(连接Node-A)
participant Client2 as 客户端2
(连接Node-B)
participant NodeA as Node-A
(WebSocket服务器)
participant NodeB as Node-B
(WebSocket服务器)
participant Redis as Redis Pub/Sub
Note over Client1,NodeA: 客户端1连接到Node-A
Client1->>NodeA: WebSocket握手
NodeA-->>Client1: 连接建立
Note over NodeA: 保存Session1到本地内存
Note over Client2,NodeB: 客户端2连接到Node-B
Client2->>NodeB: WebSocket握手
NodeB-->>Client2: 连接建立
Note over NodeB: 保存Session2到本地内存
Note over Client1,Redis: 场景:Client1给Client2发消息
Client1->>NodeA: 发送消息(目标: Client2)
NodeA->>Redis: 发布消息到频道
(目标用户: Client2)
Redis->>NodeB: 广播消息到所有订阅节点
NodeB->>NodeB: 查找本地Session2
NodeB-->>Client2: 推送消息到Client2
Note over Client1,Client2: 转发完成
4. 完整的生命周期状态图
stateDiagram-v2
[*] --> 连接中: 发起HTTP请求
(Upgrade: websocket)
连接中 --> 已连接: 收到101响应
(握手成功)
连接中 --> 握手失败: 非101响应
或超时
已连接 --> 消息交互: 发送/接收消息
消息交互 --> 已连接: 持续通信
已连接 --> 心跳检测: 定期发送Ping
心跳检测 --> 已连接: 收到Pong响应
心跳检测 --> 超时断开: 多次Ping无响应
已连接 --> 主动关闭: 发送Close Frame
主动关闭 --> 已关闭: 收到Close确认
超时断开 --> 已关闭: 触发OnClose
握手失败 --> [*]
已关闭 --> [*]
note right of 已连接
Session可用
isOpen() == true
end note
note right of 已关闭
Session失效
isOpen() == false
清理内存引用
end note
5. WebSocket 帧结构交互(协议层细节)
sequenceDiagram
participant Client as 客户端
participant Server as 服务端
Note over Client,Server: 数据帧交互(RFC 6455)
Client->>Server: 文本帧 (Opcode=0x1)
FIN=1, Payload="Hello"
Server-->>Client: 文本帧 (Opcode=0x1)
FIN=1, Payload="World"
Note over Client,Server: 分片传输(大消息)
Client->>Server: 文本帧 (Opcode=0x1, FIN=0)
片段1
Client->>Server: 连续帧 (Opcode=0x0, FIN=0)
片段2
Client->>Server: 结束帧 (Opcode=0x0, FIN=1)
片段3
Note over Client,Server: 控制帧(心跳)
Server->>Client: Ping帧 (Opcode=0x9)
Client-->>Server: Pong帧 (Opcode=0xA)
Note over Client,Server: 关闭连接
Client->>Server: 关闭帧 (Opcode=0x8)
状态码=1000, Reason="Bye"
Server-->>Client: 关闭帧 (Opcode=0x8)
状态码=1000
6. 异常处理的交互流程
sequenceDiagram
participant Client as 客户端
participant Server as 服务端
participant Session as Session对象
participant Cleaner as 清理器
Note over Client,Cleaner: 异常场景处理
rect
Note over Client,Server: 场景1: 网络突然断开
Client--xServer: 网络中断(未发送Close)
Note over Server: 心跳超时检测(60秒)
Server->>Session: 检测到超时
Session->>Cleaner: 触发@OnError
Cleaner->>Cleaner: 移除本地Session引用
Cleaner->>Cleaner: 关闭资源(线程池等)
Note over Cleaner: Session.isOpen() == false
end
rect
Note over Client,Server: 场景2: 业务异常
Client->>Server: 发送非法消息
Server->>Server: 解析异常
Server->>Session: 触发@OnError
Session-->>Client: 发送关闭帧(状态码: 1002)
Session->>Cleaner: 触发@OnClose
Cleaner->>Cleaner: 清理资源
end
7. 总结:WebSocket Session 在流程中的位置
graph TB
subgraph "连接建立阶段"
A[HTTP握手] --> B[创建Session实例]
B --> C[触发OnOpen]
C --> D[保存到内存Map]
end
subgraph "通信阶段"
E[接收消息] --> F[通过Session获取连接]
F --> G[触发OnMessage]
G --> H[业务处理]
H --> I[通过Session.send回复]
end
subgraph "连接关闭阶段"
J[收到Close帧或超时] --> K[触发OnClose]
K --> L[从Map移除Session]
L --> M[关闭底层Channel]
M --> N[Session失效]
end
D --> E
I --> J
Spring 中编写代码的流程
- @Configuration WebSocketConfig
- addEndpoint - 接收来自哪里的URL
- setApplicationDestinationPrefixes - 接收来自哪个路径的消息
- enableSimpleBroker - 开启哪些路径的广播
- JwtHandshakeInterceptor - handshake时验证用户的token,放ws的attributes
- CustomHandshakeHandler extends DefaultHandshakeHandler - handshake时将attributes中的用户信息存入 Principal
- @Component WebSocketEventListener - 处理连接、订阅、取消订阅、失去连接时的行为
- 注意失去连接不会自动取消所有订阅!
- @ControllerAdvice WebSocketExceptionHandleController - 处理异常
- @MessageExceptionHandler(Exception.class) 指定哪些Exception
- @SendToUser("/queue/errors") 指定错误消息广播给哪个路径,注意spring会加上 “/user”
- @Controller ChatController - 接收到消息如何处理
- @MessageMapping("/chat.send") - 不用写前缀,setApplicationDestinationPrefixes已有默认前缀
MQ
基于 WebSocket 的微服务聊天处理流程
一、核心设计原则
- 先入库,再分发 —— 消息可靠性的底线
- 离线判定只做一次 —— 在发布端判定,避免 N 个实例重复推送
- 本地内存是唯一的 Session 真相源 —— Redis 只是"路由表",不存 Session
- Redis 在线状态可能撒谎 —— 必须容忍脏数据(宕机残留),靠 TTL + 心跳兜底
二、连接建立阶段
1. 用户 A 发起 WS 连接
ws://gateway/chat?token=xxx
- 走 Gateway,握手阶段(HTTP Upgrade)就完成鉴权
- 鉴权失败直接拒绝握手,不要建立连接后再踢
2. 建立连接通道
// 握手拦截器
public boolean beforeHandshake(...) {
String token = extractToken(request);
Claims claims = jwtUtil.verify(token); // RSA 验签
attributes.put("userId", claims.getSubject());
return true;
}
3. 记录本地在线状态(JVM 内存)
// key: userId, value: WebSocketSession
ConcurrentHashMap<Long, WebSocketSession> LOCAL_SESSIONS;
只有本实例能操作自己的 Session,这是不可跨进程的资源。
4. 记录 Redis 全局在线状态
HSET online:user:{userId} instanceId "chat-svc-192.168.1.10:8081"
EXPIRE online:user:{userId} 90s ← 关键!必须有 TTL
为什么必须有 TTL:
- 实例被
kill -9时来不及清理,Redis 会残留"僵尸在线" - 客户端会周期性发心跳(如 30s),服务端收到后
EXPIRE续期 - 90s 内没心跳 → 自动过期 → 视为离线 → 走推送
多端登录场景: 用 Hash 存 deviceId -> instanceId,一个用户可对应多个实例。
三、消息发送阶段(核心链路)
1. 客户端发消息
{
"msgId": "客户端生成的UUID", // 幂等去重用
"conversationId": 12345, // 会话/群组ID
"content": "hello",
"type": "TEXT"
}
2. 服务端接收 → 持久化(第一优先级)
Message msg = buildMessage(dto);
msg.setSeq(seqGenerator.next(conversationId)); // 会话内递增序号
messageMapper.insert(msg); // 必须先入库!
为什么必须先入库:
- 入库成功 = 消息已被系统接受,后续所有环节失败都可以补救(客户端拉历史)
- 若先推送后入库,推送成功但入库失败 → 消息永久丢失,且不同用户看到的历史不一致
这里要加:
msgId唯一索引 → 客户端重试时天然幂等seq会话内单调递增 → 客户端可据此发现空洞并主动补拉
3. 查询群成员 + 在发布端完成离线判定
// Step 1: 取群成员(走 Redis 缓存,miss 再查 DB)
List<Long> members = groupService.getMembers(conversationId);
// Step 2: 批量查 Redis 在线状态(用 Pipeline,不要 N 次 RTT)
Map<Long, String> onlineMap = redisTemplate.executePipelined(...);
// Step 3: 分流
List<Long> onlineUsers = members ∩ onlineMap.keySet();
List<Long> offlineUsers = members - onlineMap.keySet();
// Step 4: 离线的 → 直接发 Kafka(此刻就决定,只发一次)
if (!offlineUsers.isEmpty()) {
kafkaTemplate.send("push-topic", new PushEvent(msg, offlineUsers));
}
// Step 5: 在线的 → Publish 到 Redis Channel
redisTemplate.convertAndSend("chat:broadcast",
new BroadcastEvent(msg, onlineUsers));
为什么离线判定必须在这里做? 如果放到订阅端做,3 个实例订阅同一个 Channel,每个实例都会发现"用户 B 不在线" → 发 3 条 Kafka → 用户收到 3 条重复推送。
发布端只有 1 个,判定天然只做一次。
四、消息投递阶段(订阅端)
每个实例都订阅 chat:broadcast,收到后:
public void onMessage(BroadcastEvent event) {
// 只做一件事:找自己内存里有的 Session
for (Long uid : event.getTargetUsers()) {
WebSocketSession session = LOCAL_SESSIONS.get(uid);
if (session != null && session.isOpen()) {
session.sendMessage(msg); // 命中,发送
}
// 不在本地 → 直接忽略,交给别的实例
}
}
四种情况对照表:
| Redis 状态 | 本地内存 | 处理方式 |
|---|---|---|
| 在线 | 有 Session | ✅ 直接推送 |
| 在线 | 无 Session | 🔇 忽略(别的实例会处理) |
| 离线 | 有 Session | ⚠️ 脏数据!Redis 过期但连接还在 → 补写 Redis 并推送 |
| 离线 | 无 Session | 📮 已在发布端走 Kafka,此处不重复处理 |
兜底做法:本地有 Session 就发,同时
EXPIRE续期。用户可能收到「WS 消息 + 推送」两份,客户端用msgId去重即可。宁可重复,不可丢失。
五、离线推送阶段(Kafka 消费者)
消费流程
@KafkaListener(topics = "push-topic")
public void handle(PushEvent event, Acknowledgment ack) {
try {
// 二次确认:Kafka 有延迟,用户可能已经上线了
if (redisTemplate.hasKey("online:user:" + uid)) {
return; // 已上线,WS 会送达,跳过推送
}
pushClient.send(uid, event.getMsg()); // APNs / FCM / 厂商通道
ack.acknowledge();
} catch (Exception e) {
throw e; // 交给错误处理器
}
}
失败处理三级火箭
@Bean
public DefaultErrorHandler errorHandler(KafkaTemplate<?,?> template) {
// 1. 重试:指数退避,避免打垮下游
ExponentialBackOff backOff = new ExponentialBackOff(1000L, 2.0);
backOff.setMaxElapsedTime(30_000L); // 最多重试 30s
// 2. 重试耗尽 → 进死信队列
DeadLetterPublishingRecoverer recoverer =
new DeadLetterPublishingRecoverer(template,
(r, e) -> new TopicPartition("push-topic.DLT", r.partition()));
DefaultErrorHandler handler = new DefaultErrorHandler(recoverer, backOff);
// 3. 不可重试异常直接进 DLT(如参数错误、用户不存在)
handler.addNotRetryableExceptions(IllegalArgumentException.class);
return handler;
}
分层策略:
| 异常类型 | 处理 |
|---|---|
| 网络抖动、下游 5xx | 重试(指数退避) |
| Token 失效、用户注销 | 不重试,直接 DLT |
| 重试耗尽 | DLT + 告警 |
| DLT 堆积 | 监控指标 → 报警 → 人工介入 |
DLT 的运维闭环:
- DLT 消息保留 7 天
- Prometheus 监控
kafka_consumergroup_lag{topic="push-topic.DLT"} - 提供管理后台,支持查看 + 手动重投
补充:推送失败其实不算致命。消息已经在库里了,用户下次打开 App 主动拉取即可看到。推送只是"提醒",不是"投递"。这个认知能让你的降级策略轻松很多。
六、连接断开阶段
正常断开(onClose)
public void afterConnectionClosed(WebSocketSession session, CloseStatus status) {
Long uid = getUserId(session);
LOCAL_SESSIONS.remove(uid);
// 关键:CAS 删除,只删自己写的那条
// 防止用户已在别的实例重连,把新记录误删了
String owner = redis.opsForHash().get("online:user:"+uid, deviceId);
if (INSTANCE_ID.equals(owner)) {
redis.opsForHash().delete("online:user:"+uid, deviceId);
}
}
CAS 校验
用户从 4G 切 WiFi → 在实例 B 重连成功 → 实例 A 才收到断开事件 → 若无脑
DEL,会把实例 B 刚写的记录删掉 → 用户明明在线却被判离线。
异常断开(进程被杀)
- 依赖 Redis TTL 自动过期(90s)
- 实例启动时可扫描并清理属于自己
instanceId的残留 key
心跳保活
客户端 --ping--> 服务端(每 30s)
服务端 --pong--> 客户端 + EXPIRE online:user:{uid} 90s
IdleTimeout 设 60s,超时未收到 ping 主动关闭连接。
七、完整流程图
flowchart TD
A[用户A: WS握手请求] --> B{Gateway 鉴权}
B -->|失败| B1[拒绝握手]
B -->|成功| C[建立 WS 连接]
C --> D[写入 JVM 本地
ConcurrentHashMap]
D --> E["写入 Redis
online:user:uid → instanceId
TTL 90s"]
E --> F[心跳循环
每30s续期TTL]
F --> G[用户A 发送消息]
G --> H["1️⃣ 消息入库 (必须最先)
msgId唯一索引 + seq递增"]
H --> I[查询群成员列表]
I --> J["Pipeline 批量查
Redis 在线状态"]
J --> K{发布端分流
判定只做一次}
K -->|离线成员| L["Kafka → push-topic"]
K -->|在线成员| M["Redis Publish
chat:broadcast"]
L --> L1{推送服务消费}
L1 --> L2{二次确认
是否已上线?}
L2 -->|已上线| L3[跳过,WS会送达]
L2 -->|仍离线| L4[调用 APNs/FCM]
L4 -->|成功| L5[ACK 提交位点]
L4 -->|失败| L6{可重试异常?}
L6 -->|是| L7[指数退避重试
1s→2s→4s...最长30s]
L6 -->|否| L8[直接进 DLT]
L7 -->|重试耗尽| L8
L8 --> L9[DLT告警 → 管理后台
人工介入/手动重投]
L3 --> L10[消息已入库
客户端可主动拉取]
M --> N[所有实例订阅收到]
N --> O{查本地 LOCAL_SESSIONS}
O -->|命中且Open| P[✅ session.send 推送]
O -->|未命中| Q[🔇 忽略
交给持有Session的实例]
O -->|命中但Redis已过期| R["⚠️ 脏数据兜底
补写Redis + 推送"]
P --> S[客户端按 msgId 去重]
R --> S
S --> T[用户断开连接]
T --> U{断开类型}
U -->|正常 onClose| V[移除本地 Session]
V --> W{"CAS校验
owner == 本实例?"}
W -->|是| X[删除 Redis 记录]
W -->|否| Y[跳过
用户已在别处重连]
U -->|进程被kill| Z[Redis TTL 90s
自动过期兜底]
style H fill:#ff6b6b,color:#fff
style K fill:#ff6b6b,color:#fff
style W fill:#ffa94d,color:#000
style R fill:#ffa94d,color:#000