MQTT是一种基于发布订阅模式的轻量级消息协议,最初为网络带宽有限、设备性能不高的物联网场景设计。它协议头小、支持三种QoS等级、内置心跳保活机制,在设备消息上报、服务端指令下发、实时消息推送等场景中表现出色。Spring Boot虽然没有官方starter直接支持MQTT,但借助Eclipse Paho客户端可以很方便地完成集成。本文将从概念讲起,带你完成一个完整的集成案例,并汇总常见坑点。

一、MQTT核心概念速览
在动手写代码之前,有必要先把几个核心概念理清楚。MQTT的通信模型是发布订阅模式,消息的发送方称为发布者(Publisher),接收方称为订阅者(Subscriber),两者通过主题(Topic)解耦,发布者不需要知道谁在消费消息,只需要把消息发到某个主题上,中间的Broker负责转发。常见的Broker有EMQX、Mosquitto、HiveMQ等,生产环境推荐使用EMQX,它提供了完善的管理控制台和集群能力。
主题采用层级结构,例如device/001/status表示001号设备的状态消息。订阅时可以使用通配符:+匹配单层,如device/+/status可以匹配所有设备的状态主题;#匹配多层,如device/#可以匹配device下所有子主题。QoS等级分为三级:QoS0最多发送一次,可能丢消息;QoS1至少发送一次,可能重复;QoS2恰好发送一次,开销最大。多数业务场景用QoS1即可,配合业务侧幂等设计来处理重复消息。
还有一个容易忽略的概念是遗嘱消息(Will Message)。客户端连接时可以预先设置一条遗嘱消息,当客户端异常掉线时,Broker会自动向指定主题发布这条消息,订阅者就能感知到设备离线。这在物联网监控中非常实用,建议每个设备客户端都配置遗嘱消息,主题可以约定为device/{id}/offline。
二、Spring Boot集成MQTT完整步骤
首先在pom.xml中引入Paho客户端依赖。注意版本建议选择1.2.5及以上,低版本存在已知的连接稳定性问题。同时可以引入spring-integration-mqtt,它提供了与Spring消息体系的整合能力,如果只是简单收发,直接用Paho原生API更轻量。
<dependency>
<groupId>org.eclipse.paho</groupId>
<artifactId>org.eclipse.paho.client.mqttv3</artifactId>
<version>1.2.5</version>
</dependency>接着在application.yml中配置连接参数。Broker地址、用户名、密码、客户端ID是四个必备项,其中客户端ID尤其重要,同一个Broker上不允许两个相同ID的客户端同时在线,后连接的会把先连接的踢下线,这是新手最常踩的坑之一。因此建议客户端ID拼接上随机后缀或设备唯一标识。
mqtt:
broker-url: tcp://127.0.0.1:1883
username: admin
password: public
client-id: springboot-provider-${random.uuid}
topic: test/topic
qos: 1然后编写配置类,创建MQTT客户端并完成连接。这里用@ConfigurationProperties接收配置,通过MqttClient建立连接,并注册一个回调处理消息接收和断线事件。代码中演示了自动重连的设置,setAutomaticReconnect(true)开启后Paho会自动尝试恢复连接。
@Configuration
public class MqttConfig {
@Value("${mqtt.broker-url}")
private String brokerUrl;
@Value("${mqtt.username}")
private String username;
@Value("${mqtt.password}")
private String password;
@Value("${mqtt.client-id}")
private String clientId;
@Bean
public MqttClient mqttClient() throws Exception {
MqttClient client = new MqttClient(brokerUrl, clientId,
new MemoryPersistence());
MqttConnectOptions options = new MqttConnectOptions();
options.setUserName(username);
options.setPassword(password.toCharArray());
options.setAutomaticReconnect(true); // 断线自动重连
options.setCleanSession(true); // 干净会话
options.setConnectionTimeout(10); // 连接超时时间,单位秒
options.setKeepAliveInterval(30); // 心跳间隔,单位秒
client.connect(options);
return client;
}
}发布消息非常简单,调用publish方法即可。需要留意MqttMessage的payload是字节数组,发送对象时可以先序列化为JSON字符串再转字节。订阅则调用subscribe方法并传入回调,Spring Boot项目中可以把订阅逻辑放在@PostConstruct或监听器中,服务启动后自动订阅目标主题。
@Service
public class MqttService {
@Autowired
private MqttClient mqttClient;
// 发布消息
public void publish(String topic, String content) throws Exception {
MqttMessage message = new MqttMessage();
message.setPayload(content.getBytes(StandardCharsets.UTF_8));
message.setQos(1);
mqttClient.publish(topic, message);
}
// 订阅消息
@PostConstruct
public void subscribe() throws Exception {
mqttClient.subscribe("test/#", (topic, message) -> {
String body = new String(message.getPayload(),
StandardCharsets.UTF_8);
System.out.println("收到消息,主题:" + topic + ",内容:" + body);
});
}
}写一个简单的Controller测试一下,访问接口后消息被发到Broker,订阅端回调立即打印出内容,整条链路就通了。如果使用Spring Integration,还可以通过MqttPahoMessageDrivenChannelAdapter和MqttPahoMessageHandler以更Spring化的方式收发,适合消息需要进入消息通道做流转的复杂场景。
三、常见问题与注意事项
第一类问题是连接反复断开。排查思路:先确认客户端ID是否重复,多实例部署时尤其要注意给每个实例配置不同的ID;再检查心跳间隔,网络不稳定的环境下可以把keepAliveInterval调小一些让掉线被更快感知;最后看Broker端的连接数限制,EMQX默认有最大连接数配置,超出后新连接会被拒绝。
第二类问题是消息丢失或重复。QoS0场景下丢消息是正常现象,重要业务至少用QoS1;QoS1的重复特性要求消费端做幂等处理,常见做法是给每条消息附带唯一ID,消费时先查Redis判断是否处理过。另外cleanSession设为true时,客户端离线期间的QoS1消息不会被保留,如果希望离线消息补发,需要将该参数设为false,但这样会占用Broker端的会话存储,要权衡使用。
第三类是消息顺序问题。MQTT只保证同一主题内QoS1、QoS2消息的顺序,跨主题无法保证。如果业务对顺序敏感,建议单主题串行发布,或者在消息体内携带序号和时间戳,由消费端重排。此外,订阅回调中不要执行耗时操作,Paho的回调是单线程串行执行的,阻塞回调线程会拖慢后续消息的处理,重逻辑应交给线程池异步处理。
最后给出两个练习建议:一是动手实现一个设备上下线监控小项目,利用遗嘱消息和device/+/offline通配订阅感知设备状态;二是尝试把QoS、retain消息、离线会话三个参数分别调整并观察消息行为差异,通过对比加深对协议的理解。实践出真知,跑通这些例子后你对MQTT的掌握会扎实很多。
Spring Boot MQTTMQTT集成消息队列修改时间:2026-09-08 19:23:07