MQTT 这个协议,我第一次接触是在做一个远程环境监测的小项目。当时的需求很朴素:几十个分布在城郊不同位置的采集节点,要把温湿度、PM2.5 这些数据实时传回中心服务器,同时中心还能反向下发一些控制指令。最开始想用 HTTP 轮询,写了两天就发现不对劲——设备端耗电快、服务器压力大、实时性还差。后来换成 MQTT,整个链路一下子清爽了。这篇文章就把我从零搭 MQTT 到跑通订阅发布、再到踩坑排错的完整过程拆开讲,适合刚接触物联网协议、想快速把 MQTT 用起来的开发者,也适合已经会连 Broker 但没搞明白背后机制的朋友。
1. 为什么物联网场景里 MQTT 比 HTTP 更合适
1.1 从一次轮询翻车说起
先讲个真实的翻车现场。我最早那版采集程序,设备端每 5 秒发一次 HTTP GET,把数据塞进 URL 参数里。单台设备测试没问题,等到 50 台设备同时上线,服务器那边 Nginx 的并发连接数直接飙到几百,CPU 占用居高不下。更麻烦的是,很多设备用的是电池供电,HTTP 每次请求都要重新建立 TCP 连接(就算开了 keep-alive,超时后还是要重连),电量掉得肉眼可见。
这个问题的本质在于:HTTP 是"请求-响应"模型,客户端不主动问,服务器就没法主动推。设备想知道有没有新指令,只能不停地问。而物联网场景里,绝大多数时间是没有指令的,这些轮询全是无效开销。
MQTT 换了个思路。它基于**发布/订阅(Publish/Subscribe)**模型,设备连上 Broker 之后,只需要订阅自己关心的主题(Topic),有消息时 Broker 主动推过来。设备平时保持一条长连接就行,不用反复问。这一下子把无效请求砍掉了,功耗和服务器压力都降下来了。
1.2 发布订阅模型到底解决了什么问题
用一个生活化的类比。HTTP 轮询像是你每隔五分钟给快递站打个电话问"我的包裹到了吗",打一百次可能九十九次都是"还没到"。MQTT 则像是你留了个手机号给快递站,包裹一到,快递员直接给你打电话。
技术上,这个模型把通信双方解耦了:
- 空间解耦:发布者不需要知道订阅者是谁,只管往某个 Topic 发消息。
- 时间解耦:发布者发消息时,订阅者可以不在线(取决于 QoS 和是否保留消息)。
- 同步解耦:双方不需要同时运行,也不需要互相等待。
这三点对物联网特别关键。传感器只管上报数据,根本不用关心是谁在消费;后台服务可以随时重启,重启期间的消息靠 Broker 缓存或保留机制兜底。
1.3 MQTT 的几个核心概念先理清楚
在动手之前,有几个词必须先搞明白,不然后面配置会一头雾水。
Broker(代理服务器):消息的中转站,所有客户端都连它。常见的有 EMQX、Mosquitto、HiveMQ 等。它负责接收发布者的消息,再按订阅关系转发给订阅者。
Client(客户端):任何连到 Broker 的设备或程序,既可以是发布者,也可以是订阅者,还可以两者都是。
Topic(主题):消息的分类标签,用斜杠分层,比如home/livingroom/temperature。订阅时可以用通配符:+匹配单层,#匹配多层。比如home/+/temperature能匹配home/livingroom/temperature和home/bedroom/temperature。
QoS(服务质量等级):分 0、1、2 三档。QoS 0 是"发出去就不管",可能丢;QoS 1 是"至少送达一次",可能重复;QoS 2 是"恰好送达一次",开销最大。选哪档要看业务能容忍丢还是能容忍重。
Retained Message(保留消息):Broker 会为某个 Topic 保存最后一条保留消息,新订阅者一订阅就能立刻收到。这个对"设备上线后想知道当前状态"的场景特别有用。
Will Message(遗嘱消息):客户端异常断开时,Broker 代为发布的一条消息。常用来做设备离线告警。
把这几个概念串起来,MQTT 的工作流就清楚了:客户端连 Broker,订阅 Topic,发布者往 Topic 发消息,Broker 按 QoS 转发,订阅者收到消息。
2. 本地把 Broker 跑起来:选型与安装的取舍
2.1 Broker 选哪个:Mosquitto 还是 EMQX
这一步很多人纠结。我的建议是分场景:
| Broker | 适用场景 | 优点 | 注意点 |
|---|---|---|---|
| Mosquitto | 本地开发、小规模部署 | 轻量、安装快、配置简单 | 集群和高并发能力弱 |
| EMQX | 生产环境、大规模设备 | 高并发、支持集群、有管理面板 | 资源占用相对高 |
| HiveMQ | 企业级、商业支持 | 生态完善、插件丰富 | 社区版功能受限 |
我本地开发一般用 Mosquitto,因为它足够轻,装完改两行配置就能跑。等要压测或者上生产,再换 EMQX。这个思路的好处是:开发阶段不被复杂配置干扰,专注把业务逻辑跑通。
2.2 Windows 下安装 Mosquitto 的完整步骤
网上很多教程只给个下载链接就完事,实际装的时候坑不少。我把完整流程写清楚。
第一步,去 Mosquitto 官网下载 Windows 安装包(.exe)。注意选对版本,64 位系统选 x64。
第二步,安装时它会问要不要装服务,勾上。装完后默认路径一般在C:\Program Files\mosquitto。
第三步,关键来了——默认配置只监听本地回环地址,别的机器连不上。打开安装目录下的mosquitto.conf,找到listener相关配置,改成:
listener 1883 0.0.0.0 allow_anonymous truelistener 1883 0.0.0.0表示监听所有网卡的 1883 端口,allow_anonymous true表示允许匿名连接。注意:这两项只适合本地开发,生产环境必须关掉匿名并配认证,否则等于把门敞开。
第四步,重启服务。用管理员权限打开命令行:
net stop mosquitto net start mosquitto如果启动失败,多半是端口被占用或者配置文件语法错误。可以先用mosquitto -c mosquitto.conf -v前台启动,看详细日志。
2.3 验证 Broker 是否真的在工作
装完别急着写代码,先用命令行工具验证一下。Mosquitto 自带mosquitto_sub和mosquitto_pub两个工具。
开一个终端订阅:
mosquitto_sub -h localhost -p 1883 -t "test/topic" -v再开一个终端发布:
mosquitto_pub -h localhost -p 1883 -t "test/topic" -m "hello mqtt"订阅端如果打印出test/topic hello mqtt,说明 Broker 工作正常。这一步看着简单,但能帮你排除掉一大半"代码连不上"的问题——先确认 Broker 没问题,再去怀疑代码。
提示:如果订阅端收不到消息,先检查两个终端的 Topic 是否完全一致(大小写敏感),再检查是否连的同一个 Broker 地址和端口。
3. 用 Java 把订阅发布跑通:从依赖到可运行代码
3.1 客户端库怎么选
Java 生态里 MQTT 客户端主流有两个:Eclipse Paho 和 HiveMQ MQTT Client。Paho 是老牌选手,资料多、稳定;HiveMQ 的客户端 API 更现代,异步支持更好。
我选 Paho,原因是它在各种老项目里兼容性好,遇到问题搜到的答案也多。Maven 依赖:
<dependency> <groupId>org.eclipse.paho</groupId> <artifactId>org.eclipse.paho.client.mqttv3</artifactId> <version>1.2.5</version> </dependency>版本别乱选,1.2.5 是相对稳定的版本。有些新版本在特定 JDK 上会有兼容问题,踩过。
3.2 发布端代码:把一条消息发出去
先看发布端。核心就三步:建客户端、连 Broker、发消息。
import org.eclipse.paho.client.mqttv3.*; import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence; public class Publisher { public static void main(String[] args) throws MqttException { String broker = "tcp://localhost:1883"; String clientId = "java-publisher-001"; MemoryPersistence persistence = new MemoryPersistence(); MqttClient client = new MqttClient(broker, clientId, persistence); MqttConnectOptions options = new MqttConnectOptions(); options.setCleanSession(true); options.setConnectionTimeout(10); options.setKeepAliveInterval(20); client.connect(options); String topic = "sensor/temperature"; String content = "{\"deviceId\":\"dev-01\",\"value\":26.5}"; MqttMessage message = new MqttMessage(content.getBytes()); message.setQos(1); client.publish(topic, message); client.disconnect(); client.close(); } }几个参数值得说清楚:
cleanSession(true):表示不保留会话状态。每次连接都是全新的,之前的订阅和未收消息都清掉。开发阶段用 true 省事,生产环境如果要保证离线消息不丢,得设 false。connectionTimeout(10):连接超时 10 秒。网络差的环境可以适当调大。keepAliveInterval(20):心跳间隔 20 秒。客户端会在这个周期内没发消息时发心跳包,Broker 靠它判断客户端是否还活着。设太小费流量,设太大断线发现慢。
3.3 订阅端代码:把消息收回来
订阅端稍微复杂一点,因为要处理回调。
import org.eclipse.paho.client.mqttv3.*; import org.eclipse.paho.client.mqttv3.persist.MemoryPersistence; public class Subscriber { public static void main(String[] args) throws MqttException { String broker = "tcp://localhost:1883"; String clientId = "java-subscriber-001"; MemoryPersistence persistence = new MemoryPersistence(); MqttClient client = new MqttClient(broker, clientId, persistence); MqttConnectOptions options = new MqttConnectOptions(); options.setCleanSession(true); options.setAutomaticReconnect(true); client.setCallback(new MqttCallback() { @Override public void connectionLost(Throwable cause) { System.out.println("连接断开: " + cause.getMessage()); } @Override public void messageArrived(String topic, MqttMessage message) { System.out.println("收到消息 [" + topic + "]: " + new String(message.getPayload())); } @Override public void deliveryComplete(IMqttDeliveryToken token) { // 订阅端一般用不到 } }); client.connect(options); client.subscribe("sensor/#", 1); // 保持主线程不退出 try { Thread.sleep(Long.MAX_VALUE); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }这里有个新手常犯的错:main方法里subscribe之后直接结束了,程序退出,自然收不到消息。必须让主线程保持存活,或者用CountDownLatch之类的机制阻塞。
setAutomaticReconnect(true)是 Paho 提供的自动重连,网络抖动时很管用。但要注意,自动重连后订阅关系不一定自动恢复,取决于cleanSession的设置。如果设了cleanSession(false),Broker 会记住订阅关系;设了 true,重连后得重新订阅。
3.4 一个容易忽略的细节:ClientId 必须唯一
Paho 要求 ClientId 唯一。如果两个客户端用同一个 ClientId 连同一个 Broker,后连的会把先连的踢下线。我早期调试时,本地开了两个订阅端用了默认生成的 ClientId,结果互相踢,消息时有时无,排查了半天。
生产环境建议用有意义的 ClientId,比如设备类型-设备ID-随机后缀,既方便排查,又避免冲突。
4. QoS、保留消息与遗嘱:把可靠性真正落地
4.1 QoS 三档到底怎么选
QoS 是 MQTT 可靠性的核心,但很多人要么全用 0,要么全用 2,都不对。我的经验是按数据重要性分档:
- QoS 0:适合高频、可容忍丢失的数据。比如每秒上报一次的传感器读数,丢一两条无所谓。
- QoS 1:适合大多数业务数据。比如设备状态变更、订单事件,允许偶尔重复,但业务侧要做幂等。
- QoS 2:适合绝对不能丢也不能重的场景。比如计费、支付指令。开销大,别滥用。
QoS 1 的"至少一次"意味着可能重复。我做过一个设备控制项目,下发"开阀"指令用了 QoS 1,结果网络抖动时设备收到两次,阀门开了又开。后来在指令里加了唯一 ID,设备端做去重才解决。用 QoS 1 就必须考虑幂等。
4.2 保留消息解决"上线即知状态"
场景是这样的:一个监控大屏,启动后想立刻显示所有设备的最新状态。如果不用保留消息,大屏得等设备下一次上报才有数据,可能要等几十秒。
保留消息的用法很简单,发布时设retained = true:
MqttMessage message = new MqttMessage(content.getBytes()); message.setQos(1); message.setRetained(true); client.publish(topic, message);Broker 会为这个 Topic 保存最后一条保留消息。之后任何新订阅者订阅这个 Topic,都会立刻收到它。
但有个坑:保留消息会一直存在 Broker 上,直到被新的保留消息覆盖,或者发布一条空消息清除。如果设备下线了,它的最后一条保留消息还在,新订阅者会收到过时数据。所以设备离线时,最好主动发一条空保留消息清掉,或者用遗嘱消息标记离线。
4.3 遗嘱消息做离线告警
遗嘱消息在连接时配置:
MqttConnectOptions options = new MqttConnectOptions(); options.setWill("device/status/dev-01", "offline".getBytes(), 1, true);意思是:如果这个客户端异常断开(不是主动 disconnect),Broker 就替它往device/status/dev-01发一条offline的保留消息。
配合设备上线时主动发online,就能实现设备在线状态监控。注意遗嘱消息只在异常断开时触发,主动disconnect()不会触发。这个区别很关键,我见过有人主动断开后纳闷为什么没收到离线消息。
5. 踩坑实录:那些文档里不会写的排查过程
5.1 连不上 Broker 的排查链路
现象:Java 客户端报Connection refused或Timed out。
我的排查顺序是这样的:
- 先确认 Broker 在跑:命令行
mosquitto_sub能不能连上。连不上就是 Broker 的问题。 - 确认监听地址:Broker 是不是只监听了
127.0.0.1。如果是,远程连必然失败,改成0.0.0.0。 - 确认端口和防火墙:1883 端口有没有被防火墙拦。Windows 上经常是防火墙默认拦截。
- 确认协议前缀:Paho 的 broker 地址要带协议,
tcp://或ssl://。漏了会报奇怪的错。 - 确认 ClientId 唯一:被踢下线时表现也可能是连接异常。
这个顺序的逻辑是:从服务端往客户端查,先排除最外层的问题。很多人一上来就怀疑代码,结果绕一大圈发现是防火墙。
5.2 消息收不到的几个典型原因
现象:订阅端连着,但收不到消息。
- Topic 不匹配:发布用
sensor/temp,订阅用sensor/temperature,差一个字母就收不到。通配符用错也常见,sensor/+只匹配一层,sensor/#匹配多层。 - QoS 不匹配:订阅 QoS 低于发布 QoS 时,实际生效的是两者中较低的那个,但不会导致收不到,只会影响可靠性。
- cleanSession 导致订阅丢失:设了
cleanSession(true),重连后订阅关系没了,自然收不到。要么重连后重新订阅,要么设 false。 - 回调没设:忘了
setCallback,消息到了也没人处理。
5.3 内存和连接泄漏
Paho 的MqttClient用完要close(),否则底层线程和连接不释放。我在一个批量测试脚本里忘了关,跑了几百次之后 JVM 报内存溢出。后来改成 try-with-resources 或者 finally 里 close 才解决。
另外,MqttClient是同步阻塞的,高并发场景建议用MqttAsyncClient,避免一个慢操作卡住整个线程。
6. 从能跑到好用:几个进阶优化方向
6.1 主题设计要有层次
Topic 设计得好,后期扩展省事。我的习惯是按业务域/设备类型/设备ID/数据项分层,比如factory/plc/plc-001/temperature。这样订阅时可以用通配符灵活组合:订阅所有 PLC 温度用factory/plc/+/temperature,订阅某台设备所有数据用factory/plc/plc-001/#。
别用太扁平的主题,比如temp001、temp002,后期想按类型订阅就抓瞎了。
6.2 认证与安全别偷懒
本地开发用匿名没问题,上线必须配认证。Mosquitto 支持用户名密码,配置password_file即可。EMQX 还支持更细粒度的 ACL,控制哪个客户端能发布/订阅哪些 Topic。
传输层如果走公网,建议上 TLS,把tcp://换成ssl://,配好证书。这一步不做,数据等于裸奔。
6.3 大规模设备的连接管理
设备数量上千后,Broker 的连接数、内存、文件句柄都会成为瓶颈。这时候要考虑:
- Broker 集群部署,EMQX 原生支持。
- 客户端心跳间隔调大,减少无效流量。
- 用共享订阅(Shared Subscription)做消费端负载均衡,多个后台服务订阅同一主题,Broker 轮流投递。
共享订阅的写法是在主题前加$share/组名/,比如$share/group1/sensor/#。这个特性在 EMQX 上支持得很好,做水平扩展时非常有用。
6.4 和 485 设备打交道的实际做法
热词里有人问 MQTT 怎么给 485 设备发指令。实际链路是这样的:MQTT 客户端收到指令后,通过串口(RS485)把 Modbus 报文发给设备,再把设备返回的数据通过 MQTT 上报。中间需要一个网关程序做协议转换。
关键点是:485 是半双工、一问一答的,而 MQTT 是异步的。网关里要维护一个请求队列,保证同一时刻只有一个 485 请求在途,否则会串数据。这个坑我在实际项目里踩过,两个指令同时下发,设备返回的数据对不上号,排查了很久才发现是并发问题。
7. 我个人的几点实操体会
MQTT 上手快,但要用好,核心在于理解它的模型而不是死记 API。发布订阅、QoS、保留消息、遗嘱这几个概念吃透了,剩下的都是配置问题。
调试阶段,命令行工具mosquitto_sub和mosquitto_pub比写代码快得多,遇到问题先用它们验证 Broker,能省大量时间。生产环境则一定要把认证、TLS、幂等、重连这些补上,别拿开发配置直接上线。
最后分享一个小技巧:调试时把订阅端的通配符设成#,能收到所有消息,快速确认消息到底有没有发出来、发到了哪个 Topic。等确认链路通了,再收窄到具体主题。这个笨办法在排查"消息去哪了"的时候特别有效。