Vue3 + WebSocket + RocketMQ:实时消息推送全栈实战

简介: 客服系统消息延迟 5 秒,用户无法实时感知订单状态变更——这是我在电商项目中遇到的真实问题。通过 RocketMQ + WebSocket 的全栈方案,消息端到端延迟从 5000ms 降到 50ms,在线用户 10000+ 连接稳定运行。本文从实时推送的 5 大痛点出发,详解 Vue3 WebSocket 客户端 + Spring Boot + 阿里云 RocketMQ 5.x + Redis 的全栈架构设计与代码实战,涵盖 STOMP 协议适配、有序消息与事务消息、离线消息补拉、消息去重与幂等、连接安全认证等核心环节,并给出 5 个生产踩坑实录和推送方案选型决策树。

摘要:客服系统消息延迟 5 秒,用户无法实时感知订单状态变更——这是我在电商项目中遇到的真实问题。通过 RocketMQ + WebSocket 的全栈方案,消息端到端延迟从 5000ms 降到 50ms,在线用户 10000+ 连接稳定运行。本文从实时推送的 5 大痛点出发,详解 Vue3 WebSocket 客户端 + Spring Boot + 阿里云 RocketMQ 5.x + Redis 的全栈架构设计与代码实战,涵盖 STOMP 协议适配、有序消息与事务消息、离线消息补拉、消息去重与幂等、连接安全认证等核心环节,并给出 5 个生产踩坑实录和推送方案选型决策树。

1. 场景:客服系统的"5 秒延迟"困局

2025 年,我负责的电商平台客服系统上线后,用户反馈最强烈的问题就是"消息太慢"。客服回复一条消息,用户端要等 5 秒才能看到。订单状态从"已付款"变为"已发货",用户刷新好几次才知道。

问题根源:系统采用 HTTP 轮询方案,前端每 5 秒请求一次新消息接口。高并发时接口响应变慢,轮询间隔被迫拉长,延迟雪上加霜。

数据回顾:

指标 轮询方案 优化后 WebSocket+MQ
消息端到端延迟 2000~5000ms 30~80ms
服务端无效请求/分钟 12000 次 0 次(长连接推送)
单机连接数 200(HTTP 短连接) 8000+(WebSocket)
消息丢失率 1.2%(轮询间隙丢失) 0.01%(MQ 持久化)
带宽消耗 8Mbps(重复拉取) 0.5Mbps(增量推送)

改造后,RocketMQ 作为消息中枢保证消息可靠投递,WebSocket 实现服务端到客户端的实时推送,Redis 管理在线状态和离线消息队列,端到端延迟稳定在 50ms 以内,10000+ WebSocket 连接同时在线稳定运行。

RocketMQ+WebSocket 实时推送对比

下面从痛点分析开始,逐层展开全栈实战。


2. 实时推送的 5 大痛点

很多团队在实现实时推送时,低估了"可靠"二字的含金量。轮询看起来简单,但问题远不止"慢"。

痛点 1:轮询效率极低

HTTP 轮询是"拉"模型,客户端定时请求服务端获取新数据。问题是:大部分请求返回空结果,白白消耗带宽和连接资源。

对比项 HTTP 轮询 WebSocket
通信模式 客户端拉取 服务端推送
空请求占比 95%+ 0%
延迟 轮询间隔 + 接口耗时 接近 0(实时推送)
服务端压力 与轮询频率正相关 与消息量正相关

痛点 2:消息丢失无感知

轮询间隙产生的消息,如果客户端没有及时拉取,可能被新消息覆盖或因服务端缓冲区溢出而丢失。没有任何机制保证"至少投递一次"。

痛点 3:连接不稳定

移动端网络频繁切换(WiFi↔4G),WebSocket 连接容易断开。如果没有自动重连 + 心跳保活 + 断线检测,用户体验会频繁中断。

痛点 4:广播与定向推送难以兼得

系统通知需要广播给所有在线用户,客服消息需要定向推送给特定用户,订单变更需要推送给关联用户。不同场景的推送策略完全不同,统一设计困难。

痛点 5:离线消息补拉

用户断线重连或重新上线后,离线期间的消息如何补齐?需要一个离线消息队列,在用户上线时精准补拉,既不能漏也不能重复。

核心思路:用 RocketMQ 解决消息可靠性和投递问题,用 WebSocket 解决实时推送问题,用 Redis 解决在线状态和离线消息问题。三者各司其职,组合出完整的实时推送方案。


3. 整体架构设计

3.1 架构全景图

cloud-native_mermaid_1

3.2 消息流转链路

一条消息从产生到推送到客户端,完整链路如下:

cloud-native_mermaid_2

3.3 技术选型依据

组件 选型 为什么
消息队列 阿里云 RocketMQ 5.x 支持有序/事务/延迟消息,消息轨迹可追踪,与阿里云生态深度集成
实时通信 WebSocket + STOMP STOMP 提供订阅语义,SockJS 兼容旧浏览器降级
在线状态 Redis 高性能读写,Set 结构天然适合在线用户管理
前端框架 Vue3 + Composition API 响应式 + 组合式 API 天然适合消息状态管理

4. 阿里云 RocketMQ 5.x 配置

4.1 实例创建

在阿里云控制台创建 RocketMQ 5.x 实例,选择标准版(支持消息轨迹和重试策略):

# 实例配置建议
实例类型: 标准版(5.x系列)
地域: 与应用服务器同地域(降低网络延迟)
规格: 根据TPS选择,建议预留50%余量
  - 入门: 500 TPS(日均消息量 < 1000万)
  - 生产: 2000 TPS(日均消息量 < 5000万)
消息保留: 72小时(满足离线补拉窗口期)
消息轨迹: 开启(排查消息丢失必备)

4.2 Topic 与 Tag 规划

为什么用 Topic + Tag 二级分类?Topic 是物理隔离,Tag 是逻辑过滤。同一业务域的消息共享 Topic,用 Tag 区分消息类型,减少 Topic 数量降低成本。

Topic 用途 Tag 说明
ORDER_MSG 订单消息 status 订单状态变更(已付款、已发货、已完成)
ORDER_MSG 订单消息 refund 退款消息
NOTIFY_MSG 通知消息 system 系统公告
NOTIFY_MSG 通知消息 promotion 营销推送
CHAT_MSG 聊天消息 text 文本消息
CHAT_MSG 聊天消息 image 图片消息
CHAT_MSG 聊天消息 file 文件消息

4.3 消费者 Group 与重试策略

# 消费者组规划
Group ID: 
  - ws-order-consumer:  订单消息推送消费者
  - ws-notify-consumer: 通知消息推送消费者
  - ws-chat-consumer:   聊天消息推送消费者

# 重试策略(控制台配置)
最大重试次数: 16
重试间隔: 递增策略(10s, 30s, 1min, 2min, 3min, 4min, 5min, 6min, 7min, 8min, 9min, 10min, 20min, 30min, 1h, 2h)
消费超时: 15分钟(超过未ACK自动重试)
死信队列: 开启(重试16次后进入 %DLQ% Group ID)

4.4 消息轨迹与监控

消息轨迹是排查问题的利器。开启后,每条消息从发送到消费的全链路都可追踪:

// 生产者开启消息轨迹
DefaultMQProducer producer = new DefaultMQProducer(
    "ws-producer-group",
    true  // enableTrace = true,开启消息轨迹
);

监控告警建议配置:

监控指标 告警阈值 说明
消息堆积量 > 10000 条 消费速度跟不上生产速度
消费延迟 > 5s 消息从产生到被消费的时间
发送失败率 > 0.1% 可能是网络或Broker问题
死信队列消息数 > 0 有消息消费失败16次

5. 后端 Spring Boot 实战

5.1 项目依赖与配置

为什么选择 STOMP 协议?原生 WebSocket 只有帧级别的通信,没有订阅、目的地等语义。STOMP 在 WebSocket 之上提供发布/订阅模型,前端可以按 Topic 订阅,服务端可以定向推送。

<!-- pom.xml 核心依赖 -->
<dependencies>
    <!-- WebSocket + STOMP -->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-websocket</artifactId>
    </dependency>

    <!-- RocketMQ 5.x 客户端 -->
    <dependency>
        <groupId>org.apache.rocketmq</groupId>
        <artifactId>rocketmq-client-java</artifactId>
        <version>5.0.7</version>
    </dependency>

    <!-- Redis -->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-data-redis</artifactId>
    </dependency>
</dependencies>
# application.yml
rocketmq:
  name-server: ${
   ROCKETMQ_NAMESRV:rmq-xxxx.cn-hangzhou.rmq.aliyuncs.com:8080}
  producer:
    group: ws-producer-group
    send-message-timeout: 3000
    retry-times-when-send-failed: 3
  consumer:
    topics: ORDER_MSG,NOTIFY_MSG,CHAT_MSG

spring:
  data:
    redis:
      host: ${
   REDIS_HOST:127.0.0.1}
      port: 6379
      password: ${
   REDIS_PASSWORD:}
      lettuce:
        pool:
          max-active: 50
          max-idle: 20

websocket:
  max-connections: 10000
  heartbeat-interval: 30000
  offline-queue-max-size: 500

5.2 WebSocket 配置(STOMP + SockJS 降级)

/**
 * WebSocket配置:STOMP协议 + SockJS降级方案
 * 
 * Why:SockJS在WebSocket不可用时自动降级为HTTP长轮询,
 * 兼容公司内网限制WebSocket的老旧浏览器和代理服务器
 */
@Configuration
@EnableWebSocketMessageBroker
public class WebSocketConfig implements WebSocketMessageBrokerConfigurer {
   

    @Value("${websocket.max-connections:10000}")
    private int maxConnections;

    @Override
    public void configureMessageBroker(MessageBrokerRegistry config) {
   
        // 服务端推送目的地前缀:/topic(广播)、/queue(定向)
        config.enableSimpleBroker("/topic", "/queue");
        // 客户端发送消息目的地前缀
        config.setApplicationDestinationPrefixes("/app");
        // 定向推送用户目的地前缀
        config.setUserDestinationPrefix("/user");
    }

    @Override
    public void registerStompEndpoints(StompEndpointRegistry registry) {
   
        registry.addEndpoint("/ws")
                .setAllowedOriginPatterns("*")
                .withSockJS();  // SockJS降级支持
    }

    @Override
    public void configureWebSocketTransport(WebSocketTransportRegistration registry) {
   
        // 消息大小限制:图片/文件消息可能较大
        registry.setMessageSizeLimit(512 * 1024);
        // 发送缓冲区大小
        registry.setSendBufferSizeLimit(1024 * 1024);
        // 发送时间限制
        registry.setSendTimeLimit(20 * 1000);
    }
}

5.3 用户连接管理(Redis 在线状态 + Session 映射)

为什么用 Redis 管理在线状态?多实例部署时,本地 Map 无法跨节点查询用户在线状态。Redis 作为共享存储,任何节点都能判断用户是否在线。

/**
 * WebSocket连接事件监听器
 * 
 * Why:监听连接建立/断开事件,维护Redis在线状态和Session映射,
 * 这是消息能否精准推送的关键
 */
@Component
@Slf4j
@RequiredArgsConstructor
public class WebSocketEventListener {
   

    private final RedisTemplate<String, String> redisTemplate;
    private final SimpUserRegistry userRegistry;

    private static final String ONLINE_USERS_KEY = "ws:online";
    private static final String SESSION_MAP_KEY = "ws:session";
    private static final long ONLINE_TTL = 120; // 秒

    /**
     * 连接建立:记录在线状态
     */
    @EventListener
    public void handleConnect(SessionConnectEvent event) {
   
        StompHeaderAccessor accessor = StompHeaderAccessor.wrap(event.getMessage());
        String userId = accessor.getUser().getName();
        String sessionId = accessor.getSessionId();

        // 添加到在线用户Set
        redisTemplate.opsForSet().add(ONLINE_USERS_KEY, userId);
        // 记录Session映射(用户->会话ID),支持同用户多端登录
        redisTemplate.opsForHash().put(SESSION_MAP_KEY, userId + ":" + sessionId, 
            String.valueOf(System.currentTimeMillis()));

        log.info("用户上线: userId={}, sessionId={}, 当前在线={}", 
            userId, sessionId, redisTemplate.opsForSet().size(ONLINE_USERS_KEY));
    }

    /**
     * 连接断开:移除在线状态
     */
    @EventListener
    public void handleDisconnect(SessionDisconnectEvent event) {
   
        StompHeaderAccessor accessor = StompHeaderAccessor.wrap(event.getMessage());
        String userId = accessor.getUser().getName();
        String sessionId = accessor.getSessionId();

        // 移除Session映射
        redisTemplate.opsForHash().delete(SESSION_MAP_KEY, userId + ":" + sessionId);

        // 检查该用户是否还有其他Session(多端登录场景)
        boolean hasOtherSession = redisTemplate.opsForHash().keys(SESSION_MAP_KEY)
            .stream().anyMatch(k -> ((String) k).startsWith(userId + ":"));

        if (!hasOtherSession) {
   
            // 无其他Session,从在线Set移除
            redisTemplate.opsForSet().remove(ONLINE_USERS_KEY, userId);
            log.info("用户离线: userId={}", userId);
        }
    }

    /**
     * 判断用户是否在线
     */
    public boolean isUserOnline(String userId) {
   
        return Boolean.TRUE.equals(redisTemplate.opsForSet().isMember(ONLINE_USERS_KEY, userId));
    }
}

5.4 RocketMQ 生产者(有序消息 + 事务消息 + 延迟消息)

/**
 * RocketMQ消息生产者
 * 
 * Why:三种消息类型覆盖不同业务场景:
 * - 有序消息:订单状态必须按顺序推送(已付款->已发货->已完成)
 * - 事务消息:下单+发消息必须原子成功,避免下单成功但消息没发
 * - 延迟消息:下单30分钟未支付自动取消
 */
@Component
@Slf4j
@RequiredArgsConstructor
public class MessageProducer {
   

    private final RocketMQTemplate rocketMQTemplate;

    /**
     * 发送有序消息(订单状态变更)
     * 
     * Why:同一订单的状态变更消息必须按顺序消费,
     * 否则可能出现"已发货"先于"已付款"到达客户端
     * 
     * @param orderId 订单ID,作为sharding key保证同一订单消息进入同一队列
     */
    public SendResult sendOrderlyMessage(String orderId, String status, String userId) {
   
        OrderMessage msg = OrderMessage.builder()
            .orderId(orderId)
            .status(status)
            .userId(userId)
            .timestamp(System.currentTimeMillis())
            .build();

        // hashKey = orderId,保证同一订单消息有序
        SendResult result = rocketMQTemplate.syncSendOrderly(
            "ORDER_MSG:status",
            MessageBuilder.withPayload(msg).build(),
            orderId  // sharding key
        );
        log.info("有序消息发送成功: orderId={}, status={}, msgId={}", 
            orderId, status, result.getMsgId());
        return result;
    }

    /**
     * 发送事务消息(下单成功后发消息)
     * 
     * Why:事务消息保证"本地事务+消息发送"的原子性。
     * 如果本地事务(创建订单)失败,消息不会发出;
     * 如果消息发送失败,本地事务回滚
     */
    public void sendTransactionMessage(String orderId, String userId) {
   
        OrderMessage msg = OrderMessage.builder()
            .orderId(orderId)
            .status("CREATED")
            .userId(userId)
            .timestamp(System.currentTimeMillis())
            .build();

        rocketMQTemplate.sendMessageInTransaction(
            "ORDER_MSG:status",
            MessageBuilder.withPayload(msg).build(),
            orderId  // 传递给本地事务监听器的参数
        );
    }

    /**
     * 发送延迟消息(超时未支付自动取消)
     * 
     * Why:延迟消息不需要定时任务轮询扫描,
     * RocketMQ在指定延迟级别后自动投递,可靠性更高
     */
    public SendResult sendDelayMessage(String orderId, String userId) {
   
        OrderMessage msg = OrderMessage.builder()
            .orderId(orderId)
            .status("CANCEL_TIMEOUT")
            .userId(userId)
            .timestamp(System.currentTimeMillis())
            .build();

        // delayLevel=16对应30分钟延迟
        // 1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
        SendResult result = rocketMQTemplate.syncSend(
            "ORDER_MSG:status",
            MessageBuilder.withPayload(msg).build(),
            3000,  // timeout
            16     // delayLevel=16,即30分钟
        );
        log.info("延迟消息发送成功: orderId={}, 30分钟后自动取消", orderId);
        return result;
    }
}

/**
 * 事务消息监听器
 */
@RocketMQTransactionListener
@RequiredArgsConstructor
public class OrderTransactionListener implements RocketMQLocalTransactionListener {
   

    private final OrderService orderService;

    @Override
    public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
   
        String orderId = (String) arg;
        try {
   
            // 执行本地事务:创建订单
            orderService.createOrder(orderId);
            return RocketMQLocalTransactionState.COMMIT;
        } catch (Exception e) {
   
            log.error("本地事务执行失败: orderId={}", orderId, e);
            return RocketMQLocalTransactionState.ROLLBACK;
        }
    }

    @Override
    public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
   
        String orderId = (String) msg.getHeaders().get("orderId");
        // 回查本地事务状态
        return orderService.existsById(orderId) 
            ? RocketMQLocalTransactionState.COMMIT
            : RocketMQLocalTransactionState.ROLLBACK;
    }
}

5.5 RocketMQ 消费者(广播/集群 + 重试 + 死信)

/**
 * 聊天消息消费者
 * 
 * Why:聊天消息使用集群模式(Clustering),同一Group只消费一次,
 * 由推送Service决定推送给哪个用户
 */
@Component
@RocketMQMessageListener(
    topic = "CHAT_MSG",
    selectorExpression = "text || image || file",
    consumerGroup = "ws-chat-consumer",
    consumeMode = ConsumeMode.CONCURRENTLY,
    maxReconsumeTimes = 16
)
@Slf4j
@RequiredArgsConstructor
public class ChatMessageConsumer implements RocketMQListener<MessageExt> {
   

    private final MessagePushService pushService;

    @Override
    public void onMessage(MessageExt messageExt) {
   
        try {
   
            ChatMessage chatMsg = JSON.parseObject(
                new String(messageExt.getBody(), StandardCharsets.UTF_8),
                ChatMessage.class
            );

            // 消息去重:用业务幂等键判断
            if (pushService.isDuplicate(chatMsg.getMsgId())) {
   
                log.warn("重复消息,跳过: msgId={}", chatMsg.getMsgId());
                return;
            }

            // 推送给目标用户
            pushService.pushToUser(chatMsg.getToUserId(), "/queue/chat", chatMsg);

        } catch (Exception e) {
   
            log.error("聊天消息消费失败: msgId={}", messageExt.getMsgId(), e);
            // 抛出异常触发重试
            throw new RuntimeException("消费失败,等待重试", e);
        }
    }
}

/**
 * 通知消息消费者(广播模式)
 * 
 * Why:系统公告需要推送到所有在线用户,
 * 广播模式下每个消费者实例都会收到消息
 */
@Component
@RocketMQMessageListener(
    topic = "NOTIFY_MSG",
    selectorExpression = "system || promotion",
    consumerGroup = "ws-notify-consumer",
    messageModel = MessageModel.BROADCASTING
)
@Slf4j
@RequiredArgsConstructor
public class NotifyMessageConsumer implements RocketMQListener<MessageExt> {
   

    private final MessagePushService pushService;

    @Override
    public void onMessage(MessageExt messageExt) {
   
        NotifyMessage notifyMsg = JSON.parseObject(
            new String(messageExt.getBody(), StandardCharsets.UTF_8),
            NotifyMessage.class
        );
        // 广播推送给所有在线用户
        pushService.broadcastToAll("/topic/notify", notifyMsg);
    }
}

5.6 消息推送 Service(定向/广播/标签过滤)

/**
 * 消息推送核心Service
 * 
 * Why:统一封装推送逻辑,处理在线推送和离线存储的分支,
 * 所有消费者都通过此Service推送,避免逻辑散落各处
 */
@Service
@Slf4j
@RequiredArgsConstructor
public class MessagePushService {
   

    private final SimpMessagingTemplate messagingTemplate;
    private final WebSocketEventListener connectionManager;
    private final RedisTemplate<String, String> redisTemplate;

    private static final String OFFLINE_QUEUE_PREFIX = "ws:offline:";
    private static final String DEDUP_PREFIX = "ws:dedup:";
    private static final int OFFLINE_QUEUE_MAX_SIZE = 500;

    /**
     * 定向推送给指定用户
     * 
     * 在线:通过WebSocket推送
     * 离线:写入Redis离线队列
     */
    public void pushToUser(String userId, String destination, Object payload) {
   
        if (connectionManager.isUserOnline(userId)) {
   
            try {
   
                // STOMP定向推送:/user/{userId}{destination}
                messagingTemplate.convertAndSendToUser(userId, destination, payload);
                log.debug("推送成功: userId={}, dest={}", userId, destination);
            } catch (Exception e) {
   
                // 推送失败(Session可能已失效),转离线存储
                log.warn("推送失败,转离线: userId={}", userId, e);
                saveOfflineMessage(userId, payload);
            }
        } else {
   
            // 用户离线,写入离线队列
            saveOfflineMessage(userId, payload);
        }
    }

    /**
     * 广播推送给所有在线用户
     */
    public void broadcastToAll(String destination, Object payload) {
   
        messagingTemplate.convertAndSend(destination, payload);
        log.debug("广播推送: dest={}", destination);
    }

    /**
     * 按标签过滤推送(只推送给订阅了指定Tag的用户)
     * 
     * Why:营销推送不应该推送给关闭了营销通知的用户
     */
    public void pushByTag(String tag, String destination, Object payload) {
   
        Set<String> subscribers = redisTemplate.opsForSet()
            .members("ws:subscribe:" + tag);
        if (subscribers != null) {
   
            subscribers.forEach(userId -> pushToUser(userId, destination, payload));
        }
    }

    /**
     * 保存离线消息到Redis List
     */
    private void saveOfflineMessage(String userId, Object payload) {
   
        String key = OFFLINE_QUEUE_PREFIX + userId;
        String msgJson = JSON.toJSONString(payload);

        // LPUSH写入队列头部
        redisTemplate.opsForList().leftPush(key, msgJson);
        // 限制队列长度,防止内存泄漏
        redisTemplate.opsForList().trim(key, 0, OFFLINE_QUEUE_MAX_SIZE - 1);
        // 设置过期时间7天
        redisTemplate.expire(key, 7, TimeUnit.DAYS);
    }

    /**
     * 获取离线消息(用户上线后调用)
     */
    public List<String> getOfflineMessages(String userId) {
   
        String key = OFFLINE_QUEUE_PREFIX + userId;
        // LRANGE获取所有离线消息
        List<String> messages = redisTemplate.opsForList().range(key, 0, -1);
        // 获取后删除队列
        redisTemplate.delete(key);
        return messages != null ? messages : Collections.emptyList();
    }

    /**
     * 消息去重判断
     * 
     * Why:RocketMQ至少投递一次,重复消费是正常现象,
     * 必须在业务层做幂等
     */
    public boolean isDuplicate(String msgId) {
   
        String key = DEDUP_PREFIX + msgId;
        // SETNX:不存在才设置成功,返回true表示重复
        Boolean isDuplicate = redisTemplate.opsForValue()
            .setIfAbsent(key, "1", 5, TimeUnit.MINUTES);
        return Boolean.TRUE.equals(isDuplicate) == false;
    }
}

6. 前端 Vue3 实战

6.1 WebSocket 连接管理(自动重连 + 心跳 + 断线检测)

为什么用 STOMP.js 而不是原生 WebSocket?STOMP 提供了订阅语义、事务支持和 ACK 机制,配合 SockJS 自动降级,开发体验远优于原生 WebSocket。

// composables/useWebSocket.ts

import {
    ref, onMounted, onUnmounted } from 'vue'
import SockJS from 'sockjs-client'
import {
    Client, IMessage } from '@stomp/stompjs'

// Why:连接状态枚举,用于UI展示连接状态指示器
export type ConnectionStatus = 'connecting' | 'connected' | 'disconnected' | 'reconnecting'

export function useWebSocket(userId: string) {
   
  const status = ref<ConnectionStatus>('disconnected')
  const reconnectCount = ref(0)
  let stompClient: Client | null = null
  let heartbeatTimer: ReturnType<typeof setInterval> | null = null

  const MAX_RECONNECT_ATTEMPTS = 10
  const HEARTBEAT_INTERVAL = 30000 // 30秒心跳

  /**
   * 建立STOMP连接
   * 
   * Why:SockJS作为传输层自动降级,STOMP作为协议层提供订阅语义
   * 前端无需关心底层是WebSocket还是HTTP长轮询
   */
  function connect() {
   
    status.value = 'connecting'

    stompClient = new Client({
   
      webSocketFactory: () => new SockJS('/ws'),
      reconnectDelay: 0, // 手动控制重连
      heartbeatIncoming: 20000,
      heartbeatOutgoing: 20000,
      onConnect: () => {
   
        status.value = 'connected'
        reconnectCount.value = 0
        startHeartbeat()
        // 连接成功后拉取离线消息
        fetchOfflineMessages()
      },
      onDisconnect: () => {
   
        status.value = 'disconnected'
        stopHeartbeat()
      },
      onStompError: (frame) => {
   
        console.error('STOMP错误:', frame.headers['message'])
        scheduleReconnect()
      },
      onWebSocketClose: () => {
   
        status.value = 'disconnected'
        stopHeartbeat()
        scheduleReconnect()
      }
    })

    stompClient.activate()
  }

  /**
   * 自动重连(指数退避)
   * 
   * Why:网络波动时立即重连会雪崩,指数退避避免服务端被重连请求打爆
   */
  function scheduleReconnect() {
   
    if (reconnectCount.value >= MAX_RECONNECT_ATTEMPTS) {
   
      status.value = 'disconnected'
      console.error('超过最大重连次数,请手动刷新页面')
      return
    }

    status.value = 'reconnecting'
    reconnectCount.value++

    // 指数退避:1s, 2s, 4s, 8s, 16s, 32s...
    const delay = Math.min(1000 * Math.pow(2, reconnectCount.value - 1), 32000)
    console.log(`${
     delay}ms 后尝试第${
     reconnectCount.value}次重连...`)

    setTimeout(() => {
   
      connect()
    }, delay)
  }

  /**
   * 心跳保活
   * 
   * Why:Nginx默认60秒超时关闭空闲连接,
   * 30秒心跳保持连接活跃
   */
  function startHeartbeat() {
   
    stopHeartbeat()
    heartbeatTimer = setInterval(() => {
   
      if (stompClient?.active) {
   
        stompClient.publish({
   
          destination: '/app/heartbeat',
          body: JSON.stringify({
    userId, timestamp: Date.now() })
        })
      }
    }, HEARTBEAT_INTERVAL)
  }

  function stopHeartbeat() {
   
    if (heartbeatTimer) {
   
      clearInterval(heartbeatTimer)
      heartbeatTimer = null
    }
  }

  /**
   * 订阅定向消息(聊天、订单等)
   */
  function subscribeUserQueue(
    destination: string,
    callback: (message: IMessage) => void
  ) {
   
    if (!stompClient?.active) return
    // /user/queue/chat -> 服务端推送给当前用户
    stompClient.subscribe(`/user${
     destination}`, callback)
  }

  /**
   * 订阅广播消息(系统公告等)
   */
  function subscribeTopic(
    destination: string,
    callback: (message: IMessage) => void
  ) {
   
    if (!stompClient?.active) return
    stompClient.subscribe(destination, callback)
  }

  /**
   * 拉取离线消息
   */
  async function fetchOfflineMessages() {
   
    try {
   
      const res = await fetch(`/api/messages/offline?userId=${
     userId}`)
      const messages = await res.json()
      // 触发离线消息回调
      messages.forEach((msg: any) => {
   
        // 根据消息类型分发到对应处理器
        window.dispatchEvent(
          new CustomEvent('offline-message', {
    detail: msg })
        )
      })
    } catch (e) {
   
      console.error('拉取离线消息失败:', e)
    }
  }

  function disconnect() {
   
    stopHeartbeat()
    stompClient?.deactivate()
    status.value = 'disconnected'
  }

  onMounted(() => connect())
  onUnmounted(() => disconnect())

  return {
   
    status,
    reconnectCount,
    connect,
    disconnect,
    subscribeUserQueue,
    subscribeTopic
  }
}

6.2 消息通知组件(弹窗 + 未读计数 + 声音提醒)

<!-- components/MessageNotification.vue -->

<script setup lang="ts">
import { ref, computed, onMounted, onUnmounted } from 'vue'
import { useWebSocket } from '../composables/useWebSocket'
import { ElNotification } from 'element-plus'

// Why:通知类型映射不同图标和样式
type NotifyType = 'order' | 'chat' | 'system'

interface PushMessage {
  id: string
  type: NotifyType
  title: string
  content: string
  timestamp: number
  read: boolean
}

const props = defineProps<{ userId: string }>()

const { status, subscribeUserQueue, subscribeTopic } = useWebSocket(props.userId)
const messages = ref<PushMessage[]>([])
const unreadCount = computed(() => messages.value.filter(m => !m.read).length)

// Why:声音提醒在客服场景很重要,新消息不响铃会漏接
const notificationAudio = new Audio('/assets/notification.mp3')
notificationAudio.volume = 0.3

/**
 * 处理推送消息
 */
function handlePushMessage(msgBody: string) {
  const msg: PushMessage = JSON.parse(msgBody)

  // 去重:避免重复展示
  if (messages.value.some(m => m.id === msg.id)) return

  messages.value.unshift(msg)

  // 弹窗通知
  ElNotification({
    title: msg.title,
    message: msg.content,
    type: msg.type === 'order' ? 'success' : msg.type === 'chat' ? 'info' : 'warning',
    duration: 3000,
    onClick: () => markAsRead(msg.id)
  })

  // 声音提醒(非系统公告时响铃)
  if (msg.type !== 'system') {
    notificationAudio.play().catch(() => {})
  }
}

function markAsRead(msgId: string) {
  const msg = messages.value.find(m => m.id === msgId)
  if (msg) msg.read = true
}

function markAllAsRead() {
  messages.value.forEach(m => (m.read = true))
}

onMounted(() => {
  // 订阅聊天消息
  subscribeUserQueue('/queue/chat', (message) => {
    handlePushMessage(message.body)
  })
  // 订阅订单消息
  subscribeUserQueue('/queue/order', (message) => {
    handlePushMessage(message.body)
  })
  // 订阅系统广播
  subscribeTopic('/topic/notify', (message) => {
    handlePushMessage(message.body)
  })
})
</script>

<template>
  <div class="notification-bell" @click="markAllAsRead">
    <el-badge :value="unreadCount" :hidden="unreadCount === 0" :max="99">
      <el-icon :size="20"><Bell /></el-icon>
    </el-badge>
    <!-- 连接状态指示 -->
    <span 
      class="status-dot" 
      :class="status"
      :title="`连接状态: ${status}`"
    />
  </div>
</template>

<style scoped>
.notification-bell {
  position: relative;
  cursor: pointer;
}
.status-dot {
  position: absolute;
  top: -2px;
  right: -2px;
  width: 8px;
  height: 8px;
  border-radius: 50%;
}
.status-dot.connected { background: #67c23a; }
.status-dot.connecting,
.status-dot.reconnecting { background: #e6a23c; animation: blink 1s infinite; }
.status-dot.disconnected { background: #f56c6c; }

@keyframes blink {
  50% { opacity: 0.3; }
}
</style>

6.3 聊天界面(消息列表 + 输入框 + 文件发送 + 已读回执)

<!-- components/ChatPanel.vue -->

<script setup lang="ts">
import { ref, nextTick, onMounted } from 'vue'
import { useWebSocket } from '../composables/useWebSocket'
import { uploadFile } from '../api/upload'

interface ChatMsg {
  msgId: string
  fromUserId: string
  toUserId: string
  content: string
  type: 'text' | 'image' | 'file'
  fileName?: string
  fileUrl?: string
  timestamp: number
  read: boolean
}

const props = defineProps<{
  userId: string
  targetUserId: string
  targetUserName: string
}>()

const { subscribeUserQueue } = useWebSocket(props.userId)
const messageList = ref<ChatMsg[]>([])
const inputText = ref('')
const isLoading = ref(false)
const chatContainer = ref<HTMLDivElement>()

/**
 * 发送文本消息
 * 
 * Why:通过HTTP接口发送而非WebSocket,因为消息需要经过RocketMQ持久化,
 * HTTP接口保证消息先落库再投递MQ,WebSocket只负责接收推送
 */
async function sendTextMessage() {
  if (!inputText.value.trim()) return

  const msg: ChatMsg = {
    msgId: `${props.userId}-${Date.now()}`,
    fromUserId: props.userId,
    toUserId: props.targetUserId,
    content: inputText.value.trim(),
    type: 'text',
    timestamp: Date.now(),
    read: false
  }

  try {
    await fetch('/api/chat/send', {
      method: 'POST',
      headers: { 'Content-Type': 'application/json' },
      body: JSON.stringify(msg)
    })
    messageList.value.push(msg)
    inputText.value = ''
    scrollToBottom()
  } catch (e) {
    console.error('发送失败:', e)
  }
}

/**
 * 发送文件消息
 */
async function sendFileMessage(file: File) {
  isLoading.value = true
  try {
    // 先上传文件到OSS
    const fileUrl = await uploadFile(file)
    const msg: ChatMsg = {
      msgId: `${props.userId}-${Date.now()}`,
      fromUserId: props.userId,
      toUserId: props.targetUserId,
      content: file.name,
      type: file.type.startsWith('image/') ? 'image' : 'file',
      fileName: file.name,
      fileUrl,
      timestamp: Date.now(),
      read: false
    }
    await fetch('/api/chat/send', {
      method: 'POST',
      headers: { 'Content-Type': 'application/json' },
      body: JSON.stringify(msg)
    })
    messageList.value.push(msg)
    scrollToBottom()
  } catch (e) {
    console.error('文件发送失败:', e)
  } finally {
    isLoading.value = false
  }
}

/**
 * 发送已读回执
 * 
 * Why:已读回执通过WebSocket发送,实时性要求高,不需要持久化
 */
function sendReadReceipt(msgId: string) {
  // 通过STOMP发送已读回执
  const msg = messageList.value.find(m => m.msgId === msgId)
  if (msg && !msg.read) {
    msg.read = true
    // 已读回执走WebSocket,轻量级无需MQ
  }
}

function scrollToBottom() {
  nextTick(() => {
    if (chatContainer.value) {
      chatContainer.value.scrollTop = chatContainer.value.scrollHeight
    }
  })
}

onMounted(() => {
  // 订阅对方发来的聊天消息
  subscribeUserQueue('/queue/chat', (message) => {
    const msg: ChatMsg = JSON.parse(message.body)
    if (msg.fromUserId === props.targetUserId) {
      messageList.value.push(msg)
      scrollToBottom()
      // 自动发送已读回执(窗口处于激活状态时)
      if (document.visibilityState === 'visible') {
        sendReadReceipt(msg.msgId)
      }
    }
  })
})
</script>

<template>
  <div class="chat-panel">
    <div class="chat-header">
      <span>{
  { targetUserName }}</span>
    </div>
    <div ref="chatContainer" class="chat-messages">
      <div 
        v-for="msg in messageList" 
        :key="msg.msgId"
        class="message-item"
        :class="{ 'self': msg.fromUserId === userId }"
      >
        <div v-if="msg.type === 'text'" class="msg-text">{
  { msg.content }}</div>
        <div v-else-if="msg.type === 'image'" class="msg-image">
          <img :src="msg.fileUrl" :alt="msg.fileName" />
        </div>
        <div v-else class="msg-file">
          <a :href="msg.fileUrl" target="_blank">{
  { msg.fileName }}</a>
        </div>
        <span class="msg-time">
          {
  { new Date(msg.timestamp).toLocaleTimeString() }}
          <span v-if="msg.fromUserId === userId" class="read-status">
            {
  { msg.read ? '已读' : '未读' }}
          </span>
        </span>
      </div>
    </div>
    <div class="chat-input">
      <input 
        v-model="inputText" 
        placeholder="输入消息..." 
        @keyup.enter="sendTextMessage" 
      />
      <input 
        type="file" 
        @change="(e) => sendFileMessage((e.target as any).files[0])" 
      />
      <el-button @click="sendTextMessage" :loading="isLoading">发送</el-button>
    </div>
  </div>
</template>

6.4 离线消息补拉

// composables/useOfflineSync.ts

import {
    useWebSocket } from './useWebSocket'

/**
 * 离线消息补拉模块
 * 
 * Why:用户断线重连或重新上线后,离线期间的消息通过HTTP接口补拉,
 * 补拉完成后再接收实时推送,避免消息乱序
 */
export function useOfflineSync(userId: string) {
   
  const {
    status, subscribeUserQueue, subscribeTopic } = useWebSocket(userId)

  async function syncOfflineMessages() {
   
    if (status.value !== 'connected') return

    try {
   
      // 1. 先拉取离线消息
      const res = await fetch(`/api/messages/offline?userId=${
     userId}`)
      const offlineMessages = await res.json()

      // 2. 按时间排序,保证消息顺序
      offlineMessages.sort((a: any, b: any) => a.timestamp - b.timestamp)

      // 3. 逐条处理离线消息
      for (const msg of offlineMessages) {
   
        window.dispatchEvent(
          new CustomEvent('offline-message', {
    detail: msg })
        )
      }

      // 4. 离线消息处理完毕后,再订阅实时推送
      // Why:避免离线消息和实时消息交叉乱序
      subscribeUserQueue('/queue/chat', handleRealtimeMessage)
      subscribeUserQueue('/queue/order', handleRealtimeMessage)
      subscribeTopic('/topic/notify', handleRealtimeMessage)

      console.log(`离线消息补拉完成: ${
     offlineMessages.length}条`)
    } catch (e) {
   
      console.error('离线消息补拉失败:', e)
    }
  }

  function handleRealtimeMessage(message: any) {
   
    const msg = JSON.parse(message.body)
    window.dispatchEvent(
      new CustomEvent('realtime-message', {
    detail: msg })
    )
  }

  return {
    syncOfflineMessages }
}

7. 进阶优化

7.1 消息可靠性保障

实时推送场景下,消息可靠性需要端到端保障,任何一个环节掉链子都会导致消息丢失:

cloud-native_mermaid_3

三重保障机制:

环节 保障手段 失败兜底
发送端 同步发送 + 重试3次 记录发送失败表,定时任务补偿
MQ端 同步刷盘 + 主从复制 消息轨迹追踪,定位丢失环节
消费端 手动ACK + 重试16次 死信队列人工处理

7.2 消息去重(Redis Set + 业务幂等键)

/**
 * 消息去重增强方案
 * 
 * Why:两层去重保证不重复消费:
 * - Redis SETNX:快速去重,5分钟过期
 * - 数据库唯一索引:持久化去重,防止Redis缓存失效后重复
 */
@Service
@RequiredArgsConstructor
public class MessageDedupService {
   

    private final RedisTemplate<String, String> redisTemplate;
    private final MessageLogMapper messageLogMapper;

    /**
     * 幂等检查(Redis + DB双层)
     */
    public boolean isDuplicate(String msgId) {
   
        // 第一层:Redis快速判断
        String dedupKey = "ws:dedup:" + msgId;
        Boolean redisResult = redisTemplate.opsForValue()
            .setIfAbsent(dedupKey, "1", 5, TimeUnit.MINUTES);

        if (Boolean.FALSE.equals(redisResult)) {
   
            return true; // Redis判定重复
        }

        // 第二层:数据库唯一索引兜底
        try {
   
            messageLogMapper.insertIfAbsent(msgId);
            return false;
        } catch (DuplicateKeyException e) {
   
            // 数据库判定重复,回补Redis
            redisTemplate.opsForValue().set(dedupKey, "1", 5, TimeUnit.MINUTES);
            return true;
        }
    }
}

7.3 连接安全(JWT 认证 + WSS 加密)

/**
 * WebSocket握手拦截器:JWT认证
 * 
 * Why:WebSocket连接建立前验证身份,防止未授权连接
 * STOMP的CONNECT帧携带JWT token
 */
@Component
@RequiredArgsConstructor
public class WebSocketAuthInterceptor implements HandshakeInterceptor {
   

    private final JwtTokenProvider jwtTokenProvider;

    @Override
    public boolean beforeHandshake(
        ServerHttpRequest request, ServerHttpResponse response,
        WebSocketHandler wsHandler, Map<String, Object> attributes
    ) {
   
        if (request instanceof ServletServerHttpRequest) {
   
            HttpServletRequest servletRequest = 
                ((ServletServerHttpRequest) request).getServletRequest();
            // 从请求参数或Header获取token
            String token = servletRequest.getParameter("token");
            if (token == null) {
   
                token = servletRequest.getHeader("Authorization");
                if (token != null && token.startsWith("Bearer ")) {
   
                    token = token.substring(7);
                }
            }

            if (token != null && jwtTokenProvider.validateToken(token)) {
   
                String userId = jwtTokenProvider.getUserIdFromToken(token);
                attributes.put("userId", userId);
                return true;
            }
        }
        return false; // 认证失败,拒绝连接
    }

    @Override
    public void afterHandshake(
        ServerHttpRequest request, ServerHttpResponse response,
        WebSocketHandler wsHandler, Exception exception
    ) {
   }
}

Nginx 配置 WSS 加密:

# Why:生产环境必须WSS加密,否则HTTP环境下WebSocket无法建立
server {
   
    listen 443 ssl;
    server_name api.example.com;

    ssl_certificate     /etc/nginx/ssl/cert.pem;
    ssl_certificate_key /etc/nginx/ssl/key.pem;

    location /ws {
   
        proxy_pass http://backend:8080/ws;
        proxy_http_version 1.1;
        # WebSocket必需的Upgrade头
        proxy_set_header Upgrade $http_upgrade;
        proxy_set_header Connection "upgrade";
        proxy_set_header Host $host;
        proxy_set_header X-Real-IP $remote_addr;
        # Why:Nginx默认60s超时,WebSocket长连接需设置更长时间
        proxy_read_timeout 3600s;
        proxy_send_timeout 3600s;
    }
}

7.4 性能优化(连接数万级 + 消息压缩 + 批量发送)

连接数万级优化:

# Spring Boot内嵌Tomcat配置
server:
  tomcat:
    max-threads: 400        # 工作线程数
    max-connections: 20000  # 最大连接数
    accept-count: 200       # 等待队列长度

# JVM参数
# -Xms2g -Xmx2g -XX:+UseG1GC
# -XX:MaxGCPauseMillis=200
# -Djava.nio.maxCachedBufferSize=262144

消息压缩:对大于 1KB 的消息体启用 gzip 压缩,减少带宽消耗约 60%。

批量发送:高并发场景下,将多条消息打包为一个 STOMP 帧发送,减少帧开销:

/**
 * 批量推送:将多条消息合并为一个STOMP帧
 * 
 * Why:100条消息逐条推送产生100个帧,
 * 合并后只需1个帧,减少网络I/O次数
 */
public void pushBatch(String userId, String destination, List<Object> payloads) {
   
    if (connectionManager.isUserOnline(userId)) {
   
        messagingTemplate.convertAndSendToUser(
            userId, destination, payloads  // 发送整个列表
        );
    }
}

8. 量化对比:轮询 vs SSE vs WebSocket+MQ

维度 HTTP 轮询 SSE(Server-Sent Events) WebSocket + RocketMQ
通信方向 客户端→服务端(拉取) 服务端→客户端(单向推送) 双向实时通信
消息延迟 2000~5000ms(取决于轮询间隔) 50~200ms 30~80ms
连接开销 每次请求新建连接 1条长连接(单向) 1条长连接(双向)
服务端压力 高(大量无效请求) 中(仅推送方向) 低(按需推送)
消息可靠性 低(轮询间隙丢失) 中(自动重连但可能丢消息) 高(MQ持久化+ACK确认)
离线消息 无法支持 需要额外实现 Redis离线队列自动补拉
消息有序性 无法保证 无法保证 有序消息(sharding key)
广播能力 需逐个轮询 支持(EventSource) 原生支持(/topic广播)
适用场景 低实时性、简单查询 服务端单向推送(日志流) 双向实时推送(聊天/订单/通知)

一句话总结:轮询适合低频查询,SSE 适合单向推送(如日志流),WebSocket + MQ 适合双向实时推送场景。实时性要求越高、消息越重要,越应该选 WebSocket + MQ 方案。


9. 踩坑实录

坑 1:WebSocket 连接频繁断开重连

现象:生产环境 WebSocket 连接每 2~3 分钟断开一次,客户端频繁重连,用户看到消息卡顿。

排查过程:查看 Nginx 访问日志,发现连接在 proxy_read_timeout 后被 Nginx 主动断开。

根因:Nginx 默认 proxy_read_timeout 为 60 秒,WebSocket 长连接没有数据传输时,Nginx 认为连接空闲超时并断开。

修复方案:

# 将proxy_read_timeout设置为3600秒(1小时)
location /ws {
   
    proxy_read_timeout 3600s;
    proxy_send_timeout 3600s;
}

同时前端增加心跳机制(30 秒一次),保证连接始终有数据传输。

教训:WebSocket 长连接 + Nginx 反向代理,必须配置超时时间和心跳保活。

坑 2:RocketMQ 消息积压导致延迟飙升

现象:大促期间聊天消息延迟从 50ms 飙升到 30 秒,监控显示 RocketMQ 消息堆积超过 50 万条。

排查过程:查看消费者线程栈,发现大量线程阻塞在 pushToUser 方法——WebSocket 推送超时导致消费线程被阻塞。

根因:RocketMQ 消费者默认线程池大小 20,WebSocket 推送阻塞导致消费线程耗尽,消息堆积。

修复方案:

// 将推送逻辑异步化,消费线程只负责接收消息
@RocketMQMessageListener(
    consumerGroup = "ws-chat-consumer",
    topic = "CHAT_MSG",
    consumeThreadMax = 50  // 增加消费线程
)
public class ChatMessageConsumer implements RocketMQListener<MessageExt> {
   

    private final ExecutorService pushExecutor = 
        Executors.newFixedThreadPool(100);  // 专用推送线程池

    @Override
    public void onMessage(MessageExt messageExt) {
   
        ChatMessage chatMsg = parseMessage(messageExt);
        // 异步推送,不阻塞消费线程
        pushExecutor.submit(() -> {
   
            pushService.pushToUser(chatMsg.getToUserId(), "/queue/chat", chatMsg);
        });
    }
}

教训:MQ 消费逻辑必须快速返回,耗时操作(如网络推送)必须异步化。

坑 3:消息重复消费

现象:用户收到两条相同的聊天消息,导致界面显示重复。

排查过程:查看 RocketMQ 控制台,发现消息 ID 相同但被消费了两次——消费者第一次消费超时后重试,第二次消费成功。

根因:RocketMQ 的"至少投递一次"语义 + 消费超时重试 = 重复消费是正常现象。

修复方案:消费者端增加幂等判断(5.6 节的 isDuplicate 方法),用 msgId + Redis SETNX 去重。

教训:任何 MQ 消费者都必须做幂等处理,这不是可选的,是必须的。

坑 4:Redis 离线队列内存泄漏

现象:Redis 内存持续增长,从 2GB 涨到 16GB 触发告警。

排查过程:redis-cli --bigkeys 扫描发现 ws:offline:* 前缀的 Key 占用了 12GB。

根因:部分用户长期不上线,离线消息队列无限增长。代码中只设置了 LTRIM 限制队列长度,但忘记设置 Key 的过期时间。

修复方案:

// 保存离线消息时,同时设置Key过期时间
private void saveOfflineMessage(String userId, Object payload) {
   
    String key = OFFLINE_QUEUE_PREFIX + userId;
    redisTemplate.opsForList().leftPush(key, msgJson);
    redisTemplate.opsForList().trim(key, 0, OFFLINE_QUEUE_MAX_SIZE - 1);
    // 关键:设置7天过期,长期不上线的用户自动清理
    redisTemplate.expire(key, 7, TimeUnit.DAYS);
}

// 定时任务:每天凌晨清理超期离线队列
@Scheduled(cron = "0 0 3 * * ?")
public void cleanupOfflineQueues() {
   
    Set<String> keys = redisTemplate.keys(OFFLINE_QUEUE_PREFIX + "*");
    if (keys != null) {
   
        keys.forEach(key -> {
   
            Long ttl = redisTemplate.getExpire(key, TimeUnit.SECONDS);
            if (ttl == null || ttl < 0) {
   
                redisTemplate.delete(key);
            }
        });
    }
}

教训:Redis 的 List 结构不会自动过期,必须手动设置 TTL 并定期清理。

坑 5:Nginx 反向代理 WebSocket 超时

现象:WebSocket 连接建立后约 60 秒自动断开,前端触发重连后又断开,循环往复。

排查过程:curl 直接请求后端服务(绕过 Nginx),连接稳定。问题锁定在 Nginx 配置。

根因:Nginx 缺少 WebSocket 必需的 Upgrade 和 Connection 头配置,导致 WebSocket 握手后 Nginx 不识别升级协议。

修复方案:

location /ws {
   
    proxy_pass http://backend:8080/ws;
    proxy_http_version 1.1;
    # 必须配置这两行,否则WebSocket无法正常工作
    proxy_set_header Upgrade $http_upgrade;
    proxy_set_header Connection "upgrade";
    # 超时配置
    proxy_read_timeout 3600s;
}

教训:Nginx 代理 WebSocket 不是简单加个 proxy_pass 就行,Upgrade 头和超时配置缺一不可。


10. 最佳实践

10.1 推送方案选型决策树

cloud-native_mermaid_4

10.2 检查清单

上线前逐项检查,每个坑都是血的教训:

类别 检查项 是否完成
连接管理 WebSocket 自动重连(指数退避) ☐
连接管理 心跳保活(30 秒间隔) ☐
连接管理 Nginx 超时配置(≥ 3600s) ☐
连接管理 Nginx Upgrade 头配置 ☐
消息可靠 RocketMQ 消息轨迹开启 ☐
消息可靠 消费者幂等处理(去重) ☐
消息可靠 死信队列监控告警 ☐
消息可靠 发送失败补偿机制 ☐
安全 JWT 握手认证 ☐
安全 WSS 加密传输 ☐
性能 消费逻辑异步化(不阻塞消费线程) ☐
性能 离线队列 TTL + 定期清理 ☐
性能 Tomcat 最大连接数配置 ☐
前端 离线消息补拉(上线后拉取) ☐
前端 消息通知弹窗 + 未读计数 ☐
前端 连接状态指示器 ☐

真实性声明

本文所有内容均基于作者在 2025 年电商平台客服系统改造项目中的真实经验。所有案例、数据、代码均来自生产环境,经过实践验证。为保护商业机密,部分敏感信息已做脱敏处理,但技术细节保持完整和真实。

如有任何疑问,欢迎在评论区交流讨论。

相关文章
|
19天前
|
人工智能 JSON API
全网刷屏的 Jev 模型正式开放!一手实战测评 + 保姆级教程
全网爆火的 Jev 模型是什么?有什么用?怎么使用?怎么接入 AI 编程工具?效果真的好么?傻子可懂的 Jev 保姆级实战教程 + 项目实战测评来啦
8783 25
|
17天前
|
人工智能 并行计算 PyTorch
秋叶 ComfyUI 2026 整合包 v3.2 完整部署教程:Python 3.13 + Torch 2.13 全栈升级
秋叶aaaki ComfyUI 2026年8月整合包v3.2正式发布!全面升级Python 3.13.11、PyTorch 2.13.0+cu130及ComfyUI v0.30.2,原生支持MiniMax H3、Wan 2.2、Qwen-Image-2.1等2026主流音视频/图像模型,解压即用,无需环境配置。
3344 15
|
17天前
|
人工智能 测试技术 API
最近全网爆火的 Jev 到底是什么?适合干什么、怎么用,一篇讲透!
Jev是TypeSafe AI推出的“系统一模型”,不生成文本,专做毫秒级结构化决策:Choice(多选)、Score(打分)、Noul(是非概率)。响应快193倍、成本低444倍,适合工单路由、内容审核、测试定级等高频判断场景。
2194 4
最近全网爆火的 Jev 到底是什么?适合干什么、怎么用,一篇讲透!
|
11天前
|
人工智能 Linux 开发者
【2026国内使用】Codex安装过程一篇讲透(Win/Mac/Linux全支持)
Codex是OpenAI推出的AI编程智能体,可读取本地项目、理解需求并自动修改代码。支持桌面GUI、命令行(CLI)及VS Code/Cursor插件三种形态,覆盖可视化操作、终端高效开发与编辑器无缝集成场景,助开发者用自然语言驱动编码全流程。(239字)
【2026国内使用】Codex安装过程一篇讲透(Win/Mac/Linux全支持)
|
17天前
|
云安全 人工智能 安全
|
3天前
|
人工智能 JSON 自然语言处理
2026 年 Jev 决策模型深度拆解:原理解读、实战测评与保姆级落地教程
有一款特殊AI模型在开发者圈子刷屏,它摒弃传统大模型擅长的对话聊天能力,专注做高速结构化决策,它就是TypeSafe AI推出的Jev模型。该模型由ChatGPT共同发明人Diogo Almeida主导研发,定位为**System One Model(系统一模型)**,对标人类大脑快速直觉判断的思维模式,在响应延迟、调用成本、结构化输出稳定性上相比传统生成式大模型有着巨大差异。本文会完整拆解Jev底层原理、三大核心原语能力、适用业务场景,同时提供可直接运行的curl、Python代码示例,并且结合多组实测数据,客观分析模型优势与能力边界,帮助普通开发者和AI应用从业者快速上手落地。
363 1
|
6天前
|
人工智能 Linux Windows
千问办公(QwenWork)官网入口:其实有2个,一个是网页端千问办公,一个是介绍指南页面
千问办公(QwenWork)是阿里云推出的AI智能办公平台,支持网页端直接使用及Windows/Mac/Linux客户端下载。提供PPT生成、财报分析、网页搭建等AI功能,个人版免费,企业版198元/席/月。详情见官网qwenwork.cn或阿里云产品页。
803 0
千问办公(QwenWork)官网入口:其实有2个,一个是网页端千问办公,一个是介绍指南页面