当前位置: 代码网 > it编程>编程语言>Java > SpringBoot整合MQTT实现连接、重连、订阅、发布功能

SpringBoot整合MQTT实现连接、重连、订阅、发布功能

2026年09月28日 • Java •我要评论
引言本文基于 spring-integration-mqtt + eclipse paho 实现一套生产可用的 mqtt 客户端方案,涵盖连接配置、断线自动重连、主题订阅、消息发布等核心能力,代码可直

引言

本文基于 spring-integration-mqtt + eclipse paho 实现一套生产可用的 mqtt 客户端方案,涵盖连接配置、断线自动重连、主题订阅、消息发布等核心能力,代码可直接落地使用。

一、前言

在物联网(iot)场景中,mqtt 协议因其轻量、低带宽占用、支持发布/订阅模型等特性,成为设备与服务端通信的首选协议。常见的应用场景包括:

  • 传感器数据采集(气体、温度、湿度等)
  • 设备指令下发
  • 服务端消息推送

本文将带你从零实现一个 spring boot 项目中的 mqtt 客户端,具备以下能力:

✅ 基于 mqttasyncclient 的异步客户端
✅ 支持用户名/密码认证
✅ 支持多主题订阅(通配符 +、#)
✅ 断线自动重连 + 重连后重新订阅
✅ 统一的发布工具类 mqttservice
✅ 基于 @value 的外部化配置

二、环境与依赖

2.1 版本信息

组件版本
spring boot2.x / 3.x
spring-integration-mqtt与 spring boot 版本对应
org.eclipse.paho.client.mqttv31.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-urlmqtt 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);
        }
    }
}

关键设计点

  1. 使用 imqttasyncclient 异步客户端:相比同步客户端,异步模型在重连、大批量发布场景下更稳定。
  2. 关闭 paho 内置自动重连:setautomaticreconnect(false),交由自定义重连逻辑统一管理,便于观察日志、控制重连策略。
  3. memorypersistence 内存持久化:适合大多数场景;若需离线消息保留,可换成 mqttdefaultfilepersistence。
  4. 订阅逻辑抽离:重连后可以复用同一份订阅代码。

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连接、重连、订阅与发布的资料请关注代码网其它相关文章!

赞 (0)

相关文章:

版权声明:本文内容由互联网用户贡献,该文观点仅代表作者本人。本站仅提供信息存储服务,不拥有所有权,不承担相关法律责任。 如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 2386932994@qq.com 举报,一经查实将立刻删除。

发表评论

验证码:
Copyright © 2017-2026  代码网 保留所有权利. 粤ICP备2024248653号
站长QQ:2386932994 | 联系邮箱:2386932994@qq.com