ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

MQTT协议入门与实战:从发布订阅原理到Java客户端开发

MQTT协议入门与实战:从发布订阅原理到Java客户端开发 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 groupIdorg.eclipse.paho/groupId artifactIdorg.eclipse.paho.client.mqttv3/artifactId version1.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 trueMqttMessage 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);意思是如果这个客户端异常断开不是主动 disconnectBroker 就替它往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。等确认链路通了再收窄到具体主题。这个笨办法在排查消息去哪了的时候特别有效。
返回列表