去年做充电桩数据采集平台的时候,我深刻体会到一件事:设备端和服务端之间,HTTP轮询真的不是长久之计。几千台充电桩走4G模块,每台每隔几秒就要上报一次状态,HTTP的握手开销、服务端的连接压力、设备端的功耗全部拉满,还没算上服务端想主动下发指令给设备这种"反向"需求,HTTP基本做不了。后来把通信层切到MQTT,整个链路立刻轻量了,一个长连接搞定双向通信,流量和响应速度都上了一个档次。
这次基于Spring Boot实现MQTT通信,我就把当时从协议认知、Broker选型、代码集成到生产环境排坑的完整过程整理出来。适合两类人看:一是想在Spring Boot项目里快速接入MQTT、但不想看一堆晦涩文档的后端同学;二是已经在用MQTT、但被重连、丢消息、Topic设计这些事反复折磨的工程师。全文会用一套能直接落地的代码和配置,把每个关键决策背后的原因也讲清楚,不只是给步骤。
1. 为什么设备端通信绕不开MQTT
1.1 这不是又一个消息队列:MQTT解决的核心问题
很多第一次接触MQTT的同学会把它和Kafka、RabbitMQ混为一个概念,其实两者定位完全不同。Kafka和RabbitMQ解决的是服务端之间的异步解耦、削峰填谷,而MQTT解决的是海量设备与服务端之间的可靠通信。
设备端场景有几个HTTP完全扛不住的特点:网络不稳定、带宽有限、设备可能随时离线、服务端还需要主动下发指令。MQTT基于TCP长连接,报文头部最小只有2字节,一个发布订阅消息的协议开销比HTTP小一个数量级。我实测过同一台4G模块,用HTTP轮询每次请求加响应大概2KB,换成MQTT上报一条JSON压缩到300字节左右,流量直接省了80%以上。
MQTT在TCP之上建立的是持久连接,只要网络不断,客户端和服务端随时可以互发消息,不需要像HTTP那样每次先建连再断开。这对设备上下行通信的实时性提升非常明显。
1.2 发布/订阅模型和Topic通配符:TCP长连接上的轻量级广播
MQTT的核心模型是发布/订阅,消息的生产者把消息发到一个叫Topic的主题上,订阅了该Topic的所有客户端都能收到。这个模型天然解耦了设备和服务端:设备不需要知道服务端地址,服务端也不需要知道设备IP,大家只认Topic。
Topic本身是层级结构,用斜杠分隔,比如:
chargepile/sh001/temperature chargepile/sh001/statusMQTT提供了两个通配符:
+:匹配单层,比如chargepile/+/temperature能匹配所有充电桩的温度主题#:匹配多层,比如chargepile/#能匹配chargepile下的所有主题
这个设计是MQTT的灵魂。服务端只需要订阅一个带通配符的主题,就能收到所有设备的上报数据,新设备上线也不需要额外注册——只要它往自己约定的Topic发消息,服务端自然就能收到。
1.3 QoS、遗嘱、保留消息:物联网场景的三个关键特性
这三个特性是HTTP完全没有的,也是MQTT在物联网场景不可替代的原因。
**QoS(服务质量)**有三个级别:
| 级别 | 含义 | 适用场景 |
|---|---|---|
| QoS 0 | 最多一次,发完就扔,不确认不重试 | 高频状态上报,丢了就丢了 |
| QoS 1 | 至少一次,有确认有重试,可能重复 | 大部分业务数据上报、指令下发 |
| QoS 2 | 恰好一次,四次握手机制,开销最大 | 极少数严格要求不重复不丢失的场景 |
我的实践经验:物联网项目里90%的消息用QoS 1就够了,QoS 2协议开销太大,只在资金、订单这类极端场景使用。
遗嘱消息(Last Will):客户端在连接时可以在Broker上留下一条遗嘱消息。如果客户端异常掉线(比如断网、断电、崩溃),Broker会代替这个客户端把遗嘱消息发到指定Topic。我在项目中用它来做设备掉线告警,省掉了服务端定时轮询设备在线状态的成本。
保留消息(Retained):Broker会帮客户端保存每个Topic的最后一条消息。新设备订阅该Topic时,可以立刻收到最新状态,不用等设备主动上报。这个特性在设备刚上线需要快速拿到服务端下发的最新配置时特别有用。
2. 环境准备:Broker选型与Spring Boot工程初始化
2.1 Broker选型:EMQX、Mosquitto还是公有云
MQTT通信必须有Broker(消息代理服务器)来中转消息,选型直接决定后续的运维体验和性能上限。
- EMQX:开源且社区活跃,支持MQTT 3.1.1和MQTT 5.0,自带Web控制台,支持集群、ACL、插件扩展,性能很强。我最终生产环境选的是它,单节点撑住上万台设备没有问题。推荐正式项目直接用。
- Mosquitto:Eclipse基金会出品,极简轻量,适合嵌入式设备、局域网小规模场景或学习体验。它的配置文件很传统,没有可视化界面,管理起来不方便。
- 云厂商托管服务:如果不想自己运维Broker,可以用云厂商提供的MQTT实例。优点是不用考虑高可用,缺点是对特定云有绑定,且成本随规模上涨。
- 公共测试Broker:EMQX官方提供的
broker.emqx.io,用于代码联调和功能验证很方便,但生产环境千万别用,速度和稳定性都没有保障。
从学习到生产的路径,我建议:先用Docker把EMQX跑起来感受一下,等摸透了再根据规模决定是继续自建还是上云。
2.2 Docker快速搭建EMQX
EMQX 5.x的镜像已经非常完善,一条命令就能跑起来:
docker run -d --name emqx \ -p 1883:1883 \ -p 8083:8083 \ -p 8084:8084 \ -p 18083:18083 \ emqx/emqx:5.8.0端口说明:
1883:MQTT over TCP 端口,服务端和设备端都走这个8883:MQTT over SSL/TLS 端口,生产环境建议启用8083:MQTT over WebSocket 端口,浏览器端调试用8084:MQTT over WSS 端口,浏览器端加密连接用18083:EMQX Dashboard 控制台端口,默认账号admin,密码public
启动后访问http://localhost:18083,在Dashboard里能看到所有连接上的客户端、订阅的Topic、消息收发速率。生产排查问题的时候,控制台里能直接看到客户端连接状态和离线原因,这是我最依赖的排查入口。
2.3 创建Spring Boot工程并引入MQTT相关依赖
集成方式上有个前提要搞清楚:Spring Boot本身没有MQTT原生能力,底层还是要用Eclipse Paho客户端。Paho是Java生态最流行的MQTT客户端库,如果你不想引入Spring Integration那套抽象,直接用Paho也是完全可行的。
但我的建议是:如果项目整体已经构建在Spring Boot上,就加上Spring Integration MQTT模块。它能让你用Spring风格的声明式方式管理连接、订阅、消息收发,还能和Spring的Channel、Gateway机制整合,代码更优雅,可维护性更好。
创建工程时,在pom.xml里加入这些依赖:
<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-integration</artifactId> </dependency> <dependency> <groupId>org.springframework.integration</groupId> <artifactId>spring-integration-mqtt</artifactId> </dependency>这里有个小细节要注意:Spring Boot 2.x的commons-logging到 Spring Boot 3.x 的迁移过程中,Spring Integration 6.x 已经把底层包从javax换成了jakarta,如果你是从旧项目升级过来的,编译报ClassNotFoundException先往这个方向排查。
3. 核心集成:通过Spring Integration把MQTT接进Spring容器
3.1 连接工厂和MqttConnectOptions:所有坑的源头
MqttConnectOptions是连接配置的核心,很多隐蔽问题都出在这里。它控制着连接是否持久、心跳间隔、自动重连策略等关键行为。
mqtt: broker: tcp://localhost:1883 client-id: charging-server username: admin password: public default-topic: chargepile/+/report default-qos: 1对应的配置类:
@Configuration @IntegrationComponentScan public class MqttConfig { @Value("${mqtt.broker}") private String broker; @Value("${mqtt.client-id}") private String clientId; @Bean public MqttPahoClientFactory mqttClientFactory() { DefaultMqttPahoClientFactory factory = new DefaultMqttPahoClientFactory(); MqttConnectOptions options = new MqttConnectOptions(); options.setServerURIs(new String[]{broker}); options.setCleanSession(false); options.setConnectionTimeout(30); options.setKeepAliveInterval(60); options.setAutomaticReconnect(true); options.setMaxInflight(100); factory.setConnectionOptions(options); return factory; } }几个参数背后的逻辑我解释一下:
- setCleanSession(false):让Broker持久化当前客户端的会话。客户端离线期间Broker会帮它保存QoS 1/2的消息,重连后自动补发。如果设为true,客户端重连后是不会收到离线期间的消息的。但对于服务端应用来说,离线补发不一定是好事,这时要权衡消息积压和业务实时性,后面我会详细说。
- setKeepAliveInterval(60):心跳间隔60秒。客户端和Broker之间通过PINGREQ/PINGRESP维持连接,默认的10秒太频繁,生产环境建议在30到120秒之间,否则浪费流量。
- setAutomaticReconnect(true):开启自动重连。不开启的话,网络抖动导致连接断开后客户端不会主动恢复连接,服务可能长时间"假死",这是生产环境最致命的问题之一。
3.2 入站通道:订阅主题并监听消息
Spring Integration MQTT把消息接收抽象成MessageProducer,最常用的是MqttPahoMessageDrivenChannelAdapter。它启动后会自动订阅指定Topic,并把收到的消息转成Spring的Message对象发送到一个输出通道。
@Bean public MessageProducer mqttInbound() { String[] topics = {"chargepile/+/report", "chargepile/+/event"}; MqttPahoMessageDrivenChannelAdapter adapter = new MqttPahoMessageDrivenChannelAdapter(clientId, mqttClientFactory(), topics); adapter.setCompletionTimeout(5000); adapter.setQos(1); adapter.setOutputChannelName("mqttInboundChannel"); return adapter; } @Bean public MessageChannel mqttInboundChannel() { return new DirectChannel(); }然后通过@ServiceActivator编写真正的消息处理逻辑:
@ServiceActivator(inputChannel = "mqttInboundChannel") public void handleMqttMessage(@Header(MqttHeaders.RECEIVED_TOPIC) String topic, @Payload String payload) { log.info("收到主题 [{}] 的消息: {}", topic, payload); // 根据Topic分发到不同的业务处理逻辑 if (topic.startsWith("chargepile/")) { chargePileService.processReport(topic, payload); } }这里有个设计点值得注意:MqttHeaders.RECEIVED_TOPIC是Spring Integration自动注入的消息头,能拿到消息来自哪个Topic。实际项目中,同一个Adapter订阅了多个Topic,消息处理时要按照Topic做路由分发,不要把所有逻辑糊在一个方法里。
3.3 出站通道:使用@MessagingGateway发布消息
消息发布用MqttPahoMessageHandler配合@MessagingGateway实现,这样业务层只依赖网关接口,不需要关心MQTT底层细节。
@Bean @ServiceActivator(inputChannel = "mqttOutboundChannel") public MessageHandler mqttOutbound() { MqttPahoMessageHandler handler = new MqttPahoMessageHandler(clientId + "-pub", mqttClientFactory()); handler.setAsync(true); handler.setDefaultQos(1); handler.setDefaultTopic("chargepile/unknown/command"); return handler; } @MessagingGateway(defaultRequestChannel = "mqttOutboundChannel") public interface MqttGateway { void publish(String payload, @Header(MqttHeaders.TOPIC) String topic); }注意这里我给出站客户端单独设置了一个clientId,加上-pub后缀。原因下面会在踩坑部分重点解释,先记住这是避免"互踢"的关键。
使用的时候,业务代码只需要注入网关接口:
@Autowired private MqttGateway mqttGateway; public void sendCommand(String deviceId, String command) { String topic = "chargepile/" + deviceId + "/command"; String payload = "{\"type\":\"start\",\"timestamp\":" + System.currentTimeMillis() + "}"; mqttGateway.publish(payload, topic); }3.4 收发消息的完整Demo
把上面几块拼在一起,就是一个完整的收发闭环。我习惯用一个简单的Controller先验证链路是否通:
@RestController @RequestMapping("/mqtt") public class MqttTestController { @Autowired private MqttGateway mqttGateway; @PostMapping("/publish") public String publish(@RequestParam String topic, @RequestParam String message) { mqttGateway.publish(message, topic); return "published to " + topic; } }设备侧或者测试客户端往chargepile/sh001/report发一条{"temperature": 26.5},服务端的handleMqttMessage就会打印日志。反向测试用Postman调一下/mqtt/publish?topic=chargepile/sh001/command&message=hello,设备侧能收到就算闭环成功。
我第一次跑通这个Demo的时候有个小插曲:往chargepile/+/report这个带通配符的Topic发布消息,结果收不到。后来反应过来,发布消息不能发到带通配符的Topic上,通配符只用于订阅匹配,Broker会拒绝这种发布。这也是新手最容易踩的概念性错误之一。
4. Topic设计与消息协议:从能用到可用
4.1 Topic层级命名:设备和产品维度怎么规划
很多初学MQTT的人把Topic当成随意起的字符串,觉得能收到消息就行。但一旦设备量上来了,Topic设计不合理会让权限控制、消息过滤、日志排查全部失控。
我参考过阿里云IoT和腾讯云IoT的Topic规范,它们的核心设计思想可以归纳为:固定前缀标识产品,中间层放设备维度,后面几层按业务功能细分。比如:
chargepile/{productKey}/{deviceId}/property/post 设备属性上报 chargepile/{productKey}/{deviceId}/event/post 设备事件上报 chargepile/{productKey}/{deviceId}/command 服务端指令下发 chargepile/{productKey}/{deviceId}/command/reply 设备指令应答 chargepile/{productKey}/{deviceId}/status/online 设备上下线状态这么设计的好处有三点:
- 权限可控:在EMQX的ACL配置里,可以明确指定某个设备只能发布到
chargepile/{自己的productKey}/{自己的deviceId}/#下的上报主题,只能订阅command主题,从Broker层面卡死越权行为。 - 隔离清晰:不同产品线(比如交流桩、直流桩、光储充一体桩)用不同的
productKey区分,服务端只订阅自己关心的产品线即可。 - 扩展性好:后续加功能只需要在设备维度后面追加层级,不会影响既有Topic。
层级层级不要设计超过4到5层,每层字段要稳定。尤其是设备端固件,如果Topic结构改了,旧设备要OTA升级才能适配,代价极大。
4.2 Payload设计:带幂等和时序的消息体
Topic只负责路由,业务数据全在Payload里。我见过不少团队直接把裸数据往Topic里丢,比如就发个26.5,解析倒是简单,但完全无法应对版本迭代和问题排查。
我的建议是统一用JSON结构,并且包含固定的公共字段:
{ "msgId": "7c9d8f6a2b1e4d5c", "timestamp": 1709900000000, "productKey": "charging-dc", "deviceId": "sh001", "type": "property", "data": { "voltage": 734.5, "current": 42.1, "temperature": 38.2 } }- msgId:全局唯一的消息ID,生成方式可以是UUID或者雪花算法。这个字段在QoS 1语义下用来做幂等去重非常关键,因为QoS 1可能出现重复投递,消费端拿msgId去Redis或者数据库去重,能保证业务不能重复执行。
- timestamp:设备采集时间的时间戳,毫秒级。设备离线补传时,服务端要根据这个时间戳判断数据时效性,而不是傻傻地按接收顺序入库,否则会覆盖新数据。
- type:消息子类型,配合Topic的业务后缀,服务端可以双保险地路由消息。设备端有时候Topic写死了不好改,加一个type字段让服务端可以灵活处理。
4.3 订阅策略:设备侧和服务端侧各自怎么订阅
订阅策略是双向的,两端要各司其职。
服务端:用通配符订阅,把某一类设备的数据全部接进来。
chargepile/+/+/property/post chargepile/+/+/event/post chargepile/+/+/status/online如果只关心某条产品线,可以再精确一点:
chargepile/dc-01/+/property/post设备端:只订阅属于自己的下行指令Topic和指令应答Topic,不需要订阅别人的。设备每次启动时订阅:
chargepile/{productKey}/{deviceId}/command chargepile/{productKey}/{deviceId}/command/reply这样做的好处是什么?设备端即使被黑客控制,在ACL约束下也只能收发自己那部分消息,没办法监听或干扰同产品线的其他设备,安全边界非常清晰。
5. 生产环境才会遇到的坑与调优
5.1 客户端ID重复导致互踢:一条消息都收不到
这是我踩过最离谱的坑。现象是服务端日志每隔几秒就出现一次连接成功又断开,消息时有时无,设备侧也是各种超时重连。查了半天发现,我在入站Adapter和出站MessageHandler的客户端配置中,用了同一个clientId。
MQTT协议规定clientId是客户端在Broker上的唯一标识,同一个clientId的新连接会把旧连接踢下线。我的入站和出站两个连接共用一个ID,两个连接在Broker眼里是同一个客户端,于是一个连上另一个就被踢,形成死循环。
解决办法很简单:出站连接使用clientId + "-pub"这类唯一后缀,同时Linux环境下的K8s多实例部署要特别小心,多个Pod不能共享同一个clientId,部署时可以通过环境变量注入实例唯一标识。
5.2 cleanSession和QoS组合下的消息可靠性
接续上文,光知道cleanSession布尔值是不够的,还要知道它和QoS怎么配合:
| cleanSession | QoS | 效果 |
|---|---|---|
| true | 0 | 离线期间消息全丢,重连后拿不到历史消息 |
| true | 1 | 离线期间消息可丢失,重连后的新消息正常 |
| false | 1 | Broker持久化会话,离线消息重连后补发 |
| false | 2 | Broker持久化且严格不重复,可靠性最高 |
服务端应用建议用cleanSession=false+QoS 1,两者配合能最大限度减少消息丢失。但也要留意另一个问题:如果服务端长时间宕机,Broker会持续为它保存离线消息,恢复上线后会瞬间涌入大量积压消息,冲击业务处理能力。我在生产环境里给EMQX配置了离线消息最大条数限制,超出部分按队列策略丢弃,确保恢复时不会雪崩。
5.3 阻塞与线程池:消息处理耗时的优化
Spring Integration的入站Adapter默认情况下消息处理是同步的,也就是说handleMqttMessage方法里如果在查数据库、调第三方接口,会阻塞后续消息的接收。
有个真实案例:设备上报的日志消息和告警消息走同一个Adapter,某次数据库慢查询导致handleMqttMessage卡住3秒,结果所有设备的实时上报全部滞后,在线状态判断全部失真。
解决方案是处理逻辑异步化。最简单的做法是把耗时操作丢进自定义线程池:
@ServiceActivator(inputChannel = "mqttInboundChannel") public void handleMqttMessage(MqttMessageWrapper wrapper) { asyncExecutor.execute(() -> { processBusiness(wrapper); }); }或者直接用Spring的@Async注解。如果你的业务需要对消息处理顺序有严格要求的场景,比如充电桩状态必须按时间顺序处理,就不要粗暴地异步化,而是用分区策略保证相同设备的消息落到同一个处理线程。
5.4 安全加固:认证、ACL与TLS
本地开发用默认配置没问题,上生产前安全这块必须补上。我遇到过不少项目直接把Broker裸奔在公网,不少还被扫描爆破过,轻则流量被刷爆,重则设备被恶意下发指令。
认证:EMQX默认开了用户名密码认证,生产环境建议启用内置数据库或接入外部认证,定期更换强密码。
ACL权限控制:用ACL限制每个客户端的发布订阅权限。核心原则是:只给最小必要的权限。设备只能往自己的Topic发,服务端只能往指令Topic下发。EMQX 5.x可以在Dashboard里配置ACL规则,也可以通过内置SQL数据库管理,建议用内置SQL方式,规则更灵活。
TLS加密:1883端口是明文协议,设备公网接入时,报文内容是裸奔的。生产环境一定用8883端口的TLS加密通信。如果是自建证书,设备端要内置CA根证书。这里有个性能注意点:TLS握手和加解密会增加设备端功耗和服务端CPU开销,如果对实时性要求高,可以减少证书校验的密码套件。
我还想提醒一个容易被忽略的安全细节:不要在Topic和Payload里暴露内网IP、数据库地址、账号密码等敏感信息。以前碰过一家企业的设备上行日志带着完整内网拓扑信息,只要拿到一台设备的通信报文,整个内网架构就暴露了。
6. 如果不用Spring Integration:直接用Paho的轻量方案
Spring Integration确实方便,但如果你只是做一个简单的监控脚本、定时上报任务,或者想彻底理解MQTT客户端的工作原理,直接使用Eclipse Paho更直接。
import org.eclipse.paho.client.mqttv3.*; public class SimpleMqttClient { public static void main(String[] args) throws Exception { String broker = "tcp://localhost:1883"; String clientId = "simple-client-demo"; MqttClient client = new MqttClient(broker, clientId); MqttConnectOptions options = new MqttConnectOptions(); options.setCleanSession(true); options.setConnectionTimeout(30); options.setKeepAliveInterval(60); options.setAutomaticReconnect(true); client.setCallback(new MqttCallback() { @Override public void connectionLost(Throwable cause) { log.warn("连接断开: {}", cause.getMessage()); } @Override public void messageArrived(String topic, MqttMessage message) { log.info("收到 {} 消息: {}", topic, new String(message.getPayload())); } @Override public void deliveryComplete(IMqttDeliveryToken token) { log.info("消息发送完成: {}", token.isComplete()); } }); client.connect(options); client.subscribe("chargepile/+/report", 1); MqttMessage message = new MqttMessage("hello mqtt".getBytes()); message.setQos(1); client.publish("test/topic", message); } }这段代码逻辑非常清楚:connect连接、subscribe订阅、setCallback注册回调、publish发布。相比Spring Integration,代码量少了一大截,也没有各种Channel和Gateway概念。
我的建议是:无论你用不用Spring Integration,都要先用Paho写一个小Demo,把连接、订阅、回调、发布这四个基本操作跑通。理解了底层的调用关系,你再回头看Spring Integration的封装,就会觉得它只是在Paho外面套了一层Spring的壳,遇到问题排查起来也更有方向感。
Paho的线程模型也要注意:messageArrived回调默认在Paho的Receiver线程里执行,耗时的业务逻辑建议同步丢给业务线程池,否则会阻塞后续消息的处理。这和Spring Integration默认同步处理踩的坑是同源的。
一些实操中的额外心得
整个接入过程中,有几个软件和技巧如果从一开始就知道,能省下很多试错时间。
MQTT客户端调试工具:我常用MQTT X作为日常调试客户端,它跨平台,支持MQTT 3.1.1和5.0,也能模拟QoS、遗嘱、保留消息这些特性。联调时,设备端没开发好之前,我都是先用它模拟设备发消息,这样服务端逻辑可以提前开发验证。
EMQX Dashboard的在线调试:EMQX 5.x控制台自带WebSocket客户端,可以直接在网页里订阅Topic并在页面里发消息,快速验证Broker和Topic配置是否正常,不需要本地装任何工具。
日志配置:Spring Integration MQTT的报错日志默认级别较高,问题排查前先调整日志级别把MQTT相关内容打全:
logging: level: org.eclipse.paho: DEBUG org.springframework.integration.mqtt: DEBUG打开DEBUG日志以后,能清楚看到连接重试、消息收发、心跳交互的每一步过程,很多看似诡异的问题其实是底层重试和确认机制在正常工作,只是你之前看不到而已。
最后再说一个容易被忽略的性能细节:如果你的服务端会大量下发指令,建议把出站MqttPahoMessageHandler的setAsync(true)打开,避免每次publish都同步等待Broker确认,造成调用线程阻塞。高并发下发场景下,这个参数对吞吐量的影响非常大。
MQTT这套东西门槛不高,但真正用好的门道不少。希望这份从选型到排坑的完整记录,能帮你少走一段弯路。