ControlRookie
返回文章

加更1_MQTT 5.0 到底值不值得上?工业现场别急着无脑升级

这一篇专门回答一个非常现实的问题:MQTT 5.0 到底是不是该直接上。重点会从能力协商、Reason Code、边界控制、Broker 支持度和实现成熟度几个维度,把支持 5.0 和跑稳 5.0 的差别说透。

这一篇专门回答一个非常现实的问题: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 的,不止协议本身。

还包括:

  1. Broker 支持度
  2. 第三方客户端支持度
  3. 你自己的实现成熟度
  4. 现场到底有没有这些需求

举几个很现实的例子:

场景 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 AliasM_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 句话

  1. MQTT 5.0 值得学,也值得用,但不该无脑上。
  2. 真正该问的不是“支不支持 5.0”,而是“支持到哪、稳不稳、用不用得上”。
  3. 5.0 最有价值的地方,是能力协商更清楚、边界更明确。
  4. 版本号能切到 5.0,不等于你的客户端已经稳定支持 5.0。
  5. 基础 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 属性是怎么真正落进字节流的。
iecst
/// =======================================================================
/// 名称      : 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 这些运行变量里。
iecst
/// =======================================================================
/// 名称      : 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 支持边界。
iecst
/// =======================================================================
/// 名称      : 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 客户端,我是怎么一步一步做出来的
评论和回复区

评论区预留

这里先保留评论和回复结构,不接入第三方服务。后续统一决定登录、匿名、审核、反垃圾和静态站兼容策略。

↑ ↓