从零吃透 MQTT 通信|第 2 章:裸机环境手动构建 MQTT 协议栈,报文组包与解析实现
专栏说明:本章不依赖 paho‑mqtt 等第三方库,完全基于 C 语言实现 MQTT3.1.1 报文封装、报文解析。面向 STM32 裸机、无操作系统 MCU,聚焦协议底层,帮助读者理解每一字节报文的构成,为后续 RTOS 工程、设备对接云平台打下底层基础。
前言
绝大多数嵌入式开发者使用 MQTT 时,直接移植成熟第三方 MQTT 库,调用 Connect、Publish、Subscribe 接口即可完成业务开发。但当设备资源受限、库裁剪困难、出现诡异报文异常时,如果不理解报文内部字节构成,调试将会非常艰难。
本章我们不调用任何 MQTT 库,手动实现核心报文的组装函数与简易解析逻辑。代码全部采用标准 C 编写,可以直接移植到 STM32、GD32 等单片机裸机工程。
前置知识:需要阅读第 1 章 MQTT 报文三段结构:固定头部、可变头部、有效载荷 Payload。MQTT 版本采用业界最通用 3.1.1。
1、开发前的基础数据结构定义
首先定义协议相关枚举、结构体。固定头部第一个字节包含报文类型、DUP 重传标志、QoS 等级、RETAIN 保留标志。
#ifndef MQTT_RAW_H
#define MQTT_RAW_H
#include <stdint.h>
#include <string.h>
/* MQTT 3.1.1 报文类型枚举 */
typedef enum
{
MQTT_MSG_CONNECT = 1,
MQTT_MSG_CONNACK = 2,
MQTT_MSG_PUBLISH = 3,
MQTT_MSG_PUBACK = 4,
MQTT_MSG_SUBSCRIBE = 8,
MQTT_MSG_SUBACK = 9,
MQTT_MSG_UNSUBSCRIBE = 10,
MQTT_MSG_UNSUBACK = 11,
MQTT_MSG_PINGREQ = 12,
MQTT_MSG_PINGRESP = 13,
MQTT_MSG_DISCONNECT = 14
}MqttMsgType_e;
/* QoS等级 */
typedef enum
{
MQTT_QOS0 = 0,
MQTT_QOS1 = 1,
MQTT_QOS2 = 2
}MqttQos_e;
/* MQTT连接参数结构体,统一管理CONNECT报文入参 */
typedef struct
{
const char* clientId;
const char* username;
const char* password;
uint16_t keepAlive;
uint8_t cleanSession;
uint8_t willFlag;
MqttQos_e willQos;
uint8_t willRetain;
const char* willTopic;
const uint8_t* willPayload;
uint16_t willPayloadLen;
}MqttConnectParam_t;
#endif
MQTT 报文里面有一个特殊字段:剩余长度 (Remaining Length)。这个字段不是固定 1/2/4 字节,采用可变字节编码。这也是手写协议栈第一个需要攻克的难点。
剩余长度 = 可变头部字节数 + Payload 字节数,不包含固定头部本身。编码规则:每字节最高 bit 代表是否后续还有字节,低 7 位存储有效数据,最大支持 256MB 报文。
1.1 剩余长度编码与解码工具函数
/**
* @brief 编码MQTT剩余长度字段
* @param buf 输出缓冲区
* @param length 需要编码的数值
* @return 编码占用字节数(1~4)
*/
uint8_t mqtt_encode_remaining_len(uint8_t* buf, uint32_t length)
{
uint8_t index = 0;
do
{
uint8_t byte = length % 128;
length = length / 128;
if(length > 0)
{
byte |= 0x80;
}
buf[index++] = byte;
}while(length > 0 && index < 4);
return index;
}
/**
* @brief 解码剩余长度
* @param buf 输入报文缓冲区
* @param outLen 解析得到的剩余长度
* @return 使用字节数量,解析失败返回0
*/
uint8_t mqtt_decode_remaining_len(const uint8_t* buf, uint32_t* outLen)
{
uint32_t multiplier = 1;
uint32_t value = 0;
uint8_t bytesRead = 0;
uint8_t byte;
do
{
if(bytesRead >= 4)
{
return 0;
}
byte = buf[bytesRead++];
value += (byte & 0x7F) * multiplier;
multiplier *= 128;
}while((byte & 0x80) != 0);
*outLen = value;
return bytesRead;
}
工程踩坑点:很多新手直接用 4 字节固定存储剩余长度,会导致 Broker 拒绝报文。MQTT 协议强制要求可变编码。
2、核心报文组包实现
所有组包函数传入输出缓存,返回生成报文的总字节长度;缓存不足返回 0,调用方需要判断缓冲区大小,防止内存越界。
2.1 CONNECT 连接报文组包
CONNECT 是设备建立 MQTT 会话的报文,包含客户端 ID、账号密码、心跳时间、遗嘱配置。
/**
* @brief 组装CONNECT报文
* @param buf 输出缓存
* @param bufLen 缓存总大小
* @param param 连接参数结构体
* @return 生成报文长度,0代表缓存不足
*/
uint16_t mqtt_build_connect(uint8_t* buf, uint16_t bufLen, const MqttConnectParam_t* param)
{
uint16_t offset = 0;
uint16_t payloadLen = 0;
uint16_t variableHeadLen;
uint8_t remainLenBuf[4];
uint8_t remainLenBytes;
/* 协议名 + 协议等级 + 连接标志 + keepalive,可变头部固定10字节 */
variableHeadLen = 10;
/* 计算Payload总长度 clientId + username + password + will */
payloadLen += 2 + strlen(param->clientId);
if(param->username != NULL)
{
payloadLen += 2 + strlen(param->username);
}
if(param->password != NULL)
{
payloadLen += 2 + strlen(param->password);
}
if(param->willFlag != 0)
{
payloadLen += 2 + strlen(param->willTopic);
payloadLen += 2 + param->willPayloadLen;
}
uint32_t remainingLen = variableHeadLen + payloadLen;
remainLenBytes = mqtt_encode_remaining_len(remainLenBuf, remainingLen);
uint16_t totalPacketLen = 1 + remainLenBytes + remainingLen;
if(totalPacketLen > bufLen)
{
return 0;
}
/* 固定头部第一字节:报文类型CONNECT=1,标志位全部0 */
buf[offset++] = (MQTT_MSG_CONNECT << 4);
/* 写入编码后的剩余长度 */
memcpy(&buf[offset], remainLenBuf, remainLenBytes);
offset += remainLenBytes;
/* ----可变头部开始---- */
/* 协议名称:MQTT,字符串前缀2字节长度 */
buf[offset++] = 0x00;
buf[offset++] = 0x04;
memcpy(&buf[offset], "MQTT", 4);
offset += 4;
/* MQTT 3.1.1协议等级 0x04 */
buf[offset++] = 0x04;
/* 连接标志位组装 */
uint8_t connectFlag = 0;
if(param->cleanSession) connectFlag |= (1 << 1);
if(param->willFlag)
{
connectFlag |= (1 << 2);
connectFlag |= (param->willQos << 3);
if(param->willRetain) connectFlag |= (1 << 5);
}
if(param->username != NULL) connectFlag |= (1 << 6);
if(param->password != NULL) connectFlag |= (1 << 7);
buf[offset++] = connectFlag;
/* keep‑alive 大端2字节 */
buf[offset++] = (param->keepAlive >> 8) & 0xFF;
buf[offset++] = param->keepAlive & 0xFF;
/* ----Payload开始---- */
/* clientId */
uint16_t strLen = strlen(param->clientId);
buf[offset++] = (strLen >> 8) & 0xFF;
buf[offset++] = strLen & 0xFF;
memcpy(&buf[offset], param->clientId, strLen);
offset += strLen;
/* will消息 */
if(param->willFlag)
{
strLen = strlen(param->willTopic);
buf[offset++] = (strLen >> 8) & 0xFF;
buf[offset++] = strLen & 0xFF;
memcpy(&buf[offset], param->willTopic, strLen);
offset += strLen;
buf[offset++] = (param->willPayloadLen >> 8) & 0xFF;
buf[offset++] = param->willPayloadLen & 0xFF;
memcpy(&buf[offset], param->willPayload, param->willPayloadLen);
offset += param->willPayloadLen;
}
/* username */
if(param->username != NULL)
{
strLen = strlen(param->username);
buf[offset++] = (strLen >> 8) & 0xFF;
buf[offset++] = strLen & 0xFF;
memcpy(&buf[offset], param->username, strLen);
offset += strLen;
}
/* password */
if(param->password != NULL)
{
strLen = strlen(param->password);
buf[offset++] = (strLen >> 8) & 0xFF;
buf[offset++] = strLen & 0xFF;
memcpy(&buf[offset], param->password, strLen);
offset += strLen;
}
return offset;
}
2.2 PUBLISH 发布报文组包
/**
* @brief 组装PUBLISH报文
* @param buf 输出缓存
* @param bufLen 缓存大小
* @param topic 主题字符串
* @param qos QoS等级
* @param retain 保留标志 0/1
* @param dup 重传标志 0/1
* @param packetId QoS>0有效,报文标识符
* @param payload 业务数据
* @param payloadLen 业务数据长度
* @return 报文总长度,0缓存不足
*/
uint16_t mqtt_build_publish(uint8_t* buf, uint16_t bufLen, const char* topic,
MqttQos_e qos, uint8_t retain, uint8_t dup,
uint16_t packetId, const uint8_t* payload, uint16_t payloadLen)
{
uint16_t offset = 0;
uint8_t remainBuf[4];
uint8_t remainBytes;
uint16_t topicLen = strlen(topic);
uint16_t variableLen = 2 + topicLen;
if(qos != MQTT_QOS0)
{
variableLen += 2; /* QoS1/QoS2需要报文ID */
}
uint32_t remainLen = variableLen + payloadLen;
remainBytes = mqtt_encode_remaining_len(remainBuf, remainLen);
uint16_t totalLen = 1 + remainBytes + remainLen;
if(totalLen > bufLen) return 0;
/* 固定头部首字节 */
uint8_t firstByte = (MQTT_MSG_PUBLISH << 4);
if(dup) firstByte |= (1 << 3);
firstByte |= (qos << 1);
if(retain) firstByte |= 1;
buf[offset++] = firstByte;
memcpy(&buf[offset], remainBuf, remainBytes);
offset += remainBytes;
/* 可变头部 topic长度+topic */
buf[offset++] = (topicLen >> 8) & 0xFF;
buf[offset++] = topicLen & 0xFF;
memcpy(&buf[offset], topic, topicLen);
offset += topicLen;
if(qos != MQTT_QOS0)
{
buf[offset++] = (packetId >> 8) & 0xFF;
buf[offset++] = packetId & 0xFF;
}
/* payload */
if(payloadLen > 0 && payload != NULL)
{
memcpy(&buf[offset], payload, payloadLen);
offset += payloadLen;
}
return offset;
}
2.3 SUBSCRIBE 订阅报文、PINGREQ 心跳、DISCONNECT 报文
/**
* @brief 组装心跳PINGREQ
*/
uint16_t mqtt_build_pingreq(uint8_t* buf, uint16_t bufLen)
{
if(bufLen < 2) return 0;
buf[0] = (MQTT_MSG_PINGREQ << 4);
buf[1] = 0x00;
return 2;
}
/**
* @brief 组装断开连接DISCONNECT
*/
uint16_t mqtt_build_disconnect(uint8_t* buf, uint16_t bufLen)
{
if(bufLen < 2) return 0;
buf[0] = (MQTT_MSG_DISCONNECT << 4);
buf[1] = 0x00;
return 2;
}
/**
* @brief SUBSCRIBE订阅报文
* @param topic 订阅主题
* @param reqQos 请求QoS
* @param packetId 报文ID
*/
uint16_t mqtt_build_subscribe(uint8_t* buf,uint16_t bufLen,uint16_t packetId,const char* topic,MqttQos_e reqQos)
{
uint16_t offset = 0;
uint8_t remainBuf[4];
uint8_t remainBytes;
uint16_t topicLen = strlen(topic);
uint32_t remainLen = 2 + 2 + topicLen + 1;
remainBytes = mqtt_encode_remaining_len(remainBuf, remainLen);
uint16_t totalLen = 1 + remainBytes + remainLen;
if(totalLen > bufLen) return 0;
buf[offset++] = (MQTT_MSG_SUBSCRIBE << 4) | 0x02;
memcpy(&buf[offset], remainBuf, remainBytes);
offset += remainBytes;
/* 报文ID */
buf[offset++] = (packetId >> 8) & 0xFF;
buf[offset++] = packetId & 0xFF;
/* 主题长度+主题 */
buf[offset++] = (topicLen >> 8) & 0xFF;
buf[offset++] = topicLen & 0xFF;
memcpy(&buf[offset], topic, topicLen);
offset += topicLen;
/* 订阅请求QoS */
buf[offset++] = reqQos;
return offset;
}
3、简易报文解析框架(接收方向)
MCU 从 TCP 缓冲区收到字节流,会遇到粘包、分包问题,原始数据流不能直接解析 MQTT 报文。工程上必须使用环形缓冲区缓存 TCP 字节,做报文边界识别。
这里给出解析逻辑伪代码框架,不实现完整环形缓存,重点说明解析流程。
typedef struct
{
MqttMsgType_e msgType;
uint8_t dup;
MqttQos_e qos;
uint8_t retain;
uint32_t remainLen;
const uint8_t* payloadPtr;
uint16_t payloadLen;
}MqttPacket_t;
/**
* @brief 解析一个完整MQTT报文
* @param rawPacket 指向一个完整MQTT报文首地址
* @param outPkt 解析输出结构体
* @return 0成功,非0失败
*/
int mqtt_parse_packet(const uint8_t* rawPacket, MqttPacket_t* outPkt)
{
uint8_t firstByte = rawPacket[0];
outPkt->msgType = (firstByte >> 4) & 0x0F;
outPkt->dup = (firstByte >> 3) & 0x01;
outPkt->qos = (firstByte >> 1) & 0x03;
outPkt->retain = firstByte & 0x01;
uint32_t remainLen;
uint8_t consumeBytes = mqtt_decode_remaining_len(&rawPacket[1], &remainLen);
if(consumeBytes == 0) return -1;
outPkt->remainLen = remainLen;
uint16_t payloadOffset = 1 + consumeBytes + (rawPacket[1+consumeBytes]?0:0);
/* 注意:不同报文可变头部长度不一样,需要根据msgType再偏移得到payload位置 */
return 0;
}
重要提示:解析函数只处理已经切分完整的单条报文,不能直接喂 TCP 流式数据。真实项目中需要环形缓冲区 + 报文边界检测,把 TCP 字节流切割成一条条完整 MQTT 报文之后,再调用解析函数。
4、裸机调用示例(伪业务逻辑)
假设已经实现 TCP 底层:
tcp_send()发送字节;tcp_recv()接收网络数据。
uint8_t mqttTxBuf[512];
MqttConnectParam_t connParam = {0};
void mqtt_demo_connect(void)
{
connParam.clientId = "mcu_device_001";
connParam.username = "user123";
connParam.password = "pass456";
connParam.keepAlive = 60;
connParam.cleanSession = 1;
connParam.willFlag = 0;
uint16_t pktLen = mqtt_build_connect(mqttTxBuf, sizeof(mqttTxBuf), &connParam);
if(pktLen > 0)
{
tcp_send(mqttTxBuf, pktLen);
}
}
/* 发布一条QoS0消息 */
void mqtt_demo_publish(void)
{
const char* topic = "device/mcu/status";
uint8_t data[] = "{\"temp\":25.1}";
uint16_t len = mqtt_build_publish(mqttTxBuf,sizeof(mqttTxBuf),
topic,MQTT_QOS0,0,0,0,data,sizeof(data)-1);
if(len > 0)
{
tcp_send(mqttTxBuf, len);
}
}
5、工程开发遇到的典型问题汇总
- 剩余长度编码错误:很多新手固定写 2 字节剩余长度,连接直接被 Broker 拒绝。务必使用上面的变长编码函数。
- 大端字节序:MQTT 协议所有 16 位数值(报文 ID、字符串长度)全部使用网络大端,单片机小端模式必须手动移位转换,不能直接
memcpy赋值 uint16_t。 - TCP 粘包分包:TCP 是流式协议,一次 recv 可能多条报文,也可能只收到半条报文。裸机开发必须环形缓冲区做报文重组,这是 MQTT 裸机开发最大坑点。
- SUBSCRIBE 报文固定标志位:订阅报文固定头部低 4 位必须等于
0x02,写错会直接返回协议错误。 - QoS0 报文不能携带 PacketId:QoS0 的 PUBLISH 报文,协议不允许报文标识符,如果强行写入,部分 Broker 会断开连接。
本章总结
本章脱离第三方 MQTT 库,手动完成 MQTT3.1.1 核心报文的封装基础代码。掌握报文组装逻辑之后,后续再阅读 paho‑mqtt 源码,理解门槛会大幅降低。
当前代码只完成报文编解码层,还缺少:环形缓冲区报文切分、状态机管理、超时重传、断线重连、应答处理。裸机完整工程不能只依靠组包函数,需要上层状态机驱动。
下一章预告:第 3 章:裸机 MQTT 状态机设计,环形缓冲区实现,完整收发工程框架
💡如果文章对你有帮助,欢迎点赞收藏,专栏持续更新嵌入式 MQTT 实战内容。
openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构
更多推荐



所有评论(0)