引言
本文基于 spring-integration-mqtt + eclipse paho 实现一套生产可用的 mqtt 客户端方案,涵盖连接配置、断线自动重连、主题订阅、消息发布等核心能力,代码可直接落地使用。
一、前言
在物联网(iot)场景中,mqtt 协议因其轻量、低带宽占用、支持发布/订阅模型等特性,成为设备与服务端通信的首选协议。常见的应用场景包括:
- 传感器数据采集(气体、温度、湿度等)
- 设备指令下发
- 服务端消息推送
本文将带你从零实现一个 spring boot 项目中的 mqtt 客户端,具备以下能力:
✅ 基于 mqttasyncclient 的异步客户端
✅ 支持用户名/密码认证
✅ 支持多主题订阅(通配符 +、#)
✅ 断线自动重连 + 重连后重新订阅
✅ 统一的发布工具类 mqttservice
✅ 基于 @value 的外部化配置
二、环境与依赖
2.1 版本信息
| 组件 | 版本 |
|---|---|
| spring boot | 2.x / 3.x |
| spring-integration-mqtt | 与 spring boot 版本对应 |
| org.eclipse.paho.client.mqttv3 | 1.2.5 |
2.2 maven 依赖
<dependency>
<groupid>org.springframework.boot</groupid>
<artifactid>spring-boot-starter</artifactid>
</dependency>
<dependency>
<groupid>org.springframework.integration</groupid>
<artifactid>spring-integration-mqtt</artifactid>
</dependency>
<dependency>
<groupid>org.eclipse.paho</groupid>
<artifactid>org.eclipse.paho.client.mqttv3</artifactid>
<version>1.2.5</version>
</dependency>说明:spring-integration-mqtt 已经传递依赖了 paho 客户端,但为避免版本冲突,建议显式指定版本。
三、配置文件
在 application.yml 中添加如下配置:
mqtt: broker-url: tcp://172.16.18.112:1883 client-id: mqtt-client username: admin password: public qos: 1 keep-alive: 60 topic: sensor/gas/srv000/+,msg/srv000/+
配置说明
| 配置项 | 说明 |
|---|---|
| broker-url | mqtt broker 地址,支持 tcp://、ssl://、ws:// |
| client-id | 客户端唯一标识,同一 broker 下不可重复 |
| username / password | 认证信息 |
| qos | 服务质量等级:0(最多一次)、1(至少一次)、2(恰好一次) |
| keep-alive | 心跳间隔(秒) |
| topic | 订阅主题,多个主题用 英文逗号 分隔,支持通配符 |
通配符说明
+:单层通配符,如sensor/gas/srv000/+可匹配sensor/gas/srv000/co、sensor/gas/srv000/ch4#:多层通配符,如sensor/#可匹配sensor/gas/srv000/co
四、核心代码实现
4.1 mqttconfig —— 客户端与连接配置
@slf4j
@configuration
public class mqttconfig {
@value("${mqtt.broker-url}")
private string brokerurl;
@value("${mqtt.client-id}")
private string clientid;
@value("${mqtt.username}")
private string username;
@value("${mqtt.password}")
private string password;
@value("${mqtt.qos}")
private integer qos;
@value("${mqtt.topic}")
private string tpoic;
@bean
public mqttconnectoptions mqttconnectoptions() {
mqttconnectoptions options = new mqttconnectoptions();
options.setserveruris(new string[]{brokerurl});
// 是否清除会话
options.setcleansession(true);
// 心跳间隔,单位为秒
options.setkeepaliveinterval(300);
// 连接超时时间,单位为秒
options.setconnectiontimeout(30);
// 是否自动重连(这里关闭,使用自定义重连机制)
options.setautomaticreconnect(false);
options.setusername(username);
options.setpassword(password.tochararray());
return options;
}
@bean
public imqttasyncclient mqttasyncclient(mqttconnectoptions options,
mqttcallbackhandler callbackhandler) {
try {
memorypersistence persistence = new memorypersistence();
imqttasyncclient client = new mqttasyncclient(brokerurl, clientid, persistence);
client.setcallback(callbackhandler);
client.connect(options).waitforcompletion();
log.info("成功连接到mqtt");
subscribe(client);
return client;
} catch (mqttexception e) {
log.error("创建mqtt客户端失败", e);
throw new runtimeexception("无法创建mqtt客户端", e);
}
}
public void subscribe(imqttasyncclient client) {
try {
string[] topics = tpoic.split(",");
for (string topic : topics) {
client.subscribe(topic, qos);
}
} catch (exception e) {
log.error("mqtt主题订阅失败:{}", e.getmessage(), e);
}
}
}关键设计点
- 使用
imqttasyncclient异步客户端:相比同步客户端,异步模型在重连、大批量发布场景下更稳定。 - 关闭 paho 内置自动重连:
setautomaticreconnect(false),交由自定义重连逻辑统一管理,便于观察日志、控制重连策略。 memorypersistence内存持久化:适合大多数场景;若需离线消息保留,可换成mqttdefaultfilepersistence。- 订阅逻辑抽离:重连后可以复用同一份订阅代码。
4.2 mqttcallbackhandler —— 回调与自动重连
@slf4j
@component
public class mqttcallbackhandler implements mqttcallback {
@value("${mqtt.qos}")
private integer qos;
@value("${mqtt.topic}")
private string tpoic;
@autowired
private dealdateworker dealdateworker;
@autowired
private dealmsgworker dealmsgworker;
@autowired
private applicationcontext applicationcontext;
private final atomicboolean reconnecting = new atomicboolean(false);
private scheduledexecutorservice reconnectscheduler;
@override
public void connectionlost(throwable throwable) {
log.warn("mqtt连接丢失", throwable);
if (reconnecting.compareandset(false, true)) {
schedulereconnect();
}
}
private void schedulereconnect() {
if (reconnectscheduler != null && !reconnectscheduler.isshutdown()) {
reconnectscheduler.shutdown();
}
reconnectscheduler = executors.newsinglethreadscheduledexecutor();
reconnectscheduler.scheduleatfixedrate(() -> {
try {
imqttasyncclient mqttclient = applicationcontext.getbean(imqttasyncclient.class);
mqttconnectoptions mqttconnectoptions = applicationcontext.getbean(mqttconnectoptions.class);
if (mqttclient != null && !mqttclient.isconnected()) {
log.info("尝试重新连接mqtt...");
mqttclient.connect(mqttconnectoptions).waitforcompletion();
if (mqttclient.isconnected()) {
log.info("mqtt重新连接成功");
resubscribetopics(mqttclient);
stopreconnect();
}
} else if (mqttclient != null && mqttclient.isconnected()) {
log.info("mqtt已连接,停止重连任务");
stopreconnect();
}
} catch (mqttexception e) {
log.error("mqtt重连失败,将在5秒后重试", e);
}
}, 0, 5, timeunit.seconds);
}
private void resubscribetopics(imqttasyncclient mqttclient) {
try {
string[] topics = tpoic.split(",");
for (string topic : topics) {
mqttclient.subscribe(topic, qos);
}
log.info("mqtt主题重新订阅完成");
} catch (exception e) {
log.error("mqtt主题订阅失败:{}", e.getmessage(), e);
}
}
private void stopreconnect() {
if (reconnectscheduler != null && !reconnectscheduler.isshutdown()) {
reconnectscheduler.shutdown();
reconnecting.set(false);
log.info("停止mqtt自动重连任务");
}
}
@override
public void messagearrived(string topic, mqttmessage message) throws exception {
string payload = new string(message.getpayload());
log.info("收到mqtt消息 - 主题: {}, qos: {}, 内容: {}", topic, message.getqos(), payload);
if (objectutils.isempty(payload)) {
return;
}
// 根据业务分发消息
// 例如:
// if (topic.startswith("sensor/gas/")) dealdateworker.handle(topic, payload);
// if (topic.startswith("msg/")) dealmsgworker.handle(topic, payload);
}
@override
public void deliverycomplete(imqttdeliverytoken token) {
try {
log.debug("消息交付完成 - 主题: {}",
token.gettopics() != null ? token.gettopics()[0] : "未知");
} catch (exception e) {
log.error("获取交付完成消息主题失败", e);
}
}
}重连机制核心思想
连接丢失(connectionlost)
│
▼
cas 判断是否已在重连中(reconnecting)
│
▼
启动 scheduledexecutorservice,每 5 秒检查一次
│
▼
判断 isconnected() → 否 → connect() → 成功后重订阅 → 停止调度器设计注意事项
使用 atomicboolean 防止重复重连connectionlost 在极端情况下可能被多次触发,cas 保证只有一次进入重连流程。
使用 scheduledexecutorservice 而不是 thread.sleep
可随时通过 shutdown() 停止重连,避免线程泄漏。
重连成功后必须重新订阅
paho 客户端在 cleansession=true 时,重连后订阅关系会丢失,必须重新订阅。
通过 applicationcontext 获取 bean 避免循环依赖mqttcallbackhandler 本身被 mqttconfig 使用,不能再反向注入 imqttasyncclient,否则会形成循环依赖。
4.3 mqttservice —— 统一发布工具
@slf4j
@service
public class mqttservice {
@value("${mqtt.qos}")
private integer qos;
private final imqttasyncclient mqttclient;
@autowired
@lazy
public mqttservice(imqttasyncclient mqttclient) {
this.mqttclient = mqttclient;
}
public void publish(string topic, string payload, int qos, boolean retained) {
try {
if (isconnected()) {
mqttclient.publish(topic, payload.getbytes(), qos, retained);
log.debug("已发布消息到主题 {}: {}", topic, payload);
} else {
log.warn("mqtt客户端未连接,无法发布消息");
}
} catch (mqttexception e) {
log.error("发布消息到主题 {} 失败", topic, e);
}
}
public void publish(string topic, string payload) {
publish(topic, payload, qos, false);
}
public boolean isconnected() {
return mqttclient != null && mqttclient.isconnected();
}
}使用示例
@restcontroller
@requestmapping("/api/mqtt")
public class mqttcontroller {
@autowired
private mqttservice mqttservice;
@postmapping("/publish")
public string publish(@requestparam string topic,
@requestparam string payload) {
mqttservice.publish(topic, payload);
return "ok";
}
}五、常见问题与踩坑记录
5.1 为什么要关闭automaticreconnect?
paho 内置的自动重连虽然方便,但存在以下问题:
- 重连成功后 不会自动重新订阅(尤其在
cleansession=true下) - 重连过程不可观测,无法定制日志、告警
- 无法与业务侧做联动(如重连后刷新缓存、上报状态)
因此推荐 关闭内置重连 + 自定义重连策略,可控性更强。
5.2cleansession=true与false的区别
| 选项 | 特点 | 适用场景 |
|---|---|---|
| true | 每次连接都是新会话,不接收离线消息 | 一般消费者 |
| false | 保留会话,可接收离线消息(qos≥1) | 关键指令通道 |
若使用 cleansession=false,需保证 clientid 固定不变,否则 broker 会视为新客户端。
5.3 消息乱序或重复
- qos 0:可能丢消息,不会重复
- qos 1:至少一次,可能重复 → 业务侧需幂等处理
- qos 2:恰好一次,性能最差
建议:传感器数据用 qos 0 或 1;指令通信用 qos 1 或 2 + 幂等。
5.4 clientid 冲突导致频繁断连
同一 broker 下,相同 clientid 的新连接会踢掉旧连接,表现为反复重连。生产环境建议:
client-id: mqtt-client-${spring.application.name}-${random.uuid}或使用服务器 ip + 应用名 + 随机后缀。
5.5 重连线程泄漏
scheduledexecutorservice 一定要在 stopreconnect() 或应用关闭时 shutdown()。
建议实现 disposablebean 或加 @predestroy:
@predestroy
public void destroy() {
stopreconnect();
try {
if (mqttclient != null && mqttclient.isconnected()) {
mqttclient.disconnect();
}
} catch (mqttexception e) {
log.error("mqtt断开失败", e);
}
}六、可优化方向
重连策略改为指数退避
例如 1s、2s、4s、8s…最大 60s,避免 broker 挂掉时大量无效重连。
使用 mqttpahomessagedrivenchanneladapter
若项目已经使用 spring integration,可基于适配器 + messagechannel 实现更声明式的订阅。
发布失败本地缓存 + 重试
对关键消息,可结合 redis / 数据库做失败落盘,连接恢复后补发。
健康检查与指标暴露
通过 actuator 暴露 mqtt 连接状态,prometheus 采集连接断开次数、重连次数、消息速率等。
tls 加密
生产环境建议用 ssl:// + ca 证书,避免明文传输用户名密码。
七、总结
本文实现了一套 完整的 spring boot mqtt 客户端方案,要点回顾:
| 模块 | 职责 |
|---|---|
mqttconfig | 连接参数、客户端实例创建、初始订阅 |
mqttcallbackhandler | 消息到达回调、断线检测、自动重连、重订阅 |
mqttservice | 对外统一的发布接口 |
核心思路是:关闭 paho 内置重连 → 自定义 cas 保护的重连调度 → 重连成功后重新订阅。
这套方案经过实际项目验证,稳定可靠,可作为通用模板直接复用。
以上就是springboot整合mqtt实现连接、重连、订阅、发布功能的详细内容,更多关于springboot mqtt连接、重连、订阅与发布的资料请关注代码网其它相关文章!
发表评论