文章目录

适用于:基于 Spring Boot + WebFlux + Redis + RabbitMQ + Redisson 的后端,在引入 WebSocket 做实时推送(点赞、关注、收藏、评论提醒、私信等)时的企业级落地实践。


1. 整体架构概览

1.1 架构角色
  • 客户端

    • 使用浏览器 / App WebSocket SDK 建立 wss://gateway/ws 长连接。
    • 负责心跳、断线重连、消息展示与本地缓存。
  • 应用节点(WebSocket Server)

    • 基于 Spring WebFlux + Netty 提供 WebSocket 服务。
    • 使用 Redis 维护分布式会话状态 / 路由信息。
    • 使用 RabbitMQ 接收业务侧生产的通知消息,再异步推送到用户 WebSocket 连接。
    • 使用 Redisson 分布式锁 保证跨节点并发更新会话状态的安全性。
  • 基础设施

    • Redis:WebSocket 会话路由表、连接信息、心跳时间、在线状态缓存、限流计数等。
    • RabbitMQ:点赞 / 评论 / 关注 / 收藏等业务产生的通知消息总线,解耦业务与 WebSocket 推送。
    • 数据库(MySQL/R2DBC):消息持久化(可选)、已读状态、历史消息查询等。
1.2 典型消息流
  1. 用户 A 对用户 B 的视频点赞。
  2. 点赞业务服务写数据库 / Redis 计数,并投递一条 LikeMessage 到 RabbitMQ。
  3. WebSocket 推送服务消费 MQ 中的 LikeMessage,根据 userId 到 Redis 查找 B 的所有活跃 sessionId 及所在节点。
  4. 如果当前节点持有对应连接,则通过 WebSocket 将通知推送给 B;否则(按设计)可以:
    • 简化模式:所有节点都订阅同一消息队列,节点只对自己持有的连接进行推送,其余直接忽略;
    • 进阶模式:根据 userId 做一致性哈希路由,将消息转发到对应节点专用队列再推送。

2. WebSocket 连接管理最佳实践

2.1 连接接入与线程模型
  • 使用 WebFlux + Netty 提供 WebSocket 服务:

    • 连接建立与收发均在 Netty EventLoop 线程上处理,必须避免阻塞操作(JDBC、同步 Redis 等)。
    • 所有阻塞调用包裹在 Schedulers.boundedElastic() 或自定义线程池中执行。
  • 核心建议

    • WebSocket Handler 内仅做:协议解析、简单路由、异步调用封装,不直接访问数据库 / 外部服务。
    • 推送消息尽量通过 RabbitMQ 异步化,避免在连接线程中做重 IO 写入。
2.2 连接池与资源配置
  • Netty/EventLoop 配置(示例,需压测后调优):

    server:
      netty:
        connection-idle-timeout: 30m   # 空闲连接超时关闭
        worker-count: 4                # 一般设为 CPU 核数或 2 * 核数
    spring:
      webflux:
        websocket:
          max-frame-payload-length: 65536  # 单帧最大长度,防止大包攻击
    
  • RabbitMQ 连接池

    • 使用 Spring AMQP 默认连接池即可,一般配置:

      spring:
        rabbitmq:
          host: xxx
          port: 5672
          username: xxx
          password: xxx
          listener:
            simple:
              concurrency: 4
              max-concurrency: 16
              prefetch: 100       # 每次预取消息数
      
  • Redis 连接池

    • 使用 Lettuce 连接池,针对高并发读写场景适当放大 max-active,同时限制 max-wait 防止排队过长。
2.3 心跳机制
  • 客户端 → 服务端心跳

    • 建议使用 应用层 Ping 消息,例如:

      { "type": "PING", "ts": 1710000000000 }
      
    • 服务端收到后更新 Redis 中 lastHeartbeatAt 字段,并返回 PONG

      { "type": "PONG", "ts": 1710000000001 }
      
  • 服务端主动检测

    • 定时任务扫描 Redis 中会话:now - lastHeartbeatAt > idleTimeout 则认为断开:
      • 关闭本地 WebSocketSession;
      • 删除 Redis 中该 session 的路由信息;
      • 记录日志与监控指标。
  • 时间间隔建议

    • 心跳间隔:10–30 秒;
    • 探活超时:心跳 3 次未收到则断开(即 30–90 秒)。
2.4 断线重连策略
  • 客户端策略

    • 使用 指数退避 + 随机抖动 重连:1s, 2s, 4s, 8s,上限 30s;
    • 重连时携带 上次会话标识 / 设备标识,便于服务端做状态恢复或幂等处理。
  • 服务端策略

    • 使用 userId + deviceId 作为连接唯一键;
    • 新连接建立时:
      • 获取 lock:ws:user:{userId}:device:{deviceId} 分布式锁;
      • 将旧连接标记为 replaced,关闭并清理 Redis 状态;
      • 写入新的 sessionId 与节点路由信息。
  • 与幂等组件集成

    • 对需要有“只推送一次”语义的业务消息,可复用项目中的幂等组件,对 messageId 做一次性消费控制。

3. 高并发下的 Redis 分布式会话存储

3.1 会话模型设计
  • 基础键设计(可按实际项目前缀调整):

    • ws:session:{sessionId}:Hash,记录单连接信息:

      • userId:用户 ID;
      • deviceId:设备标识;
      • nodeId:当前连接所在应用实例;
      • lastHeartbeatAt:最近一次心跳时间戳;
      • createdAt:连接建立时间;
      • clientType:APP / H5 / 小程序等。
    • ws:user:{userId}:Set,存储该用户所有活跃 sessionId

    • ws:node:{nodeId}:Set,存储该节点管理的所有 sessionId,节点下线时用于批量回收。

  • TTL 策略

    • ws:session:{sessionId} 设置 TTL(如 1 小时),每次心跳可 expire 延长;
    • 节点下线或连接正常关闭时,主动删除相关 key。
3.2 会话写入流程
  1. WebSocket 握手通过 JWT 鉴权,解析出 userId/ deviceId
  2. 获取 lock:ws:user:{userId}:device:{deviceId} 分布式锁。
  3. 写入 ws:session:{sessionId} Hash;
  4. sessionId 加入 ws:user:{userId}ws:node:{nodeId}
  5. 释放锁。
3.3 会话查询与路由
  • 按用户查询所有连接

    public Mono<Set<String>> getSessionIdsByUserId(Long userId) {
        String key = "ws:user:" + userId;
        return redisTemplate.opsForSet().members(key).collectList()
                .map(HashSet::new);
    }
    
  • 按 sessionId 查询节点路由

    public Mono<String> getNodeIdBySessionId(String sessionId) {
        String key = "ws:session:" + sessionId;
        return redisTemplate.opsForHash().get(key, "nodeId")
                .cast(String.class);
    }
    
  • 在本地节点维护内存映射

    • ConcurrentHashMap<String, WebSocketSession> localSessions
    • 仅保存本节点的 session,减少 Redis 访问次数。
3.4 一致性与并发安全
  • 使用 Redisson Lock 保证 userId + deviceId 下的写入互斥;
  • 对异常退出(进程 Crash)场景:
    • 依赖 Redis TTL 自动清理;
    • ws:node:{nodeId} 可在节点启动时做一次修复扫描,清理不存在的 session。

4. 使用 RabbitMQ 实现 WebSocket 消息异步推送

4.1 设计目标
  • 解耦业务请求线程与 WebSocket 推送线程;
  • 避免在 WebSocket 连接线程上执行业务逻辑和慢查询;
  • 支持多节点订阅同一类业务消息,实现水平扩展。
4.2 交换机与队列设计
  • 建议为 WebSocket 推送建立独立交换机,例如:

    exchange: ws.notify.exchange   # Topic / Direct
    queue   : ws.notify.queue      # 或按业务拆分:like/comment/follow
    routingKey: ws.notify.*
    
  • 业务服务(点赞 / 评论 / 关注 / 收藏)在关键流程中发送 MQ 消息,例如:

    @Service
    public class LikeService {
    
        private final RabbitTemplate rabbitTemplate;
    
        public void likeVideo(LikeCommand cmd) {
            // 1. 处理数据库 & 计数
            // 2. 构建 LikeMessage
            LikeMessage message = new LikeMessage();
            message.setFromUserId(cmd.getUserId());
            message.setTargetUserId(cmd.getVideoOwnerId());
            message.setVideoId(cmd.getVideoId());
            message.setTs(System.currentTimeMillis());
    
            rabbitTemplate.convertAndSend(
                    "ws.notify.exchange",
                    "ws.notify.like",
                    message
            );
        }
    }
    
4.3 WebSocket 推送服务消费 MQ
@Service
public class WsNotificationConsumer {

    private final WebSocketPushService webSocketPushService;

    @RabbitListener(queues = "ws.notify.queue")
    public void onMessage(LikeMessage msg) {
        webSocketPushService.pushToUser(msg.getTargetUserId(), msg);
    }
}
  • WebSocketPushService.pushToUser 中:
    • 根据 userId 从 Redis 查出所有 sessionId
    • 在本地 localSessions 中查找对应 WebSocketSession
    • 找不到则忽略或记录为离线(可落库离线消息)。
4.4 避免阻塞连接
  • WebSocket 消息发送使用 异步写,例如 WebFlux:

    public Mono<Void> sendToSession(WebSocketSession session, Object payload) {
        String json = objectMapper.writeValueAsString(payload);
        return session.send(Mono.just(session.textMessage(json))); // 非阻塞
    }
    
  • 不在 @RabbitListener 方法内做复杂计算,可转为 Reactor 流:

    @RabbitListener(queues = "ws.notify.queue")
    public void onMessage(LikeMessage msg) {
        Mono.just(msg)
            .publishOn(Schedulers.boundedElastic())
            .flatMap(m -> webSocketPushService.pushToUser(m.getTargetUserId(), m))
            .subscribe();
    }
    

5. Redisson 分布式锁在 WebSocket 并发场景中的应用

5.1 典型并发问题
  • 同一用户在多个端 / 多节点同时发起连接 / 重连,可能导致:

    • 多个有效连接并存,业务期望只有一个;
    • Redis 中会话信息被覆盖 / 脏数据。
  • 同一用户在短时间内收到大量通知消息,推送逻辑需要保证 顺序性幂等性

5.2 Redisson 锁使用建议
  • 锁键设计:

    • lock:ws:user:{userId}:device:{deviceId} —— 连接绑定级别;
    • lock:ws:push:{userId} —— 消息推送顺序控制(如需要)。
  • 连接绑定示例:

    public Mono<Void> registerSession(Long userId, String deviceId, WebSocketSession session) {
        String lockKey = "lock:ws:user:" + userId + ":device:" + deviceId;
        RLock lock = redissonClient.getLock(lockKey);
    
        return Mono.fromCallable(() -> {
                    // tryLock 带超时,避免死锁
                    if (!lock.tryLock(3, 10, TimeUnit.SECONDS)) {
                        throw new IllegalStateException("acquire ws lock timeout");
                    }
                    return true;
                })
                .subscribeOn(Schedulers.boundedElastic())
                .flatMap(ignored -> doRegister(userId, deviceId, session))
                .doFinally(sig -> lock.unlockAsync());
    }
    
  • 推送顺序控制(可选):

    • 对需要严格时序的通知(如同一对话的聊天消息),可在消费 MQ 前后加锁,确保同一 userId 的推送按顺序执行。

6. WebSocket 消息格式与 Jackson 序列化最佳实践

6.1 统一消息 Envelope
  • 推荐统一的消息Envelope,支持扩展与版本控制:

    {
      "type": "LIKE",           // 消息类型:LIKE/FOLLOW/COMMENT/SYSTEM 等
      "version": 1,              // 协议版本
      "traceId": "...",         // 链路追踪ID
      "bizId": "like:123",      // 业务ID(便于幂等追踪)
      "ts": 1710000000000,       // 服务器时间戳
      "payload": {               // 具体业务负载
        "fromUserId": 1,
        "targetUserId": 2,
        "videoId": 10001
      }
    }
    
6.2 Java DTO 设计(Jackson)
@Data
public class WsEnvelope<T> {
    private String type;
    private Integer version;
    private String traceId;
    private String bizId;
    private Long ts;
    private T payload;
}

@Data
public class LikePayload {
    private Long fromUserId;
    private Long targetUserId;
    private Long videoId;
}
  • 统一使用 Spring 上下文中的 ObjectMapper(与现有 REST 接口保持一致),避免重新实例化配置不一致。
  • 开启常见配置:
    • WRITE_DATES_AS_TIMESTAMPS=false
    • FAIL_ON_UNKNOWN_PROPERTIES=false
    • 全局时间格式与 REST API 对齐。
6.3 文本 vs 二进制消息
  • 推荐优先使用文本 JSON 消息

    • 可直接在浏览器调试,配合 Chrome/DevTools;
    • 方便与现有 REST 接口复用 DTO。
  • 在极端高并发 / 大消息体场景,可考虑:

    • 使用 二进制 + JSON 压缩(如 gzip)
    • 需要客户端 SDK 配合处理。

7. 安全性:认证、鉴权、消息防篡改

7.1 连接认证
  • WebSocket 握手阶段校验用户身份:

    • 客户端在连接 URL 或 Header 中携带 JWT:wss://domain/ws?token=xxx
    • 服务端使用项目中的 JwtUtil 解析 Token,校验签名与过期时间;
    • 未通过认证直接拒绝握手。
  • 示例(WebFlux HandlerMapping 拦截):

    public class AuthWebSocketHandler implements WebSocketHandler {
    
        private final JwtUtil jwtUtil;
        private final WebSocketSessionManager sessionManager;
    
        @Override
        public Mono<Void> handle(WebSocketSession session) {
            URI uri = session.getHandshakeInfo().getUri();
            String token = UriComponentsBuilder.fromUri(uri).build()
                    .getQueryParams().getFirst("token");
    
            UserInfo user = jwtUtil.parseToken(token);
            if (user == null) {
                return session.close(CloseStatus.NOT_ACCEPTABLE);
            }
    
            // 注册会话
            return sessionManager.register(user, session)
                    .thenMany(handleMessages(user, session))
                    .then();
        }
    }
    
7.2 权限验证
  • 服务端对每一类消息进行 服务端侧的权限校验,不要信任客户端:
    • 比如客户端请求订阅某聊天室 / 私信会话,必须校验用户是否为会话参与者;
    • 推送通知时,确保消息内容只发送给有权限的用户。
7.3 消息防篡改与重放防护
  • 客户端 → 服务端 的重要操作消息:

    • 使用 JWT 中的 sub / userId 作为唯一身份,所有敏感字段以服务端为准;
    • 对可被重复提交的操作,结合 幂等 Token(可复用项目幂等组件)。
  • 服务端 → 客户端 的通知消息:

    • 完整性由 TLS(wss)保障即可,一般无需额外签名;
    • 若业务要求极高安全,可追加 HMAC 签名字段。
7.4 防刷与限流
  • 使用 Redis 计数实现 IP / userId 级别限流
    • 单用户每秒发送消息数、订阅频次等均需限制;
    • 达到阈值时可返回错误或临时封禁连接。

8. 性能优化:压缩、批量推送、连接复用

8.1 消息压缩
  • 在浏览器 / App 侧开启 permessage-deflate 扩展;
  • 服务端可通过配置启用 WebSocket 压缩(具体取决于底层 Netty/Spring 版本)。
8.2 批量推送与合并
  • 对高频通知(如点赞/浏览数变化),避免每一条都实时推送:
    • 使用 窗口 + 合并 策略:对同一用户一定时间窗口内的多条通知合并为一条汇总消息;

    • 示例:

      Flux<LikeMessage> likeStream = ...; // 来自 MQ
      
      likeStream
          .groupBy(LikeMessage::getTargetUserId)
          .flatMap(group -> group.bufferTimeout(50, Duration.ofSeconds(1)))
          .flatMap(batch -> webSocketPushService.pushBatch(batch))
          .subscribe();
      
8.3 连接复用与资源复用
  • 对同一用户在同一设备上,建议 限制为单连接,避免过多连接占用资源;
  • 所有外部资源(Redis、RabbitMQ、数据库)均通过 Spring 单例 Bean 复用连接池,不在 Handler 中新建客户端实例。
8.4 参数调优建议
  • 根据压测结果调整:
    • 最大并发连接数(估算:在线用户数 × 平均连接数 / 用户);
    • RabbitMQ prefetch 配置,平衡延迟与批量效率;
    • Redis 连接池 max-activemax-idle 与 CPU/内存资源。

9. 错误处理与监控

9.1 异常捕获与优雅关闭
  • 在 WebSocket Handler 中对异常集中处理:

    • 协议解析错误 → 返回错误消息或直接关闭连接;
    • 业务处理异常 → 返回标准错误结构,避免抛堆栈给客户端。
  • 统一封装错误消息格式:

    {
      "type": "ERROR",
      "code": "BIZ_XXX",
      "message": "错误描述",
      "ts": 1710000000000
    }
    
9.2 日志记录
  • 关键日志点:

    • 连接建立 / 关闭(含原因码);
    • 鉴权失败 / 权限拒绝;
    • 推送失败(session 不存在、网络错误);
    • MQ 消费失败与重试情况。
  • 严格控制日志级别:正常推送不打印 INFO 级别日志,避免高并发下 IO 放大;

  • 对错误日志附加 traceIduserIdsessionId,方便追踪。

9.3 指标与监控
  • 建议集成 Micrometer + Prometheus,暴露以下指标:

    • 当前在线连接数(按节点 / 按业务类型);
    • 每秒推送消息数、失败率;
    • MQ 消息堆积深度、消费延迟;
    • Redis 操作 RT 与错误率。
  • 关键告警:

    • 在线连接数异常下降 / 激增;
    • MQ 堆积持续升高;
    • 推送失败率超阈值。

10. 代码示例:配置类、服务类、Handler 实现

以下示例以 WebFlux 原生 WebSocketHandler 为例,简化展示核心结构,实际项目可根据现有包结构分层。

10.1 WebSocket 配置类
@Configuration
public class WebSocketConfig {

    private final AuthWebSocketHandler authWebSocketHandler;

    @Bean
    public HandlerMapping webSocketMapping() {
        Map<String, WebSocketHandler> map = new HashMap<>();
        map.put("/ws", authWebSocketHandler);

        SimpleUrlHandlerMapping mapping = new SimpleUrlHandlerMapping();
        mapping.setOrder(10);
        mapping.setUrlMap(map);
        return mapping;
    }

    @Bean
    public WebSocketHandlerAdapter handlerAdapter() {
        return new WebSocketHandlerAdapter();
    }
}
10.2 会话管理服务(Redis + 本地缓存)
@Service
public class WebSocketSessionManager {

    private final ReactiveStringRedisTemplate redisTemplate;
    private final RedissonClient redissonClient;
    private final Map<String, WebSocketSession> localSessions = new ConcurrentHashMap<>();
    private final String nodeId = UUID.randomUUID().toString();

    public Mono<Void> register(UserInfo user, WebSocketSession session) {
        String sessionId = session.getId();
        String deviceId = Optional.ofNullable(user.getDeviceId()).orElse("default");

        String lockKey = "lock:ws:user:" + user.getId() + ":device:" + deviceId;
        RLock lock = redissonClient.getLock(lockKey);

        return Mono.fromCallable(() -> {
                    if (!lock.tryLock(3, 10, TimeUnit.SECONDS)) {
                        throw new IllegalStateException("acquire ws lock timeout");
                    }
                    return true;
                })
                .subscribeOn(Schedulers.boundedElastic())
                .flatMap(ignored -> doRegister(user, deviceId, session))
                .doFinally(sig -> lock.unlockAsync());
    }

    private Mono<Void> doRegister(UserInfo user, String deviceId, WebSocketSession session) {
        String sessionKey = "ws:session:" + session.getId();
        String userKey = "ws:user:" + user.getId();
        String nodeKey = "ws:node:" + nodeId;

        Map<String, String> sessionInfo = new HashMap<>();
        sessionInfo.put("userId", String.valueOf(user.getId()));
        sessionInfo.put("deviceId", deviceId);
        sessionInfo.put("nodeId", nodeId);
        sessionInfo.put("lastHeartbeatAt", String.valueOf(System.currentTimeMillis()));

        localSessions.put(session.getId(), session);

        return redisTemplate.opsForHash().putAll(sessionKey, sessionInfo)
                .then(redisTemplate.expire(sessionKey, Duration.ofHours(1)))
                .then(redisTemplate.opsForSet().add(userKey, session.getId()))
                .then(redisTemplate.opsForSet().add(nodeKey, session.getId()))
                .then();
    }

    public Mono<Void> unregister(WebSocketSession session) {
        String sessionId = session.getId();
        localSessions.remove(sessionId);
        String sessionKey = "ws:session:" + sessionId;

        return redisTemplate.opsForHash().get(sessionKey, "userId")
                .cast(String.class)
                .flatMap(userId -> {
                    String userKey = "ws:user:" + userId;
                    String nodeKey = "ws:node:" + nodeId;
                    return redisTemplate.opsForSet().remove(userKey, sessionId)
                            .then(redisTemplate.opsForSet().remove(nodeKey, sessionId));
                })
                .then(redisTemplate.delete(sessionKey))
                .then();
    }

    public Optional<WebSocketSession> getLocalSession(String sessionId) {
        return Optional.ofNullable(localSessions.get(sessionId));
    }
}
10.3 WebSocket 业务 Handler
@Component
public class AuthWebSocketHandler implements WebSocketHandler {

    private final JwtUtil jwtUtil;
    private final WebSocketSessionManager sessionManager;
    private final ObjectMapper objectMapper;

    @Override
    public Mono<Void> handle(WebSocketSession session) {
        URI uri = session.getHandshakeInfo().getUri();
        String token = UriComponentsBuilder.fromUri(uri).build().getQueryParams().getFirst("token");
        UserInfo user = jwtUtil.parseToken(token);
        if (user == null) {
            return session.close(CloseStatus.NOT_ACCEPTABLE);
        }

        return sessionManager.register(user, session)
                .thenMany(
                        session.receive()
                                .flatMap(message -> handleMessage(user, session, message))
                                .doOnError(ex -> log.error("ws error, user={}", user.getId(), ex))
                                .doFinally(sig -> sessionManager.unregister(session).subscribe())
                )
                .then();
    }

    private Mono<Void> handleMessage(UserInfo user, WebSocketSession session, WebSocketMessage message) {
        String text = message.getPayloadAsText();
        try {
            JsonNode node = objectMapper.readTree(text);
            String type = node.path("type").asText();
            if ("PING".equals(type)) {
                // 更新心跳
                // 此处可调用 sessionManager 更新 lastHeartbeatAt
                WsEnvelope<Void> pong = new WsEnvelope<>();
                pong.setType("PONG");
                pong.setTs(System.currentTimeMillis());
                String json = objectMapper.writeValueAsString(pong);
                return session.send(Mono.just(session.textMessage(json)));
            }
            // 其他业务消息处理...
            return Mono.empty();
        } catch (Exception e) {
            log.warn("invalid ws message, user={}, msg={}", user.getId(), text, e);
            return session.close(CloseStatus.BAD_DATA);
        }
    }
}
10.4 推送服务示例
@Service
public class WebSocketPushService {

    private final ReactiveStringRedisTemplate redisTemplate;
    private final WebSocketSessionManager sessionManager;
    private final ObjectMapper objectMapper;

    public Mono<Void> pushToUser(Long userId, Object payload) {
        String userKey = "ws:user:" + userId;
        return redisTemplate.opsForSet().members(userKey)
                .flatMap(sessionId -> {
                    Optional<WebSocketSession> opt = sessionManager.getLocalSession(sessionId);
                    if (!opt.isPresent()) {
                        return Mono.empty();
                    }
                    WebSocketSession session = opt.get();
                    return sendEnvelope(session, payload);
                })
                .then();
    }

    private Mono<Void> sendEnvelope(WebSocketSession session, Object bizPayload) {
        WsEnvelope<Object> envelope = new WsEnvelope<>();
        envelope.setType("NOTIFY");
        envelope.setVersion(1);
        envelope.setTs(System.currentTimeMillis());
        envelope.setPayload(bizPayload);

        return Mono.fromCallable(() -> objectMapper.writeValueAsString(envelope))
                .subscribeOn(Schedulers.boundedElastic())
                .flatMap(json -> session.send(Mono.just(session.textMessage(json))));
    }
}

11. 与现有业务逻辑的集成方案

11.1 点赞 / 关注 / 收藏等事件来源
  • 项目中已存在的 DTO:LikeMessageFollowMessageCollectMessage 等,可以直接作为 WebSocket payload 使用:
    • 业务服务在完成数据库写入后,向 RabbitMQ 投递对应消息;
    • WebSocket 推送服务消费这些消息,包装为 WsEnvelope 后推送给目标用户。
11.2 典型集成流程示例——点赞通知
  1. 用户 A 点赞视频 V(属用户 B);
  2. 点赞接口从 Controller → Service → Mapper 完成基础业务逻辑;
  3. Service 中发送 MQ 消息 LikeMessagews.notify.exchange
  4. WebSocket 推送服务消费 LikeMessage,构造 LikePayload
    • fromUserId = AtargetUserId = BvideoId = V
  5. 调用 WebSocketPushService.pushToUser(B, likePayload)
  6. 客户端收到 type = LIKE 的 WebSocket 消息,更新 UI(红点、弹窗、消息列表等)。
11.3 与幂等与风控系统的配合
  • WebSocket 推送只作为 最终一致的通知渠道,所有关键业务结果仍以 DB/缓存为准;
  • 对需要“只接收一次”的客户端逻辑,可在 WsEnvelope.bizId 填入业务唯一 ID,客户端按 bizId 做本地幂等;
  • 对敏感操作(如私信、转账类业务),仍需走现有幂等、风控、审计链路,仅使用 WebSocket 做状态更新或提醒。

以上指南结合现有技术栈(Spring WebFlux、Redis、RabbitMQ、Redisson、Jackson),从连接管理、分布式会话、异步推送、安全与性能优化等方面给出了系统化实践方案。落地时建议先从低风险的通知类场景(点赞/关注/收藏提醒)逐步引入 WebSocket,配合压测与监控不断调优参数和架构。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐