储能云平台集成MQTT总结
一、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支持同时接入两种版本的客户端
五、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 唯一,落实消息幂等,规避线程阻塞、会话资源堆积风险,保障储能平台设备上报与指令下发通讯稳定。