1. 为什么物联网项目都绕不开 MQTT
搞过物联网项目的兄弟应该都有体会,设备端和云端之间的通信协议选型,基本决定了整个项目的开发效率和后期维护成本。我最早做设备联网的时候用过 HTTP 轮询,那会儿设备少还没觉得有什么问题,后来设备数量一上来,服务器直接被轮询请求打爆,带宽费用也扛不住。后来换成 MQTT,才算真正找到了适合物联网场景的通信方式。
MQTT 全称 Message Queuing Telemetry Transport,翻译过来叫消息队列遥测传输协议。名字听着挺唬人,但核心逻辑其实特别简单:它就是一个基于发布/订阅模式的消息传输协议,专门为低带宽、不稳定网络环境下的设备通信设计的。你可以把它想象成一个邮局系统——设备把消息投递到某个“信箱”(主题),其他设备或者服务端只要订阅了这个信箱,就能收到消息。发消息的人不需要知道谁在收,收消息的人也不需要知道谁在发,双方完全解耦。
这个特性在物联网场景里太重要了。因为物联网项目里设备种类多、数量大、网络环境复杂,如果用传统的请求-响应模式,每台设备都要知道服务端的地址,服务端也要维护所有设备的连接状态,耦合度太高。而 MQTT 的发布/订阅模型天然适合这种多对多的通信场景,设备只管往对应主题发数据,后端服务只管订阅自己关心的主题,中间通过 MQTT Broker 来转发消息,架构一下子就清晰了。
我写这篇东西的目的很明确:把我这些年用 MQTT 做物联网项目的经验整理出来,从协议核心概念到服务器搭建,从 Java 客户端开发到实际对接 485 设备,把整个链路讲透。不管你是刚接触物联网开发的新手,还是想从 HTTP 切换到 MQTT 的老手,看完应该都能直接上手干活。文章里会涉及 MQTT 协议详解、MQTT 服务器搭建、Java 快速开发框架选型、MQTT 订阅与发布消息的实操、以及 MQTT 如何给 485 设备发指令读取数据这些实际问题的解决方案。
2. MQTT 协议核心概念拆解
2.1 发布订阅模式到底怎么运转的
MQTT 的发布/订阅模式里有三个核心角色:Publisher(发布者)、Subscriber(订阅者)和 Broker(代理服务器)。发布者负责往某个主题发消息,订阅者负责订阅自己感兴趣的主题来接收消息,Broker 则是中间人,负责接收所有消息并根据订阅关系进行转发。
这个模式最大的好处是空间解耦和时间解耦。空间解耦的意思是发布者和订阅者互相不需要知道对方的存在,也不需要知道对方的 IP 地址和端口,它们只跟 Broker 打交道。时间解耦的意思是消息的发送和接收不需要同时进行,发布者发完消息就可以干别的去了,订阅者什么时候上线什么时候收,只要订阅关系还在,消息就不会丢(当然这取决于 QoS 等级)。
我举个实际场景你就明白了。假设你有一个温度传感器,它每隔 5 秒往sensor/temperature/room1这个主题发一次当前温度。同时你有一个数据存储服务订阅了这个主题,还有一个告警服务也订阅了这个主题。温度传感器只管发,它根本不知道有几个服务在消费它的数据。后面你如果想加一个实时大屏展示服务,只需要让它订阅同一个主题就行,完全不用动传感器端的代码。这种扩展性在传统请求-响应模式里是很难做到的。
2.2 主题与通配符的匹配规则
MQTT 的主题(Topic)是一个用斜杠分隔的字符串,比如home/livingroom/temperature。主题本身不需要预先创建,发布者往哪个主题发消息,这个主题就存在了。订阅者可以用通配符来批量订阅多个主题,通配符有两种:单层通配符+和多层通配符#。
单层通配符+只能匹配一个层级。比如你订阅home/+/temperature,那么home/livingroom/temperature和home/bedroom/temperature都能匹配到,但home/livingroom/sensor1/temperature就匹配不到,因为多了一层。多层通配符#可以匹配任意多个层级,比如你订阅home/#,那么home/livingroom/temperature、home/bedroom/humidity、home/kitchen/sensor/status全都能匹配到。
这里有个坑我踩过:#必须放在主题的最后,home/#/temperature这种写法是非法的。另外+和#可以组合使用,比如home/+/sensor/#,但同样#必须在末尾。还有一点要注意,主题是大小写敏感的,Home/Temperature和home/temperature是两个完全不同的主题,这个在开发的时候特别容易搞混。
2.3 QoS 等级怎么选才不踩坑
MQTT 定义了三个 QoS(Quality of Service)等级,用来保证消息传递的可靠性。QoS 0 是最多发一次,消息发出去就不管了,可能丢也可能重复。QoS 1 是至少发一次,保证消息能到达,但可能重复。QoS 2 是恰好发一次,保证消息不丢也不重复,但开销最大。
很多新手一上来就选 QoS 2,觉得最可靠。但实际上 QoS 2 的握手过程需要四次交互,对设备和网络的负担都很大。我一般的建议是:普通传感器数据用 QoS 0 就够了,丢一两个数据点对整体趋势没影响;控制指令用 QoS 1,保证指令能到达,重复执行的问题可以在应用层做幂等处理;只有涉及计费、安全等绝对不能出错的场景才用 QoS 2。
还有一个容易忽略的点:QoS 等级是在发布和订阅两端分别设置的,最终生效的是两者中较低的那个。比如你发布用 QoS 2,订阅用 QoS 0,那实际传输就是 QoS 0。这个机制在设计的时候要特别注意,别以为发布端设了高 QoS 就万事大吉了。
2.4 会话保持与遗嘱消息的实战价值
MQTT 的会话保持机制(Clean Session / Persistent Session)是个很实用的功能。当 Clean Session 设为 false 时,Broker 会为客户端保存订阅关系和未确认的消息。客户端断线重连后,能收到断线期间积压的消息。这个功能对于网络不稳定的设备特别有用,比如移动网络下的设备,断线是常态,有了会话保持就不会丢数据。
遗嘱消息(Will Message)是另一个我觉得特别巧妙的设计。客户端在连接 Broker 的时候可以设置一条遗嘱消息,当客户端异常断开时,Broker 会自动把这条消息发布到指定的主题。这个机制可以用来做设备离线告警——设备上线时设置遗嘱消息为“offline”,正常运行时定期发心跳,如果设备异常掉线,订阅了遗嘱主题的服务就能立刻收到离线通知。
我做过一个项目,设备端设置了遗嘱消息,后端服务订阅了遗嘱主题。有一次现场一台设备被拔了电源,后端在几秒内就收到了离线告警,运维人员及时处理了问题。如果没有遗嘱消息,可能要靠心跳超时来判断,延迟会大很多。
3. MQTT 服务器搭建与选型
3.1 主流 Broker 对比与选择建议
MQTT Broker 是整套系统的核心,选型的时候要考虑性能、稳定性、功能完整度和运维成本。我用过几个主流的 Broker,这里做个对比。
| Broker 名称 | 开发语言 | 协议支持 | 集群能力 | 适用场景 |
|---|---|---|---|---|
| Mosquitto | C | MQTT 3.1/3.1.1/5.0 | 弱 | 小型项目、测试环境 |
| EMQX | Erlang | MQTT 3.1/3.1.1/5.0 | 强 | 中大型生产环境 |
| HiveMQ | Java | MQTT 3.1/3.1.1/5.0 | 强 | 企业级、商业授权 |
| VerneMQ | Erlang | MQTT 3.1/3.1.1/5.0 | 强 | 大规模分布式 |
Mosquitto 是最轻量的选择,安装包小,资源占用低,适合在树莓派或者低配服务器上跑。但它的集群能力比较弱,设备量大了之后单机扛不住。EMQX 是我目前用得最多的,开源版功能就很完整,支持百万级连接,集群部署也方便,中文文档齐全,社区活跃。HiveMQ 功能强大但商业版收费不便宜,小团队慎选。VerneMQ 性能很好,但文档和社区相对弱一些。
如果你刚开始做物联网项目,设备量在几千以内,我建议直接用 EMQX 开源版,功能足够,后期扩展也方便。如果只是本地测试或者学习用,Mosquitto 就够了,装起来快。
3.2 Windows 下快速搭建 MQTT 服务
很多兄弟开发环境是 Windows,这里说下 Windows 下怎么快速把 MQTT 服务器跑起来。最简单的方式是用 EMQX 的 Windows 安装包,去官网下载解压后,进入 bin 目录,执行:
emqx start服务就起来了。默认的 MQTT 端口是 1883,Web 管理界面端口是 18083,浏览器打开http://localhost:18083,默认账号 admin,密码 public。进去之后可以直观地看到连接数、主题数、消息吞吐量这些指标,调试的时候很方便。
如果你用 Mosquitto,Windows 下需要下载安装包,安装完成后需要手动配置。默认配置文件在安装目录下,需要添加监听端口和认证配置。Mosquitto 默认只允许本地连接,要允许外部设备连接,需要在配置文件里加上:
listener 1883 0.0.0.0 allow_anonymous true改完配置后重启服务。不过生产环境千万别开allow_anonymous,一定要配认证,不然谁都能连上来发消息。
3.3 生产环境的安全配置要点
生产环境的 MQTT 服务器绝对不能裸奔。我见过不少项目 MQTT 端口直接暴露在公网上,没有任何认证,这跟把数据库密码写在公网上没区别。基本的安全配置包括:启用用户名密码认证、启用 TLS 加密、配置 ACL 权限控制。
用户名密码认证是最基础的,EMQX 支持内置数据库认证,也支持对接外部 MySQL、Redis 等。ACL 权限控制能限制每个客户端只能发布和订阅特定主题,比如设备 A 只能往device/A/#发消息,不能往device/B/#发。这个在多租户或者多设备场景下特别重要,能防止设备越权操作。
TLS 加密这块,如果设备端资源允许,建议开启。MQTT over TLS 默认端口是 8883。证书可以用自签的,也可以买正式的。自签证书需要在设备端预置 CA 证书,稍微麻烦一点,但安全性有保障。如果设备端实在跑不动 TLS,至少也要保证内网隔离,别把 MQTT 端口直接暴露出去。
4. Java 客户端快速开发实战
4.1 客户端库选型与项目初始化
Java 生态里 MQTT 客户端库主要有 Eclipse Paho 和 HiveMQ MQTT Client。Paho 是最老牌的,稳定但 API 偏底层,用起来稍微繁琐。HiveMQ 的客户端库 API 更现代,支持响应式编程,用起来舒服很多。我目前新项目基本都用 HiveMQ Client。
用 Maven 的话,在pom.xml里加依赖:
<dependency> <groupId>com.hivemq</groupId> <artifactId>hivemq-mqtt-client</artifactId> <version>1.3.3</version> </dependency>如果项目里已经有 Spring Boot,也可以用 Spring Integration MQTT,它封装了 Paho,跟 Spring 生态集成得更好。但如果你想要更灵活的控制,还是建议直接用 HiveMQ Client。
4.2 连接 Broker 与断线重连策略
建立连接是第一步,但断线重连才是实际项目里最需要关注的部分。网络抖动、Broker 重启、设备休眠都会导致连接断开,如果没有自动重连机制,设备就失联了。
用 HiveMQ Client 建立连接的代码大概长这样:
MqttClient client = MqttClient.builder() .useMqttVersion3() .identifier("device-" + deviceId) .serverHost("your-broker-host") .serverPort(1883) .automaticReconnectWithDefaultConfig() .buildAsync(); client.connectWith() .cleanSession(false) .keepAlive(60) .willPublish() .topic("device/" + deviceId + "/status") .payload("offline".getBytes()) .qos(MqttQos.AT_LEAST_ONCE) .applyWillPublish() .send() .whenComplete((ack, throwable) -> { if (throwable != null) { log.error("连接失败", throwable); } else { log.info("连接成功"); } });这里有几个关键参数:cleanSession(false)开启会话保持,keepAlive(60)设置心跳间隔为 60 秒,automaticReconnectWithDefaultConfig()开启自动重连。遗嘱消息设置成offline,设备异常断开时 Broker 会自动发布这条消息。
自动重连的默认配置是初始延迟 1 秒,最大延迟 120 秒,指数退避。这个策略在大多数场景下够用了。但如果你的设备对实时性要求很高,可以自定义重连策略,比如固定 5 秒重连一次。
4.3 消息发布与订阅的完整实现
发布消息的代码很直接:
client.publishWith() .topic("sensor/temperature/room1") .payload(String.valueOf(temperature).getBytes()) .qos(MqttQos.AT_MOST_ONCE) .send();订阅消息稍微复杂一点,需要设置回调:
client.subscribeWith() .topicFilter("device/+/command") .qos(MqttQos.AT_LEAST_ONCE) .callback(publish -> { String topic = publish.getTopic().toString(); String payload = new String(publish.getPayloadAsBytes()); log.info("收到消息 topic={}, payload={}", topic, payload); // 处理指令 handleCommand(topic, payload); }) .send() .whenComplete((subAck, throwable) -> { if (throwable != null) { log.error("订阅失败", throwable); } else { log.info("订阅成功"); } });这里订阅的是device/+/command,用了单层通配符,能匹配所有设备的指令主题。回调里拿到消息后,根据主题解析出设备 ID,再执行对应的指令处理逻辑。
注意:回调方法是在 IO 线程里执行的,如果处理逻辑耗时较长,一定要放到业务线程池里执行,否则会阻塞消息接收,导致消息积压。
4.4 消息序列化与协议设计经验
实际项目里,消息体一般不会直接传裸字符串,而是用 JSON 或者 Protobuf 序列化。JSON 可读性好,调试方便,但体积大。Protobuf 体积小,解析快,但需要定义 schema,调试麻烦。
我的经验是:设备端资源紧张、消息量大的场景用 Protobuf;普通场景用 JSON 就够了。JSON 库推荐 Jackson 或者 Fastjson2,序列化和反序列化都很方便。
消息协议设计上,建议统一格式,比如:
{ "msgId": "uuid", "timestamp": 1700000000000, "type": "command", "data": { "action": "read", "params": {} } }msgId用于消息去重和追踪,timestamp用于判断消息时效性,type区分消息类型,data放具体业务数据。这种结构清晰,扩展也方便。
5. MQTT 对接 485 设备的完整方案
5.1 485 设备接入的整体架构
485 设备是工业场景里最常见的设备类型,很多传感器、PLC、仪表都是 485 接口。485 是一种物理层协议,本身不支持网络通信,所以要让 485 设备接入 MQTT,中间需要一个网关设备来做协议转换。
典型的架构是这样的:485 设备通过 485 总线连接到网关,网关把 485 协议(通常是 Modbus RTU)转换成 MQTT 消息发到 Broker,后端服务订阅 MQTT 主题来接收数据。反过来,后端要控制 485 设备时,往 MQTT 主题发指令,网关订阅到指令后转换成 485 信号发给设备。
网关可以选现成的工业网关,比如有人物联网、映翰通这些厂家的产品,也可以自己用树莓派或者工控机加 485 扩展板来做。现成网关的好处是稳定、开箱即用,缺点是灵活性差,有些定制协议支持不了。自己搭网关灵活度高,但需要写代码,维护成本也高。
5.2 网关端协议转换的核心逻辑
网关端的核心工作就是 485 协议和 MQTT 协议之间的双向转换。以 Modbus RTU 为例,读取设备数据的过程是:网关通过 485 总线发送 Modbus 读寄存器指令,设备返回寄存器数据,网关解析后封装成 JSON,通过 MQTT 发布出去。
Modbus RTU 的读指令格式是:设备地址(1字节)+ 功能码(1字节)+ 起始寄存器地址(2字节)+ 寄存器数量(2字节)+ CRC 校验(2字节)。比如要读取设备地址为 1 的设备,从寄存器 0x0000 开始读 2 个寄存器,指令就是:
01 03 00 00 00 02 C4 0B网关收到设备返回的数据后,解析出寄存器值,再转成 MQTT 消息。比如温度值存在寄存器 0x0000,原始值是 235,实际温度是 23.5 度(假设精度是 0.1),那 MQTT 消息就是:
{ "deviceId": "485-sensor-01", "temperature": 23.5, "timestamp": 1700000000000 }发布到device/485-sensor-01/data主题。
5.3 通过 MQTT 下发指令读取 485 设备数据
后端要主动读取 485 设备数据时,流程是这样的:后端往device/485-sensor-01/command主题发指令,网关订阅了这个主题,收到指令后转换成 Modbus 读指令发给 485 设备,设备返回数据后网关再通过 MQTT 发回来。
指令的 JSON 格式可以设计成:
{ "action": "read", "slaveId": 1, "functionCode": 3, "startAddress": 0, "quantity": 2 }网关端的处理逻辑:
client.subscribeWith() .topicFilter("device/+/command") .qos(MqttQos.AT_LEAST_ONCE) .callback(publish -> { String deviceId = extractDeviceId(publish.getTopic().toString()); String payload = new String(publish.getPayloadAsBytes()); Command cmd = objectMapper.readValue(payload, Command.class); // 转换成 Modbus RTU 指令 byte[] modbusCmd = buildModbusReadCommand( cmd.getSlaveId(), cmd.getStartAddress(), cmd.getQuantity() ); // 通过 485 发送指令并读取响应 byte[] response = serialPort.sendAndReceive(modbusCmd); // 解析响应并发布到 MQTT SensorData data = parseModbusResponse(response); client.publishWith() .topic("device/" + deviceId + "/data") .payload(objectMapper.writeValueAsBytes(data)) .qos(MqttQos.AT_LEAST_ONCE) .send(); }) .send();这里的关键点是串口通信的时序控制。485 是半双工总线,发送和接收不能同时进行,发送完指令后要等待设备响应,响应超时时间一般设 500ms 到 1 秒。另外总线上如果有多个设备,要保证同一时间只有一个设备在发送,否则会冲突。
5.4 数据采集频率与异常处理策略
数据采集频率要根据实际需求来定。温度、湿度这种变化慢的,30 秒到 1 分钟采集一次就够了。电流、电压这种变化快的,可能需要 1 秒甚至更短。但采集频率越高,485 总线的负载越大,网关的处理压力也越大。
我的经验是:先按业务需求定一个基础采集频率,然后观察 485 总线的负载情况,如果总线利用率超过 70%,就要考虑降低频率或者增加总线。另外,对于变化缓慢的数据,可以在网关端做变化上报,只有数据变化超过阈值时才发 MQTT 消息,这样能大幅减少消息量。
异常处理方面,要处理几种情况:485 设备无响应、返回数据 CRC 校验失败、返回数据格式异常。这些情况网关都要能识别并做相应处理,比如重试、上报异常状态、记录日志。重试次数一般设 2 到 3 次,超过就放弃并上报设备异常。
6. 常见问题排查与避坑指南
6.1 连接与认证类问题速查
| 问题现象 | 可能原因 | 排查方法 |
|---|---|---|
| 连接被拒绝 | 用户名密码错误 | 检查认证配置,用 MQTT 客户端工具测试 |
| 连接超时 | 端口不通或防火墙拦截 | telnet 测试端口,检查防火墙规则 |
| 频繁断线重连 | KeepAlive 设置过短 | 适当增大 KeepAlive,检查网络稳定性 |
| 客户端 ID 冲突 | 多个客户端用同一 ID | 确保每个客户端 ID 唯一 |
客户端 ID 冲突这个问题特别隐蔽,因为 Broker 的处理方式是后连接的踢掉先连接的,表现就是两个客户端轮流掉线。我之前有个项目,设备端代码里客户端 ID 写死了,批量烧录后所有设备 ID 都一样,上线后互相踢,排查了半天才发现。后来改成用设备序列号做客户端 ID 就解决了。
6.2 消息丢失与重复的排查思路
消息丢失一般有几个原因:QoS 等级设置不对、会话保持没开、订阅关系丢失。如果发现消息丢失,先确认发布和订阅的 QoS 等级,再看 Clean Session 是否设为 false,最后检查订阅是否成功。
消息重复在 QoS 1 下是正常现象,因为 QoS 1 保证的是至少一次,不保证不重复。解决方式是在应用层做幂等处理,比如用 msgId 去重,或者用业务上的唯一标识来判断。我一般会在消息体里加一个 msgId,接收端维护一个最近消息 ID 的缓存,收到重复的直接丢弃。
6.3 大量设备连接时的性能调优
设备量大了之后,Broker 的性能调优就很重要。EMQX 的话,主要调整这几个参数:最大连接数、消息队列长度、TCP 缓冲区大小。另外,操作系统的文件描述符限制也要调大,默认的 1024 肯定不够,改成 65535 或者更大。
还有一点容易被忽略:如果大量设备同时断线重连,会给 Broker 带来很大的瞬时压力。这种情况可以通过在客户端加随机延迟来缓解,比如重连延迟在 1 到 10 秒之间随机,避免所有设备同时重连。
6.4 我踩过的那些坑与经验总结
第一个坑是主题设计太随意。早期项目主题命名没有规范,有的用device/data,有的用data/device,后来设备多了之后完全乱套。建议一开始就定好主题规范,比如{产品线}/{设备类型}/{设备ID}/{数据类型},这样后期维护方便很多。
第二个坑是没做消息大小限制。MQTT 协议本身对消息大小没有硬性限制,但 Broker 一般都有配置。如果发了超大消息,可能会被 Broker 拒绝或者导致连接断开。建议单条消息控制在 1KB 以内,超过的用分片或者改用其他方式传输。
第三个坑是忽略了时间同步。设备端的时间如果不准,消息里的 timestamp 就没有参考价值。建议设备端定期通过 NTP 同步时间,或者由服务端在收到消息时打上时间戳。
第四个坑是遗嘱消息的主题和正常消息主题混在一起。遗嘱消息应该用单独的主题,比如device/{deviceId}/status,正常数据用device/{deviceId}/data,这样后端处理逻辑清晰,不会混淆。
7. 从开发到上线的完整检查清单
项目开发完了要上线,有几个事情必须确认。Broker 的认证和 ACL 配置好了没有,TLS 证书装了没有,这些安全相关的必须在上线前搞定。客户端的重连策略测试过没有,模拟断网再恢复,看设备能不能自动重连并恢复订阅。消息的 QoS 等级是否符合业务需求,关键指令有没有用 QoS 1 或 2。遗嘱消息配置了没有,设备离线告警能不能正常工作。
监控和告警也要提前配好。Broker 的连接数、消息吞吐量、系统资源使用率这些指标要监控起来,设置合理的告警阈值。客户端这边,连接状态、消息发送失败率、消息处理延迟这些也要有监控。我一般会用 Prometheus 加 Grafana 来做监控面板,EMQX 自带 Prometheus 集成,配置起来很方便。
压测也是上线前必须做的。用 JMeter 或者 emqtt_bench 模拟大量设备连接和消息收发,看看 Broker 和网关能不能扛住。压测的时候要关注几个指标:连接建立成功率、消息到达率、消息延迟、Broker 的 CPU 和内存使用率。根据压测结果来调整 Broker 配置和网关的并发处理能力。
最后再分享一个小技巧:在设备端加一个本地缓存,网络断开的时候把数据先存本地,恢复连接后再补发。这样即使网络不稳定,数据也不会丢。缓存大小根据设备存储空间来定,一般存最近几小时到几天的数据就够了。补发的时候要注意消息的时效性,太旧的数据如果业务上不需要,可以直接丢弃,避免补发大量历史数据把 Broker 冲垮。