LOGO 首页 OA教程 ERP教程 模切知识交流 PMS教程 CRM教程 技术文档 其他文档  
 
网站管理员

还在用WebSocket实现即时通讯?试试MQTT吧,真香!

admin
2026年7月30日 14:36 本文热度 286

绝大多数后端做实时推送、消息即时通讯,第一反应就是原生WebSocket。聊天对话、设备上报、消息通知、实时告警、IoT硬件上报,全部堆一套WebSocket长连接服务。但上线后接踵而来一堆棘手问题:单机连接上限低、集群同步复杂、消息堆积丢失、弱网重连丢消息、设备离线消息无留存、多端同步逻辑写满大量业务代码。

WebSocket只是单纯的长连接传输协议,只解决“双向通道”,不具备消息分发、持久化、QoS等级、离线缓存、主题订阅、负载均衡这些原生能力,所有能力都要自己从零封装。而MQTT是专为低带宽、弱网络、海量连接、实时消息设计的轻量级消息协议,天然适配即时通讯、物联网、消息推送场景,搭配EMQX、Mosquitto服务,开箱即用全套能力,大幅减少重复开发,线上稳定性远高于自研WebSocket。

一、WebSocket与MQTT底层差距,看懂为什么MQTT更适合实时业务

1.1 定位本质完全不同

WebSocket是传输层扩展协议,依托HTTP握手升级TCP长连接,仅提供双向数据流通道,没有任何消息语义定义。消息发送、订阅、重发、离线存储、消息过滤全部需要业务自行编码实现。
MQTT是应用层消息订阅协议,专门面向消息分发场景,内置一套完整消息规范:主题订阅、消息质量等级、离线消息、遗嘱消息、心跳保活、会话持久化,协议原生封装,无需业务重复开发。

1.2 海量连接与弱网适配差距巨大

原生WebSocket服务单机支撑万级连接就会出现内存、线程瓶颈;MQTT服务EMQX基于Erlang轻量进程模型,单机轻松支撑百万级并发长连接,内存占用极低,海量设备、多用户在线场景碾压WebSocket。
移动端、IoT设备弱网环境频繁断连是常态:WebSocket断连后所有未推送消息直接丢失,需要业务手动缓存重推;MQTT原生支持QoS1/QoS2消息持久化,断连重连后自动补发未接收消息,弱网场景无需自己写缓存逻辑。

1.3 集群扩容、消息同步成本天差地别

WebSocket集群天然存在“连接分片”问题:用户连接到A节点,消息推送到B节点,无法触达客户端,必须额外引入Redis发布订阅、消息中间件做跨节点同步,架构复杂度翻倍,还容易出现消息重复、漏推。
MQTT broker原生内置集群分发机制,同一主题消息自动同步至所有节点,无论客户端连在哪台服务器,订阅消息都能精准送达,集群扩容零额外开发成本。

1.4 业务场景原生能力对比

  1. 1. 离线消息推送
    WebSocket:自行Redis/数据库缓存离线消息,用户上线手动查询推送;
    MQTT:会话持久化+QoS1/2,broker自动缓存离线消息,重连自动下发。
  2. 2. 一对一私聊、群聊
    WebSocket:手动维护用户-连接映射关系,分组缓存;
    MQTT:依靠主题天然隔离,私聊使用user/{userId}/private,群聊使用group/{groupId},订阅即接收,无需维护映射表。
  3. 3. 消息可靠投递
    WebSocket:无重传机制,网络波动直接丢消息;
    MQTT三级QoS可控:最多一次、至少一次、恰好一次,金融、告警场景保证消息不丢失不重复。
  4. 4. 设备下线通知
    WebSocket:监听连接断开事件,业务手动处理;
    MQTT内置遗嘱消息,客户端异常离线自动发送预设通知,无需心跳兜底判断。
  5. 5. 带宽开销
    WebSocket自定义消息无规范,数据包冗余大;MQTT二进制轻量报文,头部极小,移动端、硬件设备流量消耗大幅降低。

二、MQTT核心关键能力详解

2.1 主题订阅模型(实现私聊、群聊、广播)

MQTT采用分层主题通配符设计,天然适配各类通讯场景,无需复杂关系存储:

  • • 个人私有消息:msg/user/{userId},仅当前用户订阅,实现一对一私聊;
  • • 群聊频道:msg/group/{groupId},群内所有用户订阅,群发消息;
  • • 全局系统广播:msg/system/all,全部在线用户接收公告;
  • • 分类通知:msg/order/{userId}msg/notice/{userId},按业务模块隔离消息。
    通配符+单层匹配、#多层匹配,灵活批量订阅多类消息。

2.2 QoS消息质量等级,解决消息丢失痛点

  • • QoS0(最多一次):消息发完即丢弃,不重试,适合实时性要求高、丢失无影响的普通在线通知;
  • • QoS1(至少一次):保证消息一定送达,可能重复接收,适合聊天消息、业务告警;
  • • QoS2(恰好一次):握手确认,只接收一次,适合订单、支付、设备上报等不能重复的核心数据。

2.3 会话持久化 + 离线消息缓存

客户端连接时开启cleanSession=false,broker会保存当前订阅关系与未接收消息。用户APP退出、断网、后台杀进程,所有未读消息全部缓存;下次重新建立连接,自动批量推送离线消息,完美实现“离线存消息,上线自动读”,不用业务层额外存储。

2.4 遗嘱消息(异常下线自动通知)

客户端建立连接时预先设置遗嘱主题与消息,当客户端无心跳、网络异常、崩溃离线,broker自动向遗嘱主题发送预设消息,服务端可监听该主题,实时感知用户离线状态,替代复杂的心跳检测逻辑。

2.5 心跳保活机制

连接时指定心跳间隔,客户端定时上报心跳,长时间无心跳broker判定离线,精准管理连接状态,相比WebSocket自定义心跳逻辑更标准、低开销。

三、本地快速部署EMQX MQTT服务

EMQX是国内最主流开源MQTT消息服务器,百万级并发、支持MQTT over WebSocket,前端网页可直接通过ws协议连接,前后端全链路打通。

3.1 Docker一键启动

docker run -d --name emqx -p 1883:1883 -p 8083:8083 -p 18083:18083 emqx/emqx:5.3.0

端口说明:
1883:标准MQTT TCP端口;
8083:MQTT over WebSocket(前端网页连接使用);
18083:后台管理面板,账号admin/public,可查看连接、主题、消息日志。

3.2 基础权限配置

登录后台管理,创建认证账号,限制客户端连接权限,避免匿名访问,生产环境必须关闭匿名连接。

四、SpringBoot整合MQTT完整实战(后端生产可用)

采用Eclipse Paho Java客户端,实现消息发布、订阅、离线消息接收、断线自动重连全套逻辑。

4.1 Maven依赖

<dependency>
    <groupId>org.eclipse.paho</groupId>
    <artifactId>org.eclipse.paho.client.mqttv3</artifactId>
    <version>1.2.5</version>
</dependency>

4.2 application.yml配置

mqtt:
  broker: tcp://127.0.0.1:1883
  username: admin
  password: public
  client-id: server-producer
  # 离线消息持久会话
  clean-session: false
  # 心跳30秒
  keep-alive: 30
  # 默认QoS1,保证消息至少送达一次
  default-qos: 1

4.3 MQTT客户端配置类,自动重连、会话持久

import org.eclipse.paho.client.mqttv3.*;
import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class MqttConfig {

    @Value("${mqtt.broker}")
    private String broker;
    @Value("${mqtt.username}")
    private String username;
    @Value("${mqtt.password}")
    private String password;
    @Value("${mqtt.client-id}")
    private String clientId;
    @Value("${mqtt.clean-session}")
    private boolean cleanSession;
    @Value("${mqtt.keep-alive}")
    private int keepAlive;

    @Bean(destroyMethod = "disconnect")
    public MqttClient mqttClient() throws MqttException {
        MemoryPersistence persistence = new MemoryPersistence();
        MqttConnectOptions options = new MqttConnectOptions();
        options.setUserName(username);
        options.setPassword(password.toCharArray());
        options.setCleanSession(cleanSession);
        options.setKeepAliveInterval(keepAlive);
        // 开启自动重连
        options.setAutomaticReconnect(true);

        MqttClient mqttClient = new MqttClient(broker, clientId, persistence);
        // 消息回调监听
        mqttClient.setCallback(new MqttMessageCallback());
        mqttClient.connect(options);
        return mqttClient;
    }
}

4.4 消息回调处理类,接收订阅消息

import org.eclipse.paho.client.mqttv3.IMqttDeliveryToken;
import org.eclipse.paho.client.mqttv3.MqttCallback;
import org.eclipse.paho.client.mqttv3.MqttMessage;
import java.nio.charset.StandardCharsets;

public class MqttMessageCallback implements MqttCallback {

    // 连接丢失,自动重连由客户端配置控制
    @Override
    public void connectionLost(Throwable cause) {
        System.err.println("MQTT连接断开,自动重连中:" + cause.getMessage());
    }

    // 收到订阅主题消息
    @Override
    public void messageArrived(String topic, MqttMessage message) throws Exception {
        String payload = new String(message.getPayload(), StandardCharsets.UTF_8);
        System.out.println("收到主题[" + topic + "]消息:" + payload);
        // 此处执行业务逻辑:聊天消息存储、通知推送、设备数据处理
    }

    // 消息发送完成回调
    @Override
    public void deliveryComplete(IMqttDeliveryToken token) {
    }
}

4.5 MQTT工具类:发布消息、订阅主题

import org.eclipse.paho.client.mqttv3.MqttClient;
import org.eclipse.paho.client.mqttv3.MqttException;
import org.eclipse.paho.client.mqttv3.MqttMessage;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;

@Component
public class MqttUtil {

    @Autowired
    private MqttClient mqttClient;

    /**
     * 发布消息到指定主题
     */

    public void publish(String topic, String content, int qos) throws MqttException {
        MqttMessage message = new MqttMessage();
        message.setPayload(content.getBytes());
        message.setQos(qos);
        mqttClient.publish(topic, message);
    }

    /**
     * 订阅指定主题
     */

    public void subscribe(String topic, int qos) throws MqttException {
        mqttClient.subscribe(topic, qos);
    }
}

4.6 Controller业务调用示例

import org.eclipse.paho.client.mqttv3.MqttException;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;

@RestController
public class MsgController {

    @Autowired
    private MqttUtil mqttUtil;

    // 发送私聊消息
    @GetMapping("/send/private")
    public String sendPrivateMsg(@RequestParam Long userId, @RequestParam String content) throws MqttException {
        String topic = "msg/user/" + userId;
        mqttUtil.publish(topic, content, 1);
        return "私聊消息发送成功";
    }

    // 发送群聊消息
    @GetMapping("/send/group")
    public String sendGroupMsg(@RequestParam Long groupId, @RequestParam String content) throws MqttException {
        String topic = "msg/group/" + groupId;
        mqttUtil.publish(topic, content, 1);
        return "群聊消息发送成功";
    }
}

五、前端网页对接MQTT(MQTT over WebSocket)

前端无需改造WebSocket业务架构,直接通过mqttws3.js连接EMQX的8083端口,实现网页实时接收消息,兼容原有网页即时通讯场景。

<script src="https://cdn.jsdelivr.net/npm/mqtt@3/dist/mqtt.min.js"></script>
<script>
    // 连接ws端口
    const client = mqtt.connect('ws://127.0.0.1:8083/mqtt');
    // 登录用户ID
    const userId = 10001;
    const privateTopic = `msg/user/${userId}`;

    // 连接成功订阅个人私聊主题
    client.on('connect'() => {
        client.subscribe(privateTopic);
    });

    // 接收推送消息
    client.on('message'(topic, payload) => {
        const msg = payload.toString();
        console.log('收到实时消息:', msg);
        // 页面渲染聊天、通知弹窗
    });
</script>

六、全文总结

WebSocket只是底层双向通道工具,所有实时通讯配套能力都需要业务从零开发;MQTT是成熟的消息实时通讯标准协议,百万级并发、弱网适配、离线消息、集群分发、消息可靠投递全部原生支持,后端仅需少量代码完成发布订阅,大幅降低开发、运维、线上故障成本。
传统即时通讯项目如果存在连接量大、移动端多端、离线消息、集群扩容需求,替换为MQTT架构后,稳定性、扩展性会得到质的提升,真正做到开箱即用,不用重复造轮子。


即时通讯、长连接消息推送是后端高频业务场景,很多团队还在自研WebSocket踩坑。后续持续分享EMQX集群高可用、MQTT消息持久化、分布式聊天系统完整架构、WebSocket改造MQTT迁移方案全套实战干货,附带完整可运行源码。


阅读原文:点击这里


该文章在 2026/7/30 14:36:17 编辑过
关键字查询
相关文章
正在查询...
点晴ERP是一款针对中小制造业的专业生产管理软件系统,系统成熟度和易用性得到了国内大量中小企业的青睐。
点晴PMS码头管理系统主要针对港口码头集装箱与散货日常运作、调度、堆场、车队、财务费用、相关报表等业务管理,结合码头的业务特点,围绕调度、堆场作业而开发的。集技术的先进性、管理的有效性于一体,是物流码头及其他港口类企业的高效ERP管理信息系统。
点晴WMS仓储管理系统提供了货物产品管理,销售管理,采购管理,仓储管理,仓库管理,保质期管理,货位管理,库位管理,生产管理,WMS管理系统,标签打印,条形码,二维码管理,批号管理软件。
点晴免费OA是一款软件和通用服务都免费,不限功能、不限时间、不限用户的免费OA协同办公管理系统。
Copyright 2010-2026 ClickSun All Rights Reserved  粤ICP备13012886号-1  粤公网安备44030602007207号