基础知识

TCP 三次握手

客户端 - 服务器

三次握手目的:客户端的收发没问题,服务器的收发也没问题。


第一次握手:客户端发 SYN(同步请求)

  1. 客户端:我啥也不知道,我发个 SYN 消息出去,看看有没有人接收到。
  2. 如果没人回复,我再发几次(超时重传)。
  3. 这个时刻,客户端知道的事实是
    • 我发出去了,但不知道有没有人收到 ❌
    • 我能不能收消息?不知道 ❌
    • 对方存不存在?不知道 ❌

第二次握手:服务器回 SYN+ACK(同步+确认)

  1. 服务器收到了客户端的 SYN。

  2. 服务器回复 SYN+ACK:“我收到你的消息了,我也发一条试试,看看你能不能收到。”

  3. 这个时刻,服务器知道的事实是

    • 客户端发消息是没有问题的 ✅(因为我收到了)
    • 服务器收消息是没有问题的 ✅(因为我收到了)
    • 服务器不知道服务器发消息有没有问题 ❌(我发出去了,但不知道对方收到没)
    • 服务器不知道客户端收消息有没有问题 ❌(我发出去了,但不知道对方收到没)
    • 服务器需要发送一条 SYN+ACK,等对方回复 ACK 来确认“我能发、对方能收”
  4. 当客户端收到 SYN+ACK 后,客户端可以确定以下事实

    • 我发消息是没有问题的 ✅(我发的 SYN 被收到了)
    • 我收消息也是没有问题的 ✅(我收到了对方的回复)
    • 服务器收消息是没有问题的 ✅(对方收到了我的 SYN)
    • 服务器发消息也是没有问题的 ✅(我收到了对方的 SYN+ACK)
    • 客户端这边已经确认全双工通信没问题了!
    • 但是服务器他还有顾虑(不知道自己能不能发、不知道我能不能收),我再发个 ACK 给他吧

第三次握手:客户端发 ACK(确认)

  1. 客户端收到 SYN+ACK 后,回复一个 ACK:“我收到你的 SYN+ACK 了,我这边收发都没问题,你也可以放心了。”

  2. 当服务器收到这个 ACK 后,服务器可以确定以下事实

    • 客户端发消息是没有问题的 ✅(第一次握手已验证)
    • 服务器收消息是没有问题的 ✅(第一次握手已验证)
    • 服务器发消息也是没有问题的 ✅(我发的 SYN+ACK 被客户端收到了)
    • 客户端收消息也是没有问题的 ✅(我发的 SYN+ACK 被客户端收到了,他能收到就说明他能收)
    • 服务器这边也确认全双工通信没问题了!
  3. 随后双方进入 ESTABLISHED 状态,可以开始传输正式数据。


⚠️ 如果第三次握手的 ACK 丢了怎么办?

  1. 客户端视角:我发了 ACK,也直接发了正式数据(因为客户端认为连接已建立)。
  2. 服务器视角:我没收到 ACK,我还处于 SYN-RCVD 状态,我不确定自己能不能发、客户端能不能收。
  3. 服务器行为:超时后重传 SYN+ACK(有限次数)。
  4. 客户端收到重传的 SYN+ACK 后
    • 客户端意识到“哦,我上次的 ACK 可能丢了”
    • 客户端再次发送 ACK,并且捎带重传之前发的正式数据(如果有的话)
  5. 如果正式数据先到达服务器
    • 服务器的重传计时器还没超时,收到了客户端发来的正式数据包
    • 这个数据包的 TCP 头部携带了 ACK 标志,正好确认了第二次握手
    • 服务器立刻将连接状态转为 ESTABLISHED,收下数据
  6. 如果双方都超时了
    • 服务器有限次数重传 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(结束请求)

  1. 客户端:我的数据都发送完了,我不想再发数据了,但我还能收数据。
  2. 我发一个 FIN 消息给服务器,告诉服务器:“我这边没话说了,我要准备关闭我的发送通道了。”
  3. 如果没人回复,我再发几次(超时重传)。
  4. 这个时刻,客户端知道的事实是
    • 我发了 FIN,但不知道服务器收到没有 ❌
    • 服务器还有没有数据要发?不知道 ❌
    • 我还能收数据,但不知道还能收多久 ❌

第二次挥手:服务器回 ACK(确认收到)

  1. 服务器收到了客户端的 FIN,知道客户端不再发数据了。

  2. 服务器回复一个 ACK:“收到你的 FIN,我知道你没话说了。”

  3. 这个时刻,服务器知道的事实是

    • 客户端不再发新数据了 ✅
    • 服务器收消息是没问题的 ✅(收到了 FIN)
    • 服务器还能继续发消息(因为服务器可能还有数据没发完)📦
    • 客户端还能继续收消息 ✅(收到了 ACK 就说明客户端还在听)
    • 服务器还没有发 FIN,因为数据没发完
  4. 当客户端收到 ACK 后,客户端可以确定以下事实

    • 我发的 FIN 服务器收到了 ✅(我发消息没问题)
    • 服务器知道我不再发数据了 ✅
    • 客户端不确定服务器还有没有数据要发 ❌
    • 客户端不确定服务器什么时候会关闭 ❌
    • 客户端还开着接收窗口,等着可能到来的数据 👂
    • 进入**半关闭(Half-close)**状态:客户端不再发数据,但还能收数据

这个 ACK 的作用:只是告诉客户端“我知道你关了”,但服务器不会立即关闭自己的发送通道,因为可能还有数据要发给客户端。


第三次挥手:服务器发 FIN(结束请求)

  1. 服务器:我的数据也全部发送完毕了,我也没话说了。

  2. 服务器发一个 FIN 给客户端:“我的数据也发完了,我也要关闭我的发送通道了。”

  3. 这个时刻,服务器知道的事实是

    • 我发了 FIN,但不知道客户端收到没有 ❌
    • 我不知道客户端会不会给我回 ACK ❌
    • 如果客户端不回 ACK,我不知道该不该彻底关闭 ❌
    • 我需要等待客户端的 ACK 来确认我的 FIN 已送达 ✅
  4. 当客户端收到 FIN 后,客户端可以确定以下事实

    • 服务器发 FIN 了,说明服务器确实没话说了 ✅
    • 客户端之前担心“服务器还有没有数据”的顾虑消除了 ✅
    • 客户端需要回复一个 ACK,让服务器安心关闭 ✅
    • 但客户端还不能立即关闭,因为要等 ACK 送达确认

第四次挥手:客户端回 ACK(最终确认)

  1. 客户端收到服务器的 FIN,回复一个 ACK:“收到你的 FIN,确认关闭。”

  2. 当服务器收到 ACK 后,服务器可以确定以下事实

    • 我发的 FIN 客户端收到了 ✅
    • 客户端知道我要关闭了 ✅
    • 服务器可以安心关闭连接了 ✅
    • 服务器释放资源,进入 CLOSED 状态 🗑️
  3. 但是客户端收到 FIN 后,不会立即关闭

    • 客户端发完 ACK 后,进入 TIME_WAIT 状态 ⏳
    • 客户端需要等待 2MSL(最大报文生存时间,通常是 2 分钟)
    • 为什么要等?
      • 如果最后一个 ACK 丢了,服务器会超时重传 FIN
      • 客户端在 TIME_WAIT 期间如果收到服务器重传的 FIN,可以再发一次 ACK
      • 确保服务器能收到最终的 ACK,避免服务器一直傻等 ❌
      • 保证网络中所有旧数据包都消失,不会干扰后续的新连接(同一个端口复用)
    • 等待结束后,客户端才真正关闭连接,进入 CLOSED 状态

⚠️ 如果第四次挥手的 ACK 丢了怎么办?

  1. 客户端视角:我发了 ACK,进入 TIME_WAIT,等着服务器可能重传的 FIN。
  2. 服务器视角:我发了 FIN,怎么没收到 ACK?是不是丢了?
  3. 服务器行为:超时后重传 FIN(有限次数)。
  4. 客户端在 TIME_WAIT 期间
    • 如果收到了服务器重传的 FIN → 我上次的 ACK 丢了
    • 客户端再次发送 ACK,并重置 TIME_WAIT 计时器(重新开始 2MSL 倒计时)
  5. 如果服务器多次重传 FIN 都收不到 ACK
    • 服务器放弃等待,强制关闭连接(进入 CLOSED),释放资源
  6. 如果客户端 TIME_WAIT 结束还没收到重传的 FIN
    • 客户端认为服务器已经收到了 ACK,正常关闭,进入 CLOSED

这就是 TIME_WAIT 的意义:客户端在发完最后一个 ACK 后不能立即关闭,必须等一段时间,确保这个 ACK 万一丢了,还能重发,让服务器安心关闭。


🤔 为什么挥手是四次,握手是三次?

三次握手 四次挥手
能否合并 服务器的 SYNACK 可以合并成 SYN+ACK 服务器的 ACKFIN 不能合并
为什么 服务器收到 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(最终确认,并傻等一会儿)

核心目的

  1. 确保双方都没话说了(数据全部送达)
  2. 确保最后一个 ACK 不会丢,避免服务器“悬着一颗心”不知道对方收没收到
  3. 确保网络中旧数据包消失,不影响后续新连接
  4. 半关闭(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 中编写代码的流程

  1. @Configuration WebSocketConfig
    1. addEndpoint - 接收来自哪里的URL
    2. setApplicationDestinationPrefixes - 接收来自哪个路径的消息
    3. enableSimpleBroker - 开启哪些路径的广播
  2. JwtHandshakeInterceptor - handshake时验证用户的token,放ws的attributes
  3. CustomHandshakeHandler extends DefaultHandshakeHandler - handshake时将attributes中的用户信息存入 Principal
  4. @Component WebSocketEventListener - 处理连接、订阅、取消订阅、失去连接时的行为
    1. 注意失去连接不会自动取消所有订阅!
  5. @ControllerAdvice WebSocketExceptionHandleController - 处理异常
    1. @MessageExceptionHandler(Exception.class) 指定哪些Exception
    2. @SendToUser("/queue/errors") 指定错误消息广播给哪个路径,注意spring会加上 “/user”
  6. @Controller ChatController - 接收到消息如何处理
    1. @MessageMapping("/chat.send") - 不用写前缀,setApplicationDestinationPrefixes已有默认前缀

MQ

基于 WebSocket 的微服务聊天处理流程

一、核心设计原则

  1. 先入库,再分发 —— 消息可靠性的底线
  2. 离线判定只做一次 —— 在发布端判定,避免 N 个实例重复推送
  3. 本地内存是唯一的 Session 真相源 —— Redis 只是"路由表",不存 Session
  4. 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