储能云平台集成MQTT总结

lyy 8 阅读

一、MQTT核心概念

MQTT(Message Queuing Telemetry Transport,消息队列遥测传输),轻量级、基于发布 / 订阅模式的物联网消息协议,运行在 TCP/IP 之上,适合带宽小、设备算力弱、网络不稳定的物联网设备。

1、三大角色

Publisher 发布者:产生消息,向 Broker 发送消息。

Broker 服务端(MQTT 服务器):中间枢纽,接收发布者消息,过滤,分发给对应订阅者。常见:EMQX、ApsaraMQ for MQTT、TDMQ MQTT版。

Subscriber 订阅者:向 Broker 订阅 Topic,接收推送过来的消息。

2、Topic 主题

消息的路由分类标签,相当于消息的地址。发布者向指定Topic发送消息,Broker依据Topic完成消息过滤分发,订阅该Topic的订阅者就可以收到消息,实现发布者和订阅者完全解耦。

3、消息 Message

一条MQTT消息的组成:

Topic:发到哪个主题

Payload:消息体,二进制,通常 JSON 字符串,MQTT 不解析内容

QoS:服务质量等级

Retain 保留标志

4、ClientId 客户端 ID

每个连接到 Broker 的设备必须有 ClientId,唯一标识一台客户端。

Clean Session(清理会话)

cleanSession=true:断开连接,服务器清除该客户端所有订阅、离线消息。一般临时调试客户端。

cleanSession=false:持久会话,断开后 Broker 保存订阅、保存离线消息(受遗嘱、离线消息最大队列限制),重连后补发离线消息。物联网设备常用 false。

5. Will 遗嘱消息(Last Will)

客户端连接 Broker 时预先设置遗嘱:

当客户端异常断开(断电、断网,没有正常发 disconnect),Broker 自动把遗嘱消息发送到指定 Topic

正常调用 disconnect 断开,遗嘱不会触发。

用途:设备离线告警,其他设备感知设备掉线。遗嘱也可以设置 QoS、Retain。

6. KeepAlive 心跳保活

连接时设置心跳时间(秒)。

客户端必须在心跳时间内发送 PINGREQ 心跳包;

Broker 超过 1.5 倍 keepalive 没收到心跳,判定客户端异常断开,触发遗嘱消息。

二、MQTT Topic 主题说明

MQTT 主题是 MQTT 协议中实现消息路由和过滤的核心机制,本质上是一个 UTF-8 编码的字符串。它充当了消息的“地址”或“标签”,使得发布者和订阅者能够完全解耦。Broker(代理服务器)根据主题将消息从发布者精准分发给所有匹配该主题的订阅者。

核心机制

‌Topic Name(主题名):发布消息时使用,‌严禁包含通配符‌,必须是确定的路径。

‌Topic Filter(主题过滤器):订阅时使用,‌可使用通配符‌,作为匹配模板。

‌无需预先创建‌:客户端发布或订阅时自动创建主题。‌‌

命名规范

区分大小写‌:myhome/temp 与 MyHome/Temp 是不同主题。

‌避免空格‌及非 ASCII 特殊字符,建议使用小写字母、数字、下划线、短横线。

‌层级建议 ≤ 7 层‌,主题长度建议 ≤ 128 字节(协议上限 65535 字节)。

‌避免尾部斜杠‌:smartlock/door/ ≠ smartlock/door。

‌唯一标识前置‌:device/{deviceID}/status 优于 device/status/{deviceID},便于通配符订阅。‌‌

使用说明

1、设备上报(设备 —>平台)

设备上报遥测数据、故障告警,发布时使用固定完整主题名,禁止使用通配符。

示例:device/bms_001/report、device/pcs_002/alarm

设备 ID 放在主题路径靠前位置,便于平台通配符批量订阅。

2、平台下发(平台—>设备)

平台下发启停、参数配置等控制指令,发布到指定设备的固定主题。

示例:device/bms_001/cmd

设备本地订阅该主题接收平台下发指令。

3、订阅侧开发约束

批量接收多设备数据:订阅使用通配符过滤器,例如 device/+/report;

多实例消费负载均衡:采用 MQTT5 共享订阅,例如 $share/emsGroup/device/+/report;

开发红线:发布消息绝对不能带 +、#通配符,通配符仅允许用于订阅过滤器。

三、MQTT QoS 通讯说明

1、MQTT与QoS通讯

(1)QoS 0(最多一次)通讯过程

QoS 0 不存在应答确认,仅单报文单向传输。

发布端直接向 Broker 发送PUBLISH报文。

发布端发送完毕即结束本次通讯,不等待 Broker 任何回复,没有重传机制。

Broker 收到消息后向订阅客户端发送PUBLISH,订阅客户端不需要返回确认报文。

(2)QoS 1(至少一次)通讯过程

发布者‑Broker、Broker‑订阅者两段链路都执行PUBLISH+PUBACK双向两次报文交互。

发布者 ↔ Broker

发布者发送PUBLISH报文。

Broker 接收消息,返回PUBACK确认报文。

发布者超时没有收到PUBACK,会设置 DUP 标记为 1,重复发送PUBLISH,直至收到PUBACK。

Broker ↔ 订阅者

Broker 向订阅者下发PUBLISH报文。

订阅者收到消息,回复PUBACK。

Broker 超时收不到PUBACK,会重复下发带 DUP 标记的PUBLISH报文。

该通讯机制会产生重复消息。

(3) QoS 2(恰好一次)通讯过程

两段链路都采用四次握手报文交互。

发布者 ↔ Broker

发布者发送 QoS2 等级PUBLISH报文。

Broker 接收消息,回复PUBREC报文。

发布者收到PUBREC,发送PUBREL释放报文。

Broker 收到PUBREL,回复PUBCOMP报文,发布链路通讯结束。

某一步报文丢失,上一方会重传对应报文。

Broker ↔ 订阅者

Broker 向订阅者发送PUBLISH报文。

订阅者接收消息,回复PUBREC。

Broker 收到PUBREC后发送PUBREL。

订阅者收到PUBREL,回复PUBCOMP,推送链路通讯结束。

2、客户端业务使用建议

在调用发布 API 时指定 QoS 等级,对应代码中 .qos(1)。

设备上报遥测数据:业务允许少量丢包场景选用 QoS0;重要状态数据、告警信息建议使用 QoS1。

平台下发控制指令:储能 PCS/BMS 控制命令,建议选用 QoS1,保证指令尽可能送达设备。

不建议业务默认使用 QoS2,会增大 Broker 与客户端的交互开销,储能物联网场景极少使用。

3、QoS与消息堆积说明

QoS 本身不会直接产生消息堆积,堆积是 QoS 等级 + MQTT 会话配置共同作用的结果。Topic 仅作为路由地址,本身不会存储消息。

(1)、离线消息堆积(设备断开离线)

产生离线消息必须满足以下3个条件:

消息 QoS 等级为 QoS1 / QoS2;

MQTT5 配置 clean‑start: false;

session‑expiry‑interval > 0,Broker 保留客户端会话。

当全部条件满足,设备离线期间到达的消息会被 Broker 缓存,待设备重连后推送。

(2)、在线消息队列积压

设备保持在线状态,但业务消费处理速度跟不上消息接收速率,会产生发送队列积压。

该情况和clean‑start无关,QoS0、QoS1、QoS2 均有可能出现;

属于业务消费处理能力瓶颈,需要优化消费逻辑、控制消息上报频率。

(3)、在线消息队列积压

设备保持在线状态,但业务消费处理速度跟不上消息接收速率,会产生发送队列积压。

该情况和clean‑start无关,QoS0、QoS1、QoS2 均有可能出现;

属于业务消费处理能力瓶颈,需要优化消费逻辑、控制消息上报频率。

四、MQTT3.1.1 与 MQTT5.0 区别

1 两大协议核心特性差异对比

MQTT3.1.1是物联网设备主流版本,MQTT5.0在其基础上扩展大量能力,两套协议本身不互通,但Broker支持同时接入两种版本的客户端

对比分类

对比项

MQTT3.1.1

MQTT3.5

储能场景业务价值

会话管理

会话控制字段

CleanSession布尔值:true 断开清除会话,false 持久保留会话

CleanStart标识 + Session Expiry Interval会话过期时间

可配置会话超时,自动回收长期离线储能设备会话,避免 Broker 内存持续堆积

会话生命周期

会话一旦建立,无超时回收机制,只能客户端主动断开释放会话资源

支持会话超时自动失效,到达设定时间自动回收会话

适配大规模储能集群,清理无效会话,降低服务器资源占用,便于运维

消息能力

自定义扩展字段

无协议原生用户属性,业务扩展元数据只能封装在 Payload 载荷 JSON 内部

UserProperties自定义键值对,独立于消息载荷存在

鉴权、路由、规则引擎无需解析 JSON,直接获取设备 SN、设备型号等元信息

消息有效期

离线队列消息永久保存,没有消息过期淘汰逻辑

Message Expiry Interval,可配置消息存活有效期

自动丢弃过期采集数据,规避下发过期、无效 PCS/BMS 控制指令

报文格式标记

协议无载荷格式标识

Payload Format IndicatorContent‑Type载荷格式标记

储能设备掉线告警可附带故障扩展信息,丰富离线事件上报内容

遗嘱消息能力

遗嘱仅支持 Topic、QoS、Payload、Retain 基础配置

遗嘱完整支持 UserProperties、消息过期、载荷格式等全套消息属性

储能设备掉线告警可附带故障扩展信息,丰富离线事件上报内容

订阅机制

共享订阅

不属于协议标准,属于各 Broker 厂商私有扩展,不同平台兼容性差

协议原生标准共享订阅,语法$share/group/topic

EMS 后端多实例部署,实现消息负载均衡消费,承载海量设备上报压力

网络传输化

Topic 别名压缩

Topic 别名压缩

Topic Alias,使用数字 ID 替代长 Topic 字符串

4G 窄带宽储能边缘设备,缩短报文长度,节约网络流量开销

异常处理

错误返回码

返回码数量少,报错信息笼统,故障定位困难

细化完整原因码,连接、订阅、发布失败均返回明确错误原因

快速定位设备连接失败、订阅拒绝、权限异常等线上问题,降低排错成本

五、Mica‑MQTT集成说明

1、Mica使用说明

mica‑mqtt 隶属于 Dromara 开源社区,是基于 Java AIO 实现的开源 MQTT 物联网组件,具备简单易用、低延迟、高性能的特性,可方便集成到已有业务服务进行二次开发,降低物联网平台自研成本。组件同时提供 MQTT 客户端 和 MQTT Broker 服务端 两套能力:客户端模式用于对接外部 MQTT 消息代理;服务端模式可作为嵌入式 MQTT 消息代理直接使用

2、集成方式

(1)、Maven依赖注入

<!-- mica‑mqtt springboot3 starter -->
<dependency>
    <groupId>net.dreamlu</groupId>
    <artifactId>mica‑mqtt‑spring‑boot‑starter</artifactId>
    <version>3.2.0</version>
</dependency>

(2)、application.yml 配置适配

mica:
  mqtt:
    client:
      # MQTT broker地址,示例:tcp://127.0.0.1:1883
      url: ${MQTT_BROKER_URL:}
      # 客户端id,业务建议结合服务实例编号,多实例避免重复
      client-id: ems‑cloud‑${spring.application.name}-${server.port}
      username: ${MQTT_BROKER_USERNAME:}
      password: ${MQTT_BROKER_PASSWORD:}
      # 使用 MQTT_5 协议;如需兼容旧设备改为 MQTT_3_1_1
      protocol-version: MQTT_5
      # MQTT5 CleanStart,替代旧版CleanSession
      clean-start: true
      # MQTT5会话过期时间,单位秒
      session-expiry-interval: 3600
      # 心跳保活时间,单位秒
      keep-alive: 60
      # 自动重连开启
      auto-reconnect: true

(3)、MQTT消息发布

import net.dreamlu.iot.mqtt.core.client.MqttClient;
import net.dreamlu.iot.mqtt.spring.client.MqttClientTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;

@Component
public class MqttPublishUtil {

    @Resource
    private MqttClientTemplate mqttClientTemplate;

    /**
     * MQTT5消息发布示例,携带UserProperties、消息过期时间
     */
    public void publishCmd(String topic,String payload){
        mqttClientTemplate.publish()
                .topic(topic)
                .payload(payload.getBytes())
                // 框架层设置QoS
                .qos(1)
                .retain(false)
                // MQTT5 用户自定义属性
                .userProperty("deviceSn","【设备唯一编号】")
                .userProperty("msgType","control")
                // MQTT5消息过期时间,单位秒
                .messageExpiryInterval(120)
                .send();
    }
}

(4)、消息订阅

import net.dreamlu.iot.mqtt.spring.client.annotation.MqttSubscribe;
import net.dreamlu.iot.mqtt.spring.client.message.MqttMessage;
import org.springframework.stereotype.Component;

@Component
public class MqttMessageListener {

    /**
     * 普通主题订阅
     */
    @MqttSubscribe(topic = "【业务上报topic】", qos = 1)
    public void onDeviceReport(MqttMessage message){
        byte[] payload = message.getPayload();
        // 获取MQTT5 UserProperties
        message.getUserProperties().forEach((k,v)->{
            System.out.println(k + ":" + v);
        });
    }

    /**
     * MQTT5原生共享订阅,实现多实例消费负载均衡
     */
    @MqttSubscribe(topic = "$share/【消费组标识】/【业务topic前缀】/#", qos = 1)
    public void onBmsSharedConsume(MqttMessage message){
        // 共享订阅消费逻辑
    }
}

3、集成使用注意事项

客户端ID唯一性:多实例部署时client‑id必须保证唯一,重复 ClientId 会导致老客户端被 Broker 强制踢下线;可结合端口、实例编号生成。

MQTT会话参数注意:clean‑start: false配合session‑expiry‑interval才会生效持久会话;若clean‑start=true会话过期配置不生效。

幂等性处理:QoS=1 场景下会重复收到消息,业务必须做消息幂等去重,不能依赖协议层保证不重复。

版本兼容策略:Broker 支持混合接入 MQTT3.1.1、MQTT5 客户端;新项目设备使用 MQTT5,存量老设备维持 MQTT3.1.1,不需要一次性全部升级

遗属消息配置:业务上设备离线告警遗嘱,建议在 MqttClient 构建时配置,遗嘱同样支持 MQTT5 的 UserProperties 属性。

五、总结

本次完成 MQTT 协议体系梳理与 Mica‑MQTT 客户端完整集成,明确 MQTT 核心概念、Topic 命名规范及 QoS 投递机制,对比 MQTT3.1.1 与 MQTT5.0 能力差异,基于 SpringBoot3.5.15 + JDK21 实现 MQTT5 全能力接入。利用会话过期、自定义属性、共享订阅适配海量设备与多实例消费;同时规范 Topic、QoS 配置,保障 ClientId 唯一,落实消息幂等,规避线程阻塞、会话资源堆积风险,保障储能平台设备上报与指令下发通讯稳定。