这一篇专门回答一个非常现实的问题:MQTT 5.0 到底是不是该直接上。重点会从能力协商、Reason Code、边界控制、Broker 支持度和实现成熟度几个维度,把支持 5.0 和跑稳 5.0 的差别说透。
适合谁收藏
- 正在做 CODESYS / PLC / MQTT 项目的人
- 想把 MQTT 从报文真正看到 ST 代码的人
- 正在排查 QoS1 / QoS2 超时、掉线、重连问题的人
只要一聊 MQTT,很多人都会下意识地问一句:
那是不是直接上 5.0 就最好?
如果站在“新版本功能更多”的角度,这句话好像没毛病。 但如果站在工业现场角度,这句话就太简单了。
先给结论:
MQTT 5.0 当然值得学,也值得用。
一、MQTT 5.0 比 3.1.1 强在哪
最有价值的增强,主要集中在下面几类:
| 能力 | 价值 |
|---|---|
| Properties | 报文语义更丰富 |
| Reason Code | 错误与拒绝原因更清晰 |
| Receive Maximum | 流量控制更明确 |
| Maximum Packet Size | 报文边界更可控 |
| Topic Alias | 降低长主题重复开销 |
| Subscription Identifier | 订阅来源更清晰 |
| Shared Subscription | 多消费者分摊 |
| Session Expiry | 会话管理更灵活 |
如果你只看这张表,会觉得:
5.0 明显更高级。
这没问题。 但高级不等于现场就一定要全开。
二、5.0 真正的工程价值是什么
我认为 5.0 最实际的价值,不是“参数变多了”,而是下面三件事:
1. 能力协商更明确
以前很多事情只能靠“试试看”。 现在 CONNACK 可以明确告诉你:
- 我支持的最大 QoS 是多少
- 我支不支持 Retain
- 我支不支持通配订阅
- 我支不支持共享订阅
- 我支不支持 Subscription Identifier
2. 错误原因更清楚
不是所有失败都只剩一个模糊“失败”。 Reason Code 能让诊断路径更清楚。
3. 大型系统更好控边界
比如:
- Receive Maximum
- Maximum Packet Size
- Session Expiry
这些对大规模系统、长连接系统、跨平台系统更有价值。
三、为什么工业现场不能无脑说“5.0 更先进,所以全上”
因为工程里真正决定你要不要开 5.0 的,不止协议本身。
还包括:
- Broker 支持度
- 第三方客户端支持度
- 你自己的实现成熟度
- 现场到底有没有这些需求
举几个很现实的例子:
场景 1
你只是 PLC 周期性上报几组数据,QoS0 / QoS1 足够,Broker 也比较传统。
这种场景,3.1.1 完全可能已经够用了。
场景 2
你要做复杂路由、大量订阅管理、多客户端能力协商、边界控制。
这时 5.0 的优势就会明显起来。
所以正确问法不是:
5.0 高不高级?
而是:
我的系统需不需要它这些能力?
四、在这套开源客户端里,5.0 体现在哪些地方
不是嘴上说支持 5.0,就叫支持。 真正要看代码里有没有落点。
当前这套 MqttClient_V1_0,5.0 的关键落点包括:
| 能力 | 代码落点 |
|---|---|
| CONNECT 属性 | M_BuildConnectPacket |
| CONNACK 属性解析 | M_HandleConnAck |
| 发布 Topic Alias | M_BuildPublishPacket |
| 订阅标识符 | M_BuildSubscribePacket |
| 服务端能力约束 | M_HandleConnAck 后续影响发布订阅逻辑 |
这说明当前这套实现不是“挂个 5.0 枚举值做样子”,而是真把核心协商链接进去了。
五、“支持 5.0” 和 “跑稳 5.0” 到底差在哪
这个差别特别重要。
支持 5.0
通常只代表:
- 能发 5.0 CONNECT
- 能收 5.0 CONNACK
- 能识别一些属性和原因码
跑稳 5.0
还要额外满足:
- 这些属性真的被后续逻辑正确使用
- 不支持的能力不会硬发
- 多客户端、多主题、高频下仍然稳定
- 边界条件和 Broker 差异被处理掉
所以如果一个项目只是:
- 版本号可选 5.0
- CONNECT 发的是 5.0
这只能叫“初步支持 5.0”,离“稳定可用的 5.0 客户端”还差一截。
六、我对工业现场的建议很直接
建议 1:先问需求,不要先问版本
如果你只是:
- 基础发布订阅
- 几个主题
- 不需要高级属性
那 3.1.1 可能完全够。
建议 2:真要上 5.0,就先把 Broker 能力摸清
重点看:
- 是否完整支持 5.0
- 支持到哪些属性
- 是否支持共享订阅
- 是否支持订阅标识符
建议 3:上 5.0 之前,先保证 QoS1 / QoS2 链已经稳定
因为 5.0 不是拿来替你补状态机的。 基础 ACK 链都没跑稳,再多属性也没用。
七、哪些场景我认为 5.0 真的值得上
| 场景 | 原因 |
|---|---|
| 多客户端复杂协同 | 能力协商和订阅增强更有价值 |
| 需要更细诊断 | Reason Code 更清楚 |
| 长主题高频发布 | Topic Alias 有价值 |
| 需要严格会话策略 | Session Expiry 更灵活 |
| 需要更明确边界控制 | Receive Maximum / Max Packet Size 更实用 |
八、哪些场景不要为了“新”而硬上
| 场景 | 建议 |
|---|---|
| 只做最基础点对点上报 | 3.1.1 先跑稳 |
| Broker 很老或环境复杂 | 先保证兼容性 |
| 客户端实现刚起步 | 先把 ACK 链和重连链站稳 |
| 团队没人懂 5.0 属性含义 | 先别把复杂度引进来 |
九、这一篇你最该记住的 5 句话
- MQTT 5.0 值得学,也值得用,但不该无脑上。
- 真正该问的不是“支不支持 5.0”,而是“支持到哪、稳不稳、用不用得上”。
- 5.0 最有价值的地方,是能力协商更清楚、边界更明确。
- 版本号能切到 5.0,不等于你的客户端已经稳定支持 5.0。
- 基础 ACK 链都没跑稳时,别指望 5.0 特性能替你兜底。
十、下篇预告
下一篇是最后一篇加更。 我们不再从协议讲,而是从项目本身讲:
这套开源 MQTT 客户端,到底是怎么一步一步做出来的
这篇会更偏源码演进、踩坑复盘和项目思路总结。
完整 ST 代码
复制使用说明
- 这部分给出的是与本篇主题直接对应的完整 ST 代码,不是零碎片段。
- 如果你只是想先跑通,优先整段复制,不要只摘几行变量或几条赋值语句。
- 如果是
METHOD,请确认它仍然属于FB_MqttClient;如果是PROGRAM,请确认相关 DUT、GVL、FB 已一并导入。
代码阅读重点
- 先按
报文结构 -> 状态机入口 -> 关键变量 -> 返回结果的顺序看。 - 再把正文里的十六进制拆解和这里的字节写入、字节解析语句一行行对上。
- 最后回到在线调试,重点盯
uiTxLength、uiRxLength、eState、xWaitingForAck这类状态量。
完整代码 1:M_BuildConnectPacket
- 对应源码路径:
10 MQTT/MqttClient_V1_0/Device/Application/MQTT/POUs/MqttClient NBS/FB_MqttClient/处理发送报文/M_BuildConnectPacket.st - 复制使用说明:MQTT 5.0 升级先从 CONNECT 看起,因为属性协商都是从这里开始的。
- 阅读重点:重点看 Session Expiry、Receive Maximum、Maximum Packet Size、Topic Alias Maximum 这些 5.0 属性是怎么真正落进字节流的。
/// =======================================================================
/// 名称 : M_BuildConnectPacket
/// 功能 : 构建 CONNECT 发送报文
/// 说明 : 按 MQTT 3.1.1/5.0 规则组装 CONNECT 报文并写入发送缓冲区。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-05
/// 版本 : V1.0
/// =======================================================================
{attribute 'hide_all_locals'}
METHOD M_BuildConnectPacket : BOOL
VAR
uiPos : UINT := 0; //字节索引
uiVarHeaderLen : UINT; //可变报头长度
uiPayloadLen : UINT; //载荷长度
uiRemainingLen : UINT; //剩余长度
uiPropsLen : UINT; //MQTT 5.0属性长度
i : DINT;
byConnectFlags : BYTE; //连接标志字节
uiTopicAliasMax : UINT; //客户端接收方向主题别名上限
END_VAR
// === IMPLEMENTATION ===
// BUG-10: 缓冲区溢出保护
IF SIZEOF(aTxBuf) < 256 THEN
M_BuildConnectPacket := FALSE;
RETURN;
END_IF
IF (sPassword <> '') AND (sUsername = '') THEN
M_SetError(
uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
sMessage := 'Password requires username');
M_BuildConnectPacket := FALSE;
RETURN;
END_IF
IF (NOT bWillFlag) AND ((sWillTopic <> '') OR (sWillMessage <> '')) THEN
M_SetError(
uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
sMessage := 'Will fields require Will Flag');
M_BuildConnectPacket := FALSE;
RETURN;
END_IF
IF (bWillRetain OR (eWillQoS <> E_MqttQoS.byQoS0)) AND (NOT bWillFlag) THEN
M_SetError(
uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
sMessage := 'Invalid will flag combination');
M_BuildConnectPacket := FALSE;
RETURN;
END_IF
IF bWillFlag AND (sWillTopic = '') THEN
M_SetError(
uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
sMessage := 'Will topic is required');
M_BuildConnectPacket := FALSE;
RETURN;
END_IF
// 清空发送缓冲区,避免数据干扰
FOR i := LOWER_BOUND(aTxBuf, 1) TO UPPER_BOUND(aTxBuf, 1) DO
aTxBuf[i] := 0;
END_FOR
/// =======================================================================
/// 长度计算
/// =======================================================================
// 可变报头长度: 协议名(6) + 协议版本(1) + 连接标志(1) + 保持连接(2)
uiVarHeaderLen := 10;
// MQTT 5.0: 连接属性长度
uiPropsLen := 0;
IF eVersion = E_MqttVersion.byMqttVersion50 THEN
// Session Expiry Interval (4字节值 + 1字节标识符)
IF udiSessionExpiry > 0 THEN
uiPropsLen := uiPropsLen + 5;
END_IF
// Receive Maximum (2字节值 + 1字节标识符)
uiPropsLen := uiPropsLen + 3;
// Maximum Packet Size (4字节值 + 1字节标识符)
uiPropsLen := uiPropsLen + 5;
// Topic Alias Maximum (2字节值 + 1字节标识符)
uiPropsLen := uiPropsLen + 3;
// Request Problem Information (1字节值 + 1字节标识符)
uiPropsLen := uiPropsLen + 2;
// Request Response Information (1字节值 + 1字节标识符)
IF bRequestResponseInfo THEN
uiPropsLen := uiPropsLen + 2;
END_IF
// 属性长度字段的编码长度
IF uiPropsLen < 128 THEN
uiVarHeaderLen := uiVarHeaderLen + 1 + uiPropsLen; //1字节属性长度
ELSE
uiVarHeaderLen := uiVarHeaderLen + 2 + uiPropsLen; //2字节属性长度
END_IF
END_IF
// 载荷长度
uiPayloadLen := TO_UINT(2 + LEN(sClientID)); //Client ID
IF sUsername <> '' THEN
uiPayloadLen := uiPayloadLen + TO_UINT(2 + LEN(sUsername));
END_IF
IF sPassword <> '' THEN
uiPayloadLen := uiPayloadLen + TO_UINT(2 + LEN(sPassword));
END_IF
IF bWillFlag THEN
IF eVersion = E_MqttVersion.byMqttVersion50 THEN
uiPayloadLen := uiPayloadLen + 1;
END_IF
uiPayloadLen := uiPayloadLen + TO_UINT(2 + LEN(sWillTopic) + 2 + LEN(sWillMessage));
END_IF
IF uiReceiveMax = 0 THEN
M_SetError(
uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
sMessage := 'Receive Maximum must not be zero');
M_BuildConnectPacket := FALSE;
RETURN;
END_IF
IF udMaxPacketSize = 0 THEN
M_SetError(
uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
sMessage := 'Maximum Packet Size must not be zero');
M_BuildConnectPacket := FALSE;
RETURN;
END_IF
// 剩余长度
uiRemainingLen := uiVarHeaderLen + uiPayloadLen;
// 缓冲区长度检查
IF uiRemainingLen + 5 > SIZEOF(aTxBuf) THEN
M_BuildConnectPacket := FALSE;
RETURN;
END_IF
/// =======================================================================
/// 创建报文
/// =======================================================================
uiPos := 0;
// ****************** 固定报文头 = 报文类型 + 剩余长度 ******************
// 报文类型
aTxBuf[0] := E_MqttPacketType.byConnect;
uiPos := uiPos + 1;
// 剩余长度
uiPos := uiPos + M_EncodeRemainingLength(uiRemainingLen, ADR(aTxBuf[uiPos]));
// ****************** 可变报文头 = 协议名(6) + 协议版本(1) + 连接标志(1) + 保持连接(2) ******************
// 协议名:"MQTT"
aTxBuf[uiPos] := 0; //长度高字节MSB
uiPos := uiPos + 1;
aTxBuf[uiPos] := 4; //长度低字节LSB
uiPos := uiPos + 1;
aTxBuf[uiPos] := 77; // 'M'
uiPos := uiPos + 1;
aTxBuf[uiPos] := 81; // 'Q'
uiPos := uiPos + 1;
aTxBuf[uiPos] := 84; // 'T'
uiPos := uiPos + 1;
aTxBuf[uiPos] := 84; // 'T'
uiPos := uiPos + 1;
// 协议版本(根据eVersion选择)
IF eVersion = E_MqttVersion.byMqttVersion50 THEN
aTxBuf[uiPos] := 5; //MQTT 5.0
ELSE
aTxBuf[uiPos] := 4; //MQTT 3.1.1
END_IF
uiPos := uiPos + 1;
// 连接标志
byConnectFlags :=
(BOOL_TO_BYTE(FALSE) * 1) OR //Bit 0: Reserved。服务端必须验证保留标志位是否为0
(BOOL_TO_BYTE(bCleanSession) * 2) OR //Bit 1: Clean Session
(BOOL_TO_BYTE(bWillFlag AND sWillTopic <> '') * 4) OR //Bit 2: Will Flag
(BOOL_TO_BYTE(eWillQoS = 1 OR eWillQoS = 3) * 8) OR //Bit 3: Will QoS低位
(BOOL_TO_BYTE(eWillQoS = 2 OR eWillQoS = 3) * 16) OR //Bit 4: Will QoS高位
(BOOL_TO_BYTE(bWillFlag AND bWillRetain) * 32) OR //Bit 5: Will Retain
(BOOL_TO_BYTE(sPassword <> '') * 64) OR //Bit 6: Password Flag
(BOOL_TO_BYTE(sUsername <> '') * 128); //Bit 7: Username Flag
aTxBuf[uiPos] := byConnectFlags;
uiPos := uiPos + 1;
// 保持连接
aTxBuf[uiPos] := UINT_TO_BYTE(SHR(uiKeepAlive, 8)); //KeepAlive高字节
uiPos := uiPos + 1;
aTxBuf[uiPos] := UINT_TO_BYTE(uiKeepAlive AND 16#FF); //KeepAlive低字节
uiPos := uiPos + 1;
// ****************** MQTT 5.0: 连接属性 ******************
IF eVersion = E_MqttVersion.byMqttVersion50 THEN
// 属性长度(可变字节编码)
IF uiPropsLen < 128 THEN
aTxBuf[uiPos] := TO_BYTE(uiPropsLen);
uiPos := uiPos + 1;
ELSE
aTxBuf[uiPos] := TO_BYTE((uiPropsLen / 128) OR 16#80);
uiPos := uiPos + 1;
aTxBuf[uiPos] := TO_BYTE(uiPropsLen MOD 128);
uiPos := uiPos + 1;
END_IF
// Session Expiry Interval
IF udiSessionExpiry > 0 THEN
aTxBuf[uiPos] := GVL_Mqtt.cnPropSessionExpiry; uiPos := uiPos + 1;
aTxBuf[uiPos] := UDINT_TO_BYTE(SHR(udiSessionExpiry, 24) AND 16#FF); uiPos := uiPos + 1;
aTxBuf[uiPos] := UDINT_TO_BYTE(SHR(udiSessionExpiry, 16) AND 16#FF); uiPos := uiPos + 1;
aTxBuf[uiPos] := UDINT_TO_BYTE(SHR(udiSessionExpiry, 8) AND 16#FF); uiPos := uiPos + 1;
aTxBuf[uiPos] := UDINT_TO_BYTE(udiSessionExpiry AND 16#FF); uiPos := uiPos + 1;
END_IF
// Receive Maximum
aTxBuf[uiPos] := GVL_Mqtt.cnPropReceiveMaximum; uiPos := uiPos + 1;
aTxBuf[uiPos] := UINT_TO_BYTE(SHR(uiReceiveMax, 8)); uiPos := uiPos + 1;
aTxBuf[uiPos] := UINT_TO_BYTE(uiReceiveMax AND 16#FF); uiPos := uiPos + 1;
// Maximum Packet Size
aTxBuf[uiPos] := GVL_Mqtt.cnPropMaxPacketSize; uiPos := uiPos + 1;
aTxBuf[uiPos] := UDINT_TO_BYTE(SHR(udMaxPacketSize, 24) AND 16#FF); uiPos := uiPos + 1;
aTxBuf[uiPos] := UDINT_TO_BYTE(SHR(udMaxPacketSize, 16) AND 16#FF); uiPos := uiPos + 1;
aTxBuf[uiPos] := UDINT_TO_BYTE(SHR(udMaxPacketSize, 8) AND 16#FF); uiPos := uiPos + 1;
aTxBuf[uiPos] := UDINT_TO_BYTE(udMaxPacketSize AND 16#FF); uiPos := uiPos + 1;
// Topic Alias Maximum
uiTopicAliasMax := GVL_Mqtt.cnMaxTopicAlias;
aTxBuf[uiPos] := GVL_Mqtt.cnPropTopicAliasMax; uiPos := uiPos + 1;
aTxBuf[uiPos] := UINT_TO_BYTE(SHR(uiTopicAliasMax, 8)); uiPos := uiPos + 1;
aTxBuf[uiPos] := UINT_TO_BYTE(uiTopicAliasMax AND 16#FF); uiPos := uiPos + 1;
// Request Problem Information
aTxBuf[uiPos] := GVL_Mqtt.cnPropRequestProblemInfo; uiPos := uiPos + 1;
aTxBuf[uiPos] := BOOL_TO_BYTE(bRequestProblemInfo); uiPos := uiPos + 1;
// Request Response Information
IF bRequestResponseInfo THEN
aTxBuf[uiPos] := GVL_Mqtt.cnPropRequestResponseInfo; uiPos := uiPos + 1;
aTxBuf[uiPos] := 1; uiPos := uiPos + 1;
END_IF
END_IF
// ****************** 载荷 = 客户端标识符 + 遗嘱主题 + 遗嘱消息 + 用户名 + 密码 ******************
// 客户端标识符
uiPos := uiPos + M_AppendString(sClientID, ADR(aTxBuf[uiPos]));
// 遗嘱主题 + 遗嘱消息
IF bWillFlag THEN
IF eVersion = E_MqttVersion.byMqttVersion50 THEN
aTxBuf[uiPos] := 0;
uiPos := uiPos + 1;
END_IF
uiPos := uiPos + M_AppendString(sWillTopic, ADR(aTxBuf[uiPos]));
uiPos := uiPos + M_AppendString(sWillMessage, ADR(aTxBuf[uiPos]));
END_IF
// 用户名
IF sUsername <> '' THEN
uiPos := uiPos + M_AppendString(sUsername, ADR(aTxBuf[uiPos]));
END_IF
// 密码
IF sPassword <> '' THEN
uiPos := uiPos + M_AppendString(sPassword, ADR(aTxBuf[uiPos]));
END_IF
uiTxLength := uiPos; //需要发送的报文总字节数
M_BuildConnectPacket := TRUE;完整代码 2:M_HandleConnAck
- 对应源码路径:
10 MQTT/MqttClient_V1_0/Device/Application/MQTT/POUs/MqttClient NBS/FB_MqttClient/处理接收报文/M_HandleConnAck.st - 复制使用说明:5.0 真正难的不是发属性,而是正确吃掉 Broker 回来的属性协商结果。
- 阅读重点:重点看服务端能力字段被解析后怎么落到
uiServerReceiveMax、byServerMaxQoS、bServerSharedSubAvail这些运行变量里。
/// =======================================================================
/// 名称 : M_HandleConnAck
/// 功能 : 处理 CONNACK 接收报文
/// 说明 : 解析连接确认结果与服务器属性,成功后清除等待确认状态。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-05
/// 版本 : V1.1
/// =======================================================================
{attribute 'hide_all_locals'}
METHOD M_HandleConnAck : BOOL
VAR
byReasonCode : BYTE; // CONNACK 原因码
byConnAckFlags : BYTE; // CONNACK 标志位
uiPropsLen : UINT; // 属性总长度
uiPropertyHeaderBytes : UINT; // 属性长度字段字节数
uiPropsStart : UINT; // 属性起始偏移
uiPropsEnd : UINT; // 属性结束偏移
uiPos : UINT; // 当前解析偏移
byPropId : BYTE; // 当前属性标识符
uiLoopGuard : UINT; // 循环安全计数器
END_VAR
// === IMPLEMENTATION ===
IF uiRxLength < 4 THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
IF NOT xWaitingForAck OR aRxBuf[0] <> byExpectedMsgType THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
byConnAckFlags := aRxBuf[2];
IF (byConnAckFlags AND 16#FE) <> 0 THEN
sDiagMsg := 'Invalid CONNACK flags';
M_HandleConnAck := FALSE;
RETURN;
END_IF
IF eVersion = E_MqttVersion.byMqttVersion50 THEN
byReasonCode := aRxBuf[3];
ELSE
byReasonCode := aRxBuf[3];
END_IF
IF byReasonCode <> 0 THEN
IF eVersion = E_MqttVersion.byMqttVersion50 THEN
CASE byReasonCode OF
E_MqttReasonCode.byUnsupportedProtocolVersion: sDiagMsg := 'Protocol version not supported';
E_MqttReasonCode.byClientIdentifierInvalid: sDiagMsg := 'Client identifier invalid';
E_MqttReasonCode.byBadUserNamePassword: sDiagMsg := 'Bad username or password';
E_MqttReasonCode.byNotAuthorized: sDiagMsg := 'Not authorized';
E_MqttReasonCode.byServerUnavailable: sDiagMsg := 'Server unavailable';
E_MqttReasonCode.byServerBusy: sDiagMsg := 'Server busy';
E_MqttReasonCode.byBanned: sDiagMsg := 'Banned';
E_MqttReasonCode.byServerShuttingDown: sDiagMsg := 'Server shutting down';
E_MqttReasonCode.byBadAuthMethod: sDiagMsg := 'Bad authentication method';
E_MqttReasonCode.byQuotaExceeded: sDiagMsg := 'Quota exceeded';
E_MqttReasonCode.byRetainNotSupported: sDiagMsg := 'Retain not supported';
E_MqttReasonCode.byQosNotSupported: sDiagMsg := 'QoS not supported';
E_MqttReasonCode.byUseAnotherServer: sDiagMsg := 'Use another server';
E_MqttReasonCode.byServerMoved: sDiagMsg := 'Server moved';
ELSE
sDiagMsg := CONCAT('CONNACK reason code: ', BYTE_TO_STRING(byReasonCode));
END_CASE
ELSE
CASE byReasonCode OF
1: sDiagMsg := 'Unacceptable protocol version';
2: sDiagMsg := 'Identifier rejected';
3: sDiagMsg := 'Server unavailable';
4: sDiagMsg := 'Bad username or password';
5: sDiagMsg := 'Not authorized';
ELSE
sDiagMsg := CONCAT('CONNACK return code: ', BYTE_TO_STRING(byReasonCode));
END_CASE
END_IF
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiServerReceiveMax := GVL_Mqtt.cnDefaultReceiveMax;
byServerMaxQoS := TO_BYTE(E_MqttQoS.byQoS2);
bServerRetainAvailable := TRUE;
udServerMaxPacketSize := GVL_Mqtt.cnMaxPacketSize;
uiServerTopicAliasMax := 0;
uiSendQuota := GVL_Mqtt.cnDefaultReceiveMax;
bServerWildcardSubAvail := TRUE;
bServerSubIdAvail := TRUE;
bServerSharedSubAvail := TRUE;
IF eVersion = E_MqttVersion.byMqttVersion50 AND uiRxLength >= 5 THEN
uiPos := 4;
uiPropertyHeaderBytes := M_DecodeRemainingLength(pBuffer := ADR(aRxBuf[uiPos]), uiLength => uiPropsLen);
IF uiPropertyHeaderBytes = 0 THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiPos := uiPos + uiPropertyHeaderBytes;
uiPropsStart := uiPos;
uiPropsEnd := uiPropsStart + uiPropsLen;
IF uiPropsEnd > uiRxLength THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiLoopGuard := 0;
WHILE uiPos < uiPropsEnd AND uiPos < uiRxLength AND uiLoopGuard < GVL_Mqtt.cnMaxPropertyLoop DO
byPropId := aRxBuf[uiPos];
uiPos := uiPos + 1;
CASE byPropId OF
GVL_Mqtt.cnPropReceiveMaximum:
IF uiPos + 1 >= uiPropsEnd THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiServerReceiveMax := SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]);
IF uiServerReceiveMax = 0 THEN
sDiagMsg := 'Receive Maximum must not be zero';
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiPos := uiPos + 2;
GVL_Mqtt.cnPropMaxQoS:
IF uiPos >= uiPropsEnd THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
byServerMaxQoS := aRxBuf[uiPos];
IF byServerMaxQoS > TO_BYTE(E_MqttQoS.byQoS2) THEN
sDiagMsg := 'Server Max QoS is invalid';
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiPos := uiPos + 1;
GVL_Mqtt.cnPropRetainAvailable:
IF uiPos >= uiPropsEnd THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
bServerRetainAvailable := BYTE_TO_BOOL(aRxBuf[uiPos]);
uiPos := uiPos + 1;
GVL_Mqtt.cnPropMaxPacketSize:
IF uiPos + 3 >= uiPropsEnd THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
udServerMaxPacketSize := SHL(BYTE_TO_UDINT(aRxBuf[uiPos]), 24) OR SHL(BYTE_TO_UDINT(aRxBuf[uiPos + 1]), 16) OR SHL(BYTE_TO_UDINT(aRxBuf[uiPos + 2]), 8) OR BYTE_TO_UDINT(aRxBuf[uiPos + 3]);
IF udServerMaxPacketSize = 0 THEN
sDiagMsg := 'Server Maximum Packet Size must not be zero';
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiPos := uiPos + 4;
GVL_Mqtt.cnPropTopicAliasMax:
IF uiPos + 1 >= uiPropsEnd THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiServerTopicAliasMax := SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]);
uiPos := uiPos + 2;
GVL_Mqtt.cnPropReasonString:
IF uiPos + 1 >= uiPropsEnd THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiPos := uiPos + 2 + (SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]));
GVL_Mqtt.cnPropServerKeepAlive:
IF uiPos + 1 >= uiPropsEnd THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiKeepAlive := SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]);
uiPos := uiPos + 2;
GVL_Mqtt.cnPropAssignedClientId:
IF uiPos + 1 >= uiPropsEnd THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiPos := uiPos + 2 + (SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]));
GVL_Mqtt.cnPropAuthMethod,
GVL_Mqtt.cnPropAuthData,
GVL_Mqtt.cnPropServerReference:
IF uiPos + 1 >= uiPropsEnd THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiPos := uiPos + 2 + (SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]));
GVL_Mqtt.cnPropResponseInfo:
IF uiPos + 1 >= uiPropsEnd THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiPos := uiPos + 2 + (SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]));
GVL_Mqtt.cnPropUserProperty:
IF uiPos + 1 >= uiPropsEnd THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiPos := uiPos + 2 + (SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]));
IF uiPos + 1 >= uiPropsEnd THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiPos := uiPos + 2 + (SHL(BYTE_TO_UINT(aRxBuf[uiPos]), 8) OR BYTE_TO_UINT(aRxBuf[uiPos + 1]));
GVL_Mqtt.cnPropWildcardSubAvail,
GVL_Mqtt.cnPropSubIdAvail,
GVL_Mqtt.cnPropSharedSubAvail:
IF uiPos >= uiPropsEnd THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
CASE byPropId OF
GVL_Mqtt.cnPropWildcardSubAvail:
bServerWildcardSubAvail := BYTE_TO_BOOL(aRxBuf[uiPos]);
GVL_Mqtt.cnPropSubIdAvail:
bServerSubIdAvail := BYTE_TO_BOOL(aRxBuf[uiPos]);
GVL_Mqtt.cnPropSharedSubAvail:
bServerSharedSubAvail := BYTE_TO_BOOL(aRxBuf[uiPos]);
END_CASE
uiPos := uiPos + 1;
GVL_Mqtt.cnPropRequestProblemInfo,
GVL_Mqtt.cnPropRequestResponseInfo:
IF uiPos >= uiPropsEnd THEN
M_HandleConnAck := FALSE;
RETURN;
END_IF
uiPos := uiPos + 1;
ELSE
M_HandleConnAck := FALSE;
RETURN;
END_CASE
uiLoopGuard := uiLoopGuard + 1;
END_WHILE
END_IF
uiRxLength := 0;
xWaitingForAck := FALSE;
IF eVersion = E_MqttVersion.byMqttVersion50 THEN
uiSendQuota := uiServerReceiveMax;
END_IF
IF NOT bCleanSession THEN
IF (byConnAckFlags AND GVL_Mqtt.cnConnAckSessionPresent) = 0 THEN
M_InflightClear();
M_SubListClear();
END_IF
ELSE
M_InflightClear();
M_SubListClear();
END_IF
M_TopicAliasClear();
M_HandleConnAck := TRUE;完整代码 3:M_BuildSubscribePacket
- 对应源码路径:
10 MQTT/MqttClient_V1_0/Device/Application/MQTT/POUs/MqttClient NBS/FB_MqttClient/处理发送报文/M_BuildSubscribePacket.st - 复制使用说明:如果你想把 5.0 的 Subscription Identifier、共享订阅这些能力用起来,订阅方法必须一起看。
- 阅读重点:重点看 5.0 订阅属性、共享订阅/通配符能力判定,以及客户端如何在组包前先检查 Broker 支持边界。
/// =======================================================================
/// 名称 : M_BuildSubscribePacket
/// 功能 : 构建 SUBSCRIBE 发送报文
/// 说明 : 根据 MQTT 版本组装订阅报文,并在发送前执行基础协议校验。
/// 编程人员 : ControlRookie
/// 时间 : 2026-05-05
/// 版本 : V1.2
/// =======================================================================
{attribute 'hide_all_locals'}
METHOD M_BuildSubscribePacket : BOOL
VAR
uiPos : UINT := 0; // 当前写入偏移
uiVarHeaderLen : UINT; // 可变报头长度
uiPayloadLen : UINT; // 载荷长度
uiRemainingLen : UINT; // 剩余长度
uiPropsLen : UINT; // MQTT 5.0 属性总长度
uiSubIdVbiBytes : UINT; // 订阅标识符 VBI 编码字节数
i : DINT; // 通用循环索引
bHasWildcard : BOOL; // 是否包含通配符
bIsSharedSubscription : BOOL; // 是否共享订阅
END_VAR
// === IMPLEMENTATION ===
IF SIZEOF(aTxBuf) < 256 THEN
M_BuildSubscribePacket := FALSE;
RETURN;
END_IF
IF sSubTopic = '' THEN
M_SetError(
uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
sMessage := 'Subscribe topic is required');
M_BuildSubscribePacket := FALSE;
RETURN;
END_IF
IF TO_UINT(LEN(sSubTopic)) > GVL_Mqtt.cnMaxTopicLen THEN
M_SetError(
uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
sMessage := 'Subscribe topic exceeds maximum length');
M_BuildSubscribePacket := FALSE;
RETURN;
END_IF
IF NOT M_IsValidUtf8String(sValue := sSubTopic) THEN
M_SetError(
uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
sMessage := 'Subscribe topic filter is not valid UTF-8');
M_BuildSubscribePacket := FALSE;
RETURN;
END_IF
IF NOT M_IsValidTopicFilter(sTopicFilter := sSubTopic) THEN
M_SetError(
uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
sMessage := 'Subscribe topic filter is invalid');
M_BuildSubscribePacket := FALSE;
RETURN;
END_IF
bHasWildcard := FALSE;
FOR i := 1 TO TO_DINT(LEN(sSubTopic)) DO
IF (sSubTopic[i - 1] = 16#2B) OR (sSubTopic[i - 1] = 16#23) THEN
bHasWildcard := TRUE;
EXIT;
END_IF
END_FOR
bIsSharedSubscription := FALSE;
IF LEN(sSubTopic) >= 7 THEN
IF (sSubTopic[0] = 16#24) AND
(sSubTopic[1] = 16#73) AND
(sSubTopic[2] = 16#68) AND
(sSubTopic[3] = 16#61) AND
(sSubTopic[4] = 16#72) AND
(sSubTopic[5] = 16#65) AND
(sSubTopic[6] = 16#2F) THEN
bIsSharedSubscription := TRUE;
END_IF
END_IF
IF (eVersion = E_MqttVersion.byMqttVersion50) AND (udiSubscriptionId > 0) AND (NOT bServerSubIdAvail) THEN
M_SetError(
uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
sMessage := 'Server does not support subscription identifiers');
M_BuildSubscribePacket := FALSE;
RETURN;
END_IF
IF (eVersion = E_MqttVersion.byMqttVersion50) AND (NOT bServerWildcardSubAvail) AND bHasWildcard THEN
M_SetError(
uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
sMessage := 'Server does not support wildcard subscriptions');
M_BuildSubscribePacket := FALSE;
RETURN;
END_IF
IF (eVersion = E_MqttVersion.byMqttVersion50) AND (NOT bServerSharedSubAvail) AND bIsSharedSubscription THEN
M_SetError(
uiErrorCode := TO_UINT(E_ReasonCode.uiErrInvalidParameter),
sMessage := 'Server does not support shared subscriptions');
M_BuildSubscribePacket := FALSE;
RETURN;
END_IF
FOR i := LOWER_BOUND(aTxBuf, 1) TO UPPER_BOUND(aTxBuf, 1) DO
aTxBuf[i] := 0;
END_FOR
uiVarHeaderLen := 2;
uiPropsLen := 0;
uiSubIdVbiBytes := 0;
IF eVersion = E_MqttVersion.byMqttVersion50 THEN
IF udiSubscriptionId > 0 THEN
IF udiSubscriptionId < 128 THEN
uiSubIdVbiBytes := 1;
ELSIF udiSubscriptionId < 16384 THEN
uiSubIdVbiBytes := 2;
ELSIF udiSubscriptionId < 2097152 THEN
uiSubIdVbiBytes := 3;
ELSE
uiSubIdVbiBytes := 4;
END_IF
uiPropsLen := uiPropsLen + 1 + uiSubIdVbiBytes;
END_IF
IF uiPropsLen < 128 THEN
uiVarHeaderLen := uiVarHeaderLen + 1 + uiPropsLen;
ELSE
uiVarHeaderLen := uiVarHeaderLen + 2 + uiPropsLen;
END_IF
END_IF
uiPayloadLen := TO_UINT(LEN(sSubTopic)) + 3;
uiRemainingLen := uiVarHeaderLen + uiPayloadLen;
IF uiRemainingLen + 5 > SIZEOF(aTxBuf) THEN
M_BuildSubscribePacket := FALSE;
RETURN;
END_IF
uiPos := 0;
aTxBuf[uiPos] := E_MqttPacketType.bySubscribe;
uiPos := uiPos + 1;
uiPos := uiPos + M_EncodeRemainingLength(udiLength := uiRemainingLen, pBuffer := ADR(aTxBuf[uiPos]));
uiPendingSubPacketId := M_GetNextPacketId();
IF uiPendingSubPacketId = 0 THEN
M_SetError(
uiErrorCode := TO_UINT(E_ReasonCode.uiErrReceiveMaxExceeded),
sMessage := 'No Packet ID available for subscribe');
M_BuildSubscribePacket := FALSE;
RETURN;
END_IF
uiExpectedPacketId := uiPendingSubPacketId;
aTxBuf[uiPos] := UINT_TO_BYTE(SHR(uiPendingSubPacketId, 8));
uiPos := uiPos + 1;
aTxBuf[uiPos] := UINT_TO_BYTE(uiPendingSubPacketId AND 16#FF);
uiPos := uiPos + 1;
IF eVersion = E_MqttVersion.byMqttVersion50 THEN
uiPos := uiPos + M_EncodeRemainingLength(udiLength := uiPropsLen, pBuffer := ADR(aTxBuf[uiPos]));
IF udiSubscriptionId > 0 THEN
aTxBuf[uiPos] := GVL_Mqtt.cnPropSubscriptionId;
uiPos := uiPos + 1;
uiPos := uiPos + M_EncodeRemainingLength(udiLength := udiSubscriptionId, pBuffer := ADR(aTxBuf[uiPos]));
END_IF
END_IF
uiPos := uiPos + M_AppendString(sStr := sSubTopic, pBuffer := ADR(aTxBuf[uiPos]));
aTxBuf[uiPos] := TO_BYTE(eSubQoS AND 16#03);
uiPos := uiPos + 1;
uiTxLength := uiPos;
M_BuildSubscribePacket := TRUE;系列导航
- 系列定位:加更篇 1
- 上一篇:第8篇 怎么把这套开源 MQTT 客户端真正用起来
- 下一篇:加更2 这套开源 MQTT 客户端,我是怎么一步一步做出来的
评论区预留
这里先保留评论和回复结构,不接入第三方服务。后续统一决定登录、匿名、审核、反垃圾和静态站兼容策略。