砍材农夫砍材农夫
  • 微信记账小程序
  • java
  • redis
  • mysql
  • 场景类
  • 框架类
  • vuepress搭建
  • hexo搭建
  • 云图
  • llm wiki

    • 基于karpathy
    • gradle
  • 常用工具

    • git
    • gradle
    • Zadig
    • it-tools
    • 开源推荐
    • curl
  • 大前端

    • 环境配置
    • 微信生态
    • 正则
    • 全栈技能
  • java圈

    • java基础
    • jdk体系
    • jvm
    • spring框架
    • 分库分表
    • zookeeper
  • python技能

    • python编程
    • python数据
  • 算法

    • 算法
  • 网关

    • spring_cloud_gateway
    • openresty
  • 高可用

    • 秒杀
    • 分布式
    • 缓存一致
  • MQ

    • MQ
    • rabbitMQ
    • rocketMQ
    • kafka
  • 其它

    • 设计模式
    • 领域驱动(ddd)
  • 关系型数据库

    • mysql5.0
    • mysql8.0
  • 非关系型数据库

    • redis
    • mongoDB
  • 分布式/其他

    • ShardingSphere
    • 区块链
  • 向量数据库

    • M3E
    • OPEN AI
  • Jmeter
  • fiddler
  • wireshark
  • AI入门
  • AI大模型
  • AI插件
  • AI集成框架
  • 相关算法
  • AI训练师
  • 量化交易
  • AIoT
  • gitee
  • github
  • infoq
  • osc
  • 砍材工具
  • 关于
  • 相关运营
  • devops
  • 元宇宙
  • 区块链
  • 物联网
  • webrtc
  • web3.0
  • gitee
  • github
  • infoq
  • osc
  • 砍材工具
  • 关于
  • 中考
  • 投资
  • 保险
  • 思
  • 微信记账小程序
  • java
  • redis
  • mysql
  • 场景类
  • 框架类
  • vuepress搭建
  • hexo搭建
  • 云图
  • llm wiki

    • 基于karpathy
    • gradle
  • 常用工具

    • git
    • gradle
    • Zadig
    • it-tools
    • 开源推荐
    • curl
  • 大前端

    • 环境配置
    • 微信生态
    • 正则
    • 全栈技能
  • java圈

    • java基础
    • jdk体系
    • jvm
    • spring框架
    • 分库分表
    • zookeeper
  • python技能

    • python编程
    • python数据
  • 算法

    • 算法
  • 网关

    • spring_cloud_gateway
    • openresty
  • 高可用

    • 秒杀
    • 分布式
    • 缓存一致
  • MQ

    • MQ
    • rabbitMQ
    • rocketMQ
    • kafka
  • 其它

    • 设计模式
    • 领域驱动(ddd)
  • 关系型数据库

    • mysql5.0
    • mysql8.0
  • 非关系型数据库

    • redis
    • mongoDB
  • 分布式/其他

    • ShardingSphere
    • 区块链
  • 向量数据库

    • M3E
    • OPEN AI
  • Jmeter
  • fiddler
  • wireshark
  • AI入门
  • AI大模型
  • AI插件
  • AI集成框架
  • 相关算法
  • AI训练师
  • 量化交易
  • AIoT
  • gitee
  • github
  • infoq
  • osc
  • 砍材工具
  • 关于
  • 相关运营
  • devops
  • 元宇宙
  • 区块链
  • 物联网
  • webrtc
  • web3.0
  • gitee
  • github
  • infoq
  • osc
  • 砍材工具
  • 关于
  • 中考
  • 投资
  • 保险
  • 思
  • 首页
    • 开发板介绍
    • micropython环境搭建
    • esp32开发板
    • 面包板
    • 万能表使用
  • 面包板
    • 点灯
  • esp32
    • 点亮开发板led灯
    • 点亮外接led
    • 点亮外接oled文字
    • 红外传感器
    • 红外传感器+olde
    • esp32+面包板
  • MQTT编程

    • MQTT入门

      • 物联网 MQTT
      • 物联网 MQTT和Socket
      • 物联网 MQTT订阅性能优势
      • 物联网 MQTT简易版Broker
    • HiveMQ

      • hivemq实战入门
    • Protobuf

      • Protobuf入门
      • Protobuf入门+梳理
      • Protobuf实战第一篇
    • emqx

      • emqx入门
      • emqx配置开启tls
      • emqx自动重连
      • emqx订阅去重
      • emqx解析QoS 0、1、2
    • mica

      • mica入门
    • netty

      • 入门

        • 基于netty构建入门
        • 理解粘包/拆包
        • 编解码器机制与自定义协议
        • 心跳和ack机制
        • mqtt服务demo演示
        • mqtt服务协议支持
        • mqtt服务udp支持
      • 协议规范

        • mqtt协议规范(发布/订阅模式)
        • mqtt协议规范(轻量级二进制协议)
        • mqtt协议规范(三种 QoS 等级)
        • mqtt协议规范(主题通配符订阅)
        • mqtt协议规范(遗嘱与保留消息)
      • 报文结构

        • 控制报文结构(报文分类)
        • 控制报文结构(连接与握手)
        • 控制报文结构(发布与接收)
      • 核心实战

        • 核心实战(握手与认证)
        • 核心实战(心跳保活机制)
        • 核心实战(会话管理)
        • 核心实战(安全)
    • mqtt-模拟器|客户端

      • 集成Paho

        • 设备模拟器设计
        • 设备模拟器演示
        • Paho拆解入门
        • Paho拆解核心
        • Paho拆解高性能
        • 其他客户端框架比较
      • NetAssist

        • 设备模拟器设计
      • hiveMq

        • hiveMq客户端
    • netty-mqtt-boot

      • 模块化设计
      • 统一接入层
      • 消息路由与流转层
      • 核心服务层
      • 业务应用层
      • 整体项目管理
      • 测试脚手架
      • 兼容支持
    • mqtt-压测

      • mqtt-jmeter

        • 模块化设计
      • ‌emqtt-bench

        • 模块化设计
    • mqtt-规则引擎

      • MQTT规则引擎
      • sql

        • 基于sql规则
      • ice

        • 模块化设计
      • Aviator

        • 模块化设计
      • Drools

        • 模块化设计
    • mqtt-实战问题

      • MQTT分布式集群

        • 分布式集群

          • 分布式集群意义
          • 如何构建分布式集群
          • 分布式集群简易版
          • mqtt入口HAProxy模拟负载
        • ignite

          • ignite入门
          • 模块化设计
      • 实战汇总

        • netty集成udp和tcp编码
        • mqtt集成tls
  • mqtt-netty源码拆解

    • mqtt-netty源码拆解

      • Pipeline双向链表
      • MqttEncoder.INSTANCE
      • 分布式集群简易版
      • mqtt入口HAProxy模拟负载
  • 物联网实战|Spring Boot MQTT平台|解决产线emqx自动重连
    • 客户端网络波动断开
    • 源码
    • pom依赖
    • 相关参数说明
      • 自动重连参数
      • 会话超时时间参数
      • cleanStart参数
      • clientId
      • 会话配置 和 重连时机
        • 会话在有效期内重连 → 订阅不丢失,无需重新订阅
        • 会话过期或 clientId 变了 → 订阅丢失,必须重新订阅
    • 部分代码
      • 结果

物联网实战|Spring Boot MQTT平台|解决产线emqx自动重连

基于emqx v6.2.2 版本

客户端网络波动断开

基于hivemq sdk,版本为1.3.15

源码

https://gitee.com/kcnf-iot/iot-sample/tree/master/emqx/emqx-sample-03

pom依赖

<dependencies>
    <!-- HiveMQ MQTT5 客户端(同时兼容 MQTT3.1.1) -->
    <dependency>
        <groupId>com.hivemq</groupId>
        <artifactId>hivemq-mqtt-client</artifactId>
        <version>1.3.15</version>
    </dependency>
</dependencies>

相关参数说明

自动重连参数

  • 自动重连频率
.automaticReconnect()
    .initialDelay(1, TimeUnit.SECONDS)   // ① 首次重连等待 1 秒
    .maxDelay(5, TimeUnit.MINUTES)       // ② 最大重连间隔封顶 5 分钟
    .applyAutomaticReconnect()           // ③ 应用上述配置
  • 自动重连频率效果
第 1 次重连 → 等待 1 秒
第 2 次重连 → 等待 2 秒(翻倍)
第 3 次重连 → 等待 4 秒
第 4 次重连 → 等待 8 秒
第 5 次重连 → 等待 16 秒
第 6 次重连 → 等待 32 秒
第 7 次重连 → 等待 64 秒
...
第 N 次重连 → 等待 5 分钟(封顶,不再增长)

会话超时时间参数

.sessionExpiryInterval(5, TimeUnit.MINUTES)
属性含义
sessionExpiryIntervalMQTT 5 的会话过期时间。客户端断开后,服务端保留该会话的订阅关系和未送达消息的最长时间
设为 5 分钟断网或客户端异常退出后,5 分钟内重连可以恢复原有会话(sessionPresent=true),无需重新订阅,也不会丢失离线期间的消息

cleanStart参数

.cleanStart(false)                              // 尝试恢复已有会话
.sessionExpiryInterval(TimeUnit.MINUTES.toSeconds(5))  // 断开后保留会话 5 分钟
参数值作用
cleanStart(false)false连接时尝试恢复与该 clientId 关联的历史会话(订阅列表、未确认消息)
sessionExpiryInterval(300)300 秒断开后服务端保留会话状态的最长时间

clientId

场景建议
希望应用重启也能恢复会话把 clientId 改为固定值(如 "server_" + mqttProperties.getClientId()),配合 cleanStart(false)
每次重启都是全新会话保持随机 clientId,但把 cleanStart 改为 true,语义更清晰
当前代码现状cleanStart(false) 对随机 clientId 无实际收益,但也不会报错(sessionPresent 始终为 false)

会话配置 和 重连时机

会话在有效期内重连 → 订阅不丢失,无需重新订阅

条件代码
clientId 固定不变❌ 随机 UUID,每次运行不同
sessionExpiryInterval > 0✅ 设了 5 分钟
在 sessionExpiryInterval 时间内重连✅ automaticReconnect 通常几秒到几十秒就重连

会话过期或 clientId 变了 → 订阅丢失,必须重新订阅

场景结果
断线超过 sessionExpiryInterval(5 分钟)服务端清除会话,订阅丢失
应用重启( clientId 是随机 UUID)新 clientId,找不到旧会话,订阅丢失
服务端主动踢掉并清除了会话订阅丢失

部分代码

package com.jysemel.iot.client;

import com.hivemq.client.mqtt.MqttClientSslConfig;
import com.hivemq.client.mqtt.lifecycle.MqttDisconnectSource;
import com.hivemq.client.mqtt.mqtt5.Mqtt5AsyncClient;
import com.hivemq.client.mqtt.mqtt5.Mqtt5Client;
import com.hivemq.client.mqtt.mqtt5.message.connect.Mqtt5Connect;
import lombok.extern.slf4j.Slf4j;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

import javax.net.ssl.TrustManagerFactory;
import java.io.InputStream;
import java.security.KeyStore;
import java.util.UUID;
import java.util.concurrent.TimeUnit;

@Slf4j
@Configuration
public class MqttConfig {

    @Bean
    public  Mqtt5AsyncClient mqtt5AsyncClient() {
        // 1. 加载CA根证书,构建SSL信任配置
        MqttClientSslConfig sslConfig = buildSslConfig();
        // 2. 构建MQTT TLS客户端
        String clientId = "server_" + UUID.randomUUID().toString().substring(0, 16);
        // 认证信息必须配置在 builder 级别,否则 automaticReconnect 重连时会丢失
        // 3. 构建连接报文
        Mqtt5AsyncClient client = Mqtt5Client.builder()
                .identifier(clientId)
                .serverHost("192.168.0.19")
                .serverPort(18085)
                .sslConfig(sslConfig)
                // 认证信息放在 builder 里,重连时会自动携带
                .simpleAuth()
                .username("emqx")
                .password("123456".getBytes())
                .applySimpleAuth()
                // 自动重连配置
                .automaticReconnect()
                .initialDelay(1, TimeUnit.SECONDS)
                .maxDelay(5, TimeUnit.MINUTES)
                .applyAutomaticReconnect()
                // 连接成功监听器(包括重连成功)
                .addConnectedListener(context -> {
                    log.info("MQTT 连接已建立 | clientId={}", clientId);
                })
                // 断开连接监听器
                .addDisconnectedListener(context -> {
                    Throwable cause = context.getCause();
                    MqttDisconnectSource source = context.getSource();
                    if (cause != null) {
                        log.error("MQTT 连接断开 | clientId={} | source={} | 原因: {} | 自动重连中...", clientId, source, cause.getMessage(), cause);
                    } else if (source == MqttDisconnectSource.SERVER) {
                        log.warn("MQTT 被服务端断开 | clientId={} | source={} | 自动重连中...", clientId, source);
                    } else if (source == MqttDisconnectSource.USER) {
                        log.info("MQTT 用户主动断开 | clientId={} | source={}", clientId, source);
                    } else {
                        log.warn("MQTT 连接断开 | clientId={} | source={} | 自动重连中...", clientId, source);
                    }
                    // 用户主动断开(如应用关闭)不推送告警,仅异常断开推送
                    if (source == MqttDisconnectSource.USER) {
                        return;
                    }
                })
                .buildAsync();
        // 4. 发起异步连接
        Mqtt5Connect connectMsg = Mqtt5Connect.builder()
                // 尝试恢复已有会话
                .cleanStart(false)
                // 会话过期时间 5 分钟:断开后服务端保留会话状态(订阅、未确认消息)300 秒
                .sessionExpiryInterval(TimeUnit.MINUTES.toSeconds(5))
                .keepAlive(60)
                .build();
        // 4. 发起异步连接
        client.connect(connectMsg)
                .whenComplete((connAck, throwable) -> {
                    if (throwable != null) {
                        log.error("MQTT 初始连接失败 | clientId={} |   | 错误: {}", clientId , throwable.getMessage(), throwable);
                        return;
                    }
                    log.info("MQTT 初始连接成功 | clientId={} | reasonCode={} | sessionPresent={}", clientId, connAck.getReasonCode(), connAck.isSessionPresent());
                });

        return client;
    }

    private MqttClientSslConfig buildSslConfig() {
        try {
            KeyStore trustStore = KeyStore.getInstance(KeyStore.getDefaultType());
            try (InputStream is = MqttConfig.class.getResourceAsStream("/cert/emqx.jks")) {
                if (is == null) {
                    throw new RuntimeException("信任库文件未找到: /cert/emqx.jks");
                }
                trustStore.load(is, "123456".toCharArray());
            }
            TrustManagerFactory tmf = TrustManagerFactory.getInstance(TrustManagerFactory.getDefaultAlgorithm());
            tmf.init(trustStore);
            return MqttClientSslConfig.builder()
                    .trustManagerFactory(tmf)
                    // 关键:自定义 HostnameVerifier,总是验证通过(仅测试用!)
                    .hostnameVerifier((hostname, session) -> true)
                    .build();
        } catch (Exception e) {
            throw new RuntimeException("SSL 初始化失败", e);
        }
    }
}

结果

  • 图一 img

  • 图二 img

  • 图三 img

  • 图四 img

  • 图五 img

  • 图六 img

最近更新: 2026/8/17 13:40
Contributors: kcnf
Prev
emqx配置开启tls
Next
emqx订阅去重