绝大多数后端做实时推送、消息即时通讯,第一反应就是原生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. 离线消息推送
WebSocket:自行Redis/数据库缓存离线消息,用户上线手动查询推送;
MQTT:会话持久化+QoS1/2,broker自动缓存离线消息,重连自动下发。 - 2. 一对一私聊、群聊
WebSocket:手动维护用户-连接映射关系,分组缓存;
MQTT:依靠主题天然隔离,私聊使用user/{userId}/private,群聊使用group/{groupId},订阅即接收,无需维护映射表。 - 3. 消息可靠投递
WebSocket:无重传机制,网络波动直接丢消息;
MQTT三级QoS可控:最多一次、至少一次、恰好一次,金融、告警场景保证消息不丢失不重复。 - 4. 设备下线通知
WebSocket:监听连接断开事件,业务手动处理;
MQTT内置遗嘱消息,客户端异常离线自动发送预设通知,无需心跳兜底判断。 - 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 编辑过