STM32F103 基于 MQTT 协议与 OneNET 平台通信
一、系统架构设计
┌─────────────────────────────────────────────────────────────┐
│ STM32F103 MQTT 物联网系统 │
├─────────────────────────────────────────────────────────────┤
│ STM32F103 │ ESP8266/ │ OneNET │ 应用端 │
│ (主控制器) │ Ethernet │ 云平台 │ (手机/网页) │
│ │ 模块 │ │ │
│ • 数据采集 │ • WiFi/ │ • 设备管理 │ • 远程监控 │
│ • 数据处理 │ TCP/IP │ • 数据存储 │ • 设备控制 │
│ • MQTT客户端 │ • MQTT桥接 │ • 规则引擎 │ • 数据分析 │
│ • 本地控制 │ │ • 告警通知 │ │
└─────────────────────────────────────────────────────────────┘
二、OneNET 平台配置
2.1 创建产品和设备
- 登录 OneNET 平台 (https://open.iot.10086.cn/)
- 创建产品:选择 MQTT 协议,填写产品信息
- 创建设备:获取设备 ID、产品 ID、鉴权信息
2.2 获取连接参数
// OneNET MQTT 连接参数
#define ONENET_SERVER "mqtt.heclouds.com"
#define ONENET_PORT 1883
#define PRODUCT_ID "1234567890" // 产品ID
#define DEVICE_ID "device_001" // 设备ID
#define AUTH_INFO "abcdef123456" // 鉴权信息(设备密钥)
#define TOPIC_DATA "$sys/" PRODUCT_ID "/" DEVICE_ID "/dp/post/json" // 数据上传主题
#define TOPIC_CMD "$sys/" PRODUCT_ID "/" DEVICE_ID "/cmd/#" // 命令下发主题
三、完整代码实现
3.1 MQTT 客户端核心代码
/**
* @file mqtt_client.c
* @brief STM32F103 OneNET MQTT 客户端
* @author AI Assistant
* @date 2024
*/
#include "stm32f10x.h"
#include "mqtt_client.h"
#include "esp8266.h"
#include "cJSON.h"
#include <string.h>
#include <stdio.h>
// MQTT 连接状态
typedef enum {
MQTT_DISCONNECTED = 0,
MQTT_CONNECTING,
MQTT_CONNECTED,
MQTT_SUBSCRIBED,
MQTT_PUBLISHED
} MQTT_State;
// MQTT 报文类型
typedef enum {
MQTT_CONNECT = 1,
MQTT_CONNACK = 2,
MQTT_PUBLISH = 3,
MQTT_PUBACK = 4,
MQTT_SUBSCRIBE = 8,
MQTT_SUBACK = 9,
MQTT_PINGREQ = 12,
MQTT_PINGRESP = 13,
MQTT_DISCONNECT = 14
} MQTT_PacketType;
// MQTT 客户端结构体
typedef struct {
char client_id[64];
char username[32];
char password[64];
uint16_t keep_alive;
uint8_t clean_session;
MQTT_State state;
uint16_t packet_id;
uint8_t connected;
} MQTT_Client;
// 全局变量
static MQTT_Client mqtt_client;
static uint8_t mqtt_buffer[512];
static uint16_t mqtt_buffer_len = 0;
/**
* @brief MQTT 客户端初始化
*/
void MQTT_Init(void) {
// 初始化 MQTT 客户端参数
sprintf(mqtt_client.client_id, "%s", DEVICE_ID);
sprintf(mqtt_client.username, "%s", PRODUCT_ID);
sprintf(mqtt_client.password, "%s", AUTH_INFO);
mqtt_client.keep_alive = 60;
mqtt_client.clean_session = 1;
mqtt_client.state = MQTT_DISCONNECTED;
mqtt_client.packet_id = 1;
mqtt_client.connected = 0;
printf("MQTT Client Initialized\r\n");
printf("Server: %s:%d\r\n", ONENET_SERVER, ONENET_PORT);
printf("Client ID: %s\r\n", mqtt_client.client_id);
printf("Username: %s\r\n", mqtt_client.username);
}
/**
* @brief 计算 MQTT 字符串长度字段
*/
static uint16_t MQTT_EncodeLength(uint8_t* buf, uint32_t len) {
uint16_t encoded = 0;
uint8_t digit;
do {
digit = len % 128;
len /= 128;
if (len > 0) {
digit |= 0x80;
}
buf[encoded++] = digit;
} while (len > 0);
return encoded;
}
/**
* @brief 构建 MQTT CONNECT 报文
*/
uint16_t MQTT_BuildConnectPacket(uint8_t* buffer, uint16_t buffer_size) {
uint16_t pos = 0;
uint8_t protocol_name[] = {0x00, 0x04, 'M', 'Q', 'T', 'T'};
uint8_t connect_flags = 0xC2; // 用户名+密码+清除会话
// 固定报头
buffer[pos++] = (MQTT_CONNECT << 4) | 0x00;
// 剩余长度
uint32_t remaining_length = 10 + 2 + strlen(mqtt_client.client_id) +
2 + strlen(mqtt_client.username) +
2 + strlen(mqtt_client.password);
pos += MQTT_EncodeLength(&buffer[pos], remaining_length);
// 可变报头
memcpy(&buffer[pos], protocol_name, 6);
pos += 6;
buffer[pos++] = 0x04; // 协议级别 MQTT 3.1.1
buffer[pos++] = connect_flags;
// 保持连接时间
buffer[pos++] = (mqtt_client.keep_alive >> 8) & 0xFF;
buffer[pos++] = mqtt_client.keep_alive & 0xFF;
// 客户端标识符
buffer[pos++] = (strlen(mqtt_client.client_id) >> 8) & 0xFF;
buffer[pos++] = strlen(mqtt_client.client_id) & 0xFF;
memcpy(&buffer[pos], mqtt_client.client_id, strlen(mqtt_client.client_id));
pos += strlen(mqtt_client.client_id);
// 用户名
buffer[pos++] = (strlen(mqtt_client.username) >> 8) & 0xFF;
buffer[pos++] = strlen(mqtt_client.username) & 0xFF;
memcpy(&buffer[pos], mqtt_client.username, strlen(mqtt_client.username));
pos += strlen(mqtt_client.username);
// 密码
buffer[pos++] = (strlen(mqtt_client.password) >> 8) & 0xFF;
buffer[pos++] = strlen(mqtt_client.password) & 0xFF;
memcpy(&buffer[pos], mqtt_client.password, strlen(mqtt_client.password));
pos += strlen(mqtt_client.password);
return pos;
}
/**
* @brief 构建 MQTT PUBLISH 报文
*/
uint16_t MQTT_BuildPublishPacket(uint8_t* buffer, uint16_t buffer_size,
const char* topic, const char* payload, uint8_t qos) {
uint16_t pos = 0;
uint16_t topic_len = strlen(topic);
uint16_t payload_len = strlen(payload);
// 固定报头
uint8_t fixed_header = (MQTT_PUBLISH << 4) | (qos << 1);
buffer[pos++] = fixed_header;
// 剩余长度
uint32_t remaining_length = 2 + topic_len + payload_len;
if (qos > 0) {
remaining_length += 2; // 报文标识符
}
pos += MQTT_EncodeLength(&buffer[pos], remaining_length);
// 可变报头 - 主题名
buffer[pos++] = (topic_len >> 8) & 0xFF;
buffer[pos++] = topic_len & 0xFF;
memcpy(&buffer[pos], topic, topic_len);
pos += topic_len;
// 报文标识符(QoS > 0 时需要)
if (qos > 0) {
buffer[pos++] = (mqtt_client.packet_id >> 8) & 0xFF;
buffer[pos++] = mqtt_client.packet_id & 0xFF;
mqtt_client.packet_id++;
}
// 载荷 - 消息体
memcpy(&buffer[pos], payload, payload_len);
pos += payload_len;
return pos;
}
/**
* @brief 构建 MQTT SUBSCRIBE 报文
*/
uint16_t MQTT_BuildSubscribePacket(uint8_t* buffer, uint16_t buffer_size, const char* topic, uint8_t qos) {
uint16_t pos = 0;
uint16_t topic_len = strlen(topic);
// 固定报头
buffer[pos++] = (MQTT_SUBSCRIBE << 4) | 0x02;
// 剩余长度
uint32_t remaining_length = 2 + 2 + topic_len + 1;
pos += MQTT_EncodeLength(&buffer[pos], remaining_length);
// 可变报头 - 报文标识符
buffer[pos++] = (mqtt_client.packet_id >> 8) & 0xFF;
buffer[pos++] = mqtt_client.packet_id & 0xFF;
mqtt_client.packet_id++;
// 载荷 - 主题过滤器
buffer[pos++] = (topic_len >> 8) & 0xFF;
buffer[pos++] = topic_len & 0xFF;
memcpy(&buffer[pos], topic, topic_len);
pos += topic_len;
// QoS 等级
buffer[pos++] = qos;
return pos;
}
/**
* @brief 构建 MQTT PINGREQ 报文
*/
uint16_t MQTT_BuildPingReqPacket(uint8_t* buffer, uint16_t buffer_size) {
uint16_t pos = 0;
// 固定报头
buffer[pos++] = (MQTT_PINGREQ << 4) | 0x00;
buffer[pos++] = 0x00; // 剩余长度 0
return pos;
}
/**
* @brief 解析 MQTT CONNACK 报文
*/
uint8_t MQTT_ParseConnAck(uint8_t* buffer, uint16_t length) {
if (length < 4) return 0;
uint8_t connack_flags = buffer[2];
uint8_t connack_rc = buffer[3];
if (connack_rc == 0x00) {
printf("MQTT Connected Successfully!\r\n");
mqtt_client.connected = 1;
mqtt_client.state = MQTT_CONNECTED;
return 1;
} else {
printf("MQTT Connection Failed! RC=%d\r\n", connack_rc);
return 0;
}
}
/**
* @brief 解析 MQTT SUBACK 报文
*/
uint8_t MQTT_ParseSubAck(uint8_t* buffer, uint16_t length) {
if (length < 5) return 0;
uint8_t suback_rc = buffer[4];
if (suback_rc == 0x00 || suback_rc == 0x01 || suback_rc == 0x02) {
printf("MQTT Subscribed Successfully!\r\n");
mqtt_client.state = MQTT_SUBSCRIBED;
return 1;
} else {
printf("MQTT Subscribe Failed! RC=%d\r\n", suback_rc);
return 0;
}
}
/**
* @brief 解析 MQTT PUBLISH 报文(命令下发)
*/
void MQTT_ParsePublish(uint8_t* buffer, uint16_t length) {
uint16_t pos = 0;
uint16_t topic_len, payload_len;
char topic[128] = {0};
char payload[256] = {0};
// 跳过固定报头
pos += 1;
uint8_t remaining_len = buffer[pos++];
// 解析主题名
topic_len = (buffer[pos] << 8) | buffer[pos+1];
pos += 2;
memcpy(topic, &buffer[pos], topic_len);
pos += topic_len;
// 跳过报文标识符(如果存在)
if ((buffer[0] >> 1) & 0x03) {
pos += 2;
}
// 解析载荷
payload_len = remaining_len - topic_len - (pos - 2);
memcpy(payload, &buffer[pos], payload_len);
printf("Received Command:\r\n");
printf("Topic: %s\r\n", topic);
printf("Payload: %s\r\n", payload);
// 处理命令
MQTT_HandleCommand(topic, payload);
}
/**
* @brief 处理下发的命令
*/
void MQTT_HandleCommand(const char* topic, const char* payload) {
cJSON *root = cJSON_Parse(payload);
if (root == NULL) {
printf("Invalid JSON Format!\r\n");
return;
}
cJSON *cmd = cJSON_GetObjectItem(root, "cmd");
if (cmd != NULL) {
if (strcmp(cmd->valuestring, "led_on") == 0) {
// 打开 LED
GPIO_SetBits(GPIOC, GPIO_Pin_13);
printf("LED Turned ON\r\n");
} else if (strcmp(cmd->valuestring, "led_off") == 0) {
// 关闭 LED
GPIO_ResetBits(GPIOC, GPIO_Pin_13);
printf("LED Turned OFF\r\n");
}
}
cJSON_Delete(root);
}
3.2 网络通信模块(ESP8266 AT 指令)
/**
* @file esp8266.c
* @brief ESP8266 WiFi 模块驱动
*/
#include "esp8266.h"
#include "usart.h"
#include <string.h>
#define ESP8266_USART USART1
// ESP8266 状态
typedef enum {
ESP8266_IDLE = 0,
ESP8266_READY,
ESP8266_CONNECTED,
ESP8266_DISCONNECTED
} ESP8266_State;
static ESP8266_State esp_state = ESP8266_IDLE;
static uint8_t esp_rx_buffer[1024];
static uint16_t esp_rx_index = 0;
/**
* @brief 发送 AT 指令
*/
uint8_t ESP8266_SendATCommand(const char* cmd, const char* expect, uint32_t timeout) {
char response[256] = {0};
uint32_t start_time = HAL_GetTick();
USART_SendString(ESP8266_USART, (char*)cmd);
USART_SendString(ESP8266_USART, "\r\n");
while ((HAL_GetTick() - start_time) < timeout) {
if (ESP8266_ReceiveResponse(response, 256, 100)) {
if (strstr(response, expect) != NULL) {
return 1;
}
}
}
return 0;
}
/**
* @brief 初始化 ESP8266
*/
uint8_t ESP8266_Init(void) {
printf("Initializing ESP8266...\r\n");
// 测试 AT 指令
if (!ESP8266_SendATCommand("AT", "OK", 2000)) {
printf("ESP8266 Not Responding!\r\n");
return 0;
}
// 设置 WiFi 模式为 Station
if (!ESP8266_SendATCommand("AT+CWMODE=1", "OK", 2000)) {
printf("Set WiFi Mode Failed!\r\n");
return 0;
}
// 连接 WiFi
char cmd[128];
sprintf(cmd, "AT+CWJAP=\"%s\",\"%s\"", WIFI_SSID, WIFI_PASSWORD);
if (!ESP8266_SendATCommand(cmd, "WIFI GOT IP", 10000)) {
printf("WiFi Connection Failed!\r\n");
return 0;
}
// 启用多连接
if (!ESP8266_SendATCommand("AT+CIPMUX=1", "OK", 2000)) {
printf("Enable Multi-Connection Failed!\r\n");
return 0;
}
esp_state = ESP8266_READY;
printf("ESP8266 Initialized Successfully!\r\n");
return 1;
}
/**
* @brief 建立 TCP 连接
*/
uint8_t ESP8266_ConnectTCP(const char* server, uint16_t port) {
char cmd[128];
sprintf(cmd, "AT+CIPSTART=0,\"TCP\",\"%s\",%d", server, port);
if (!ESP8266_SendATCommand(cmd, "CONNECT", 5000)) {
printf("TCP Connection Failed!\r\n");
return 0;
}
esp_state = ESP8266_CONNECTED;
printf("TCP Connected to %s:%d\r\n", server, port);
return 1;
}
/**
* @brief 发送 TCP 数据
*/
uint8_t ESP8266_SendTCPData(uint8_t* data, uint16_t length) {
char cmd[32];
sprintf(cmd, "AT+CIPSEND=0,%d", length);
if (!ESP8266_SendATCommand(cmd, ">", 2000)) {
return 0;
}
USART_SendData(ESP8266_USART, data, length);
if (ESP8266_SendATCommand("", "SEND OK", 3000)) {
return 1;
}
return 0;
}
/**
* @brief 接收 TCP 数据
*/
uint8_t ESP8266_ReceiveTCPData(uint8_t* buffer, uint16_t buffer_size, uint32_t timeout) {
uint32_t start_time = HAL_GetTick();
while ((HAL_GetTick() - start_time) < timeout) {
if (USART_ReceiveData(ESP8266_USART, buffer, buffer_size, 100)) {
return 1;
}
}
return 0;
}
3.3 主程序与应用逻辑
/**
* @file main.c
* @brief STM32F103 OneNET MQTT 主程序
*/
#include "stm32f10x.h"
#include "mqtt_client.h"
#include "esp8266.h"
#include "sensor.h"
#include "cJSON.h"
// 系统状态
typedef enum {
SYS_INIT = 0,
SYS_WIFI_CONNECTING,
SYS_MQTT_CONNECTING,
SYS_RUNNING,
SYS_ERROR
} SystemState;
static SystemState sys_state = SYS_INIT;
static uint32_t last_publish_time = 0;
static uint32_t last_ping_time = 0;
static uint8_t mqtt_connected = 0;
/**
* @brief 构建传感器数据 JSON
*/
void BuildSensorDataJSON(char* json_buffer, uint16_t buffer_size) {
cJSON *root = cJSON_CreateObject();
cJSON *datastreams = cJSON_CreateArray();
cJSON *datastream = cJSON_CreateObject();
// 获取传感器数据
float temperature = Sensor_GetTemperature();
float humidity = Sensor_GetHumidity();
float voltage = Sensor_GetVoltage();
// 构建数据流
cJSON_AddStringToObject(datastream, "id", "temperature");
cJSON_AddNumberToObject(datastream, "value", temperature);
cJSON_AddItemToArray(datastreams, datastream);
datastream = cJSON_CreateObject();
cJSON_AddStringToObject(datastream, "id", "humidity");
cJSON_AddNumberToObject(datastream, "value", humidity);
cJSON_AddItemToArray(datastreams, datastream);
datastream = cJSON_CreateObject();
cJSON_AddStringToObject(datastream, "id", "voltage");
cJSON_AddNumberToObject(datastream, "value", voltage);
cJSON_AddItemToArray(datastreams, datastream);
cJSON_AddItemToObject(root, "datastreams", datastreams);
char* json_str = cJSON_PrintUnformatted(root);
strncpy(json_buffer, json_str, buffer_size);
cJSON_free(json_str);
cJSON_Delete(root);
}
/**
* @brief 发布传感器数据到 OneNET
*/
void PublishSensorData(void) {
char json_buffer[512];
BuildSensorDataJSON(json_buffer, sizeof(json_buffer));
uint16_t packet_len = MQTT_BuildPublishPacket(mqtt_buffer, sizeof(mqtt_buffer),
TOPIC_DATA, json_buffer, 1);
if (ESP8266_SendTCPData(mqtt_buffer, packet_len)) {
printf("Data Published to OneNET: %s\r\n", json_buffer);
mqtt_client.state = MQTT_PUBLISHED;
}
}
/**
* @brief MQTT 连接流程
*/
void MQTT_ConnectionProcess(void) {
switch (mqtt_client.state) {
case MQTT_DISCONNECTED:
// 构建并发送 CONNECT 报文
mqtt_buffer_len = MQTT_BuildConnectPacket(mqtt_buffer, sizeof(mqtt_buffer));
ESP8266_SendTCPData(mqtt_buffer, mqtt_buffer_len);
mqtt_client.state = MQTT_CONNECTING;
break;
case MQTT_CONNECTING:
// 等待 CONNACK
if (ESP8266_ReceiveTCPData(mqtt_buffer, sizeof(mqtt_buffer), 1000)) {
if (MQTT_ParseConnAck(mqtt_buffer, mqtt_buffer_len)) {
// 连接成功,订阅命令主题
mqtt_buffer_len = MQTT_BuildSubscribePacket(mqtt_buffer, sizeof(mqtt_buffer),
TOPIC_CMD, 1);
ESP8266_SendTCPData(mqtt_buffer, mqtt_buffer_len);
}
}
break;
case MQTT_CONNECTED:
// 订阅成功,开始正常工作
mqtt_connected = 1;
sys_state = SYS_RUNNING;
break;
case MQTT_SUBSCRIBED:
// 订阅成功
mqtt_client.state = MQTT_CONNECTED;
break;
default:
break;
}
}
/**
* @brief 系统主循环
*/
int main(void) {
// 系统初始化
System_Init();
USART_Init(115200);
GPIO_Init();
Sensor_Init();
Delay_Init();
printf("STM32F103 OneNET MQTT System Starting...\r\n");
// MQTT 初始化
MQTT_Init();
while (1) {
switch (sys_state) {
case SYS_INIT:
// 初始化 ESP8266
if (ESP8266_Init()) {
sys_state = SYS_WIFI_CONNECTING;
} else {
sys_state = SYS_ERROR;
}
break;
case SYS_WIFI_CONNECTING:
// 连接 OneNET MQTT 服务器
if (ESP8266_ConnectTCP(ONENET_SERVER, ONENET_PORT)) {
sys_state = SYS_MQTT_CONNECTING;
}
break;
case SYS_MQTT_CONNECTING:
// MQTT 连接流程
MQTT_ConnectionProcess();
break;
case SYS_RUNNING:
// 定时发布传感器数据(每10秒)
if ((HAL_GetTick() - last_publish_time) >= 10000) {
PublishSensorData();
last_publish_time = HAL_GetTick();
}
// 定时发送心跳包(每30秒)
if ((HAL_GetTick() - last_ping_time) >= 30000) {
mqtt_buffer_len = MQTT_BuildPingReqPacket(mqtt_buffer, sizeof(mqtt_buffer));
ESP8266_SendTCPData(mqtt_buffer, mqtt_buffer_len);
last_ping_time = HAL_GetTick();
}
// 检查是否有命令下发
if (ESP8266_ReceiveTCPData(mqtt_buffer, sizeof(mqtt_buffer), 100)) {
MQTT_ParsePublish(mqtt_buffer, mqtt_buffer_len);
}
break;
case SYS_ERROR:
printf("System Error! Restarting...\r\n");
Delay_ms(5000);
sys_state = SYS_INIT;
break;
}
Delay_ms(10);
}
}
四、OneNET 平台数据流配置
4.1 数据流模板
在 OneNET 平台创建设备后,需要配置数据流模板:
{
"datastreams": [
{
"id": "temperature",
"name": "温度",
"unit": "℃",
"data_type": "number"
},
{
"id": "humidity",
"name": "湿度",
"unit": "%",
"data_type": "number"
},
{
"id": "voltage",
"name": "电压",
"unit": "V",
"data_type": "number"
}
]
}
4.2 命令下发格式
通过 OneNET 平台下发控制命令:
{
"cmd": "led_on",
"params": {
"duration": 5000
}
}
参考代码 stm32f103单片机基于mqtt协议和onenet平台通信 www.youwenfan.com/contentcnu/60344.html
五、调试与优化
5.1 调试技巧
- 串口调试:使用串口打印 MQTT 报文,验证协议正确性
- 网络抓包:使用 Wireshark 抓包分析 MQTT 通信过程
- OneNET 日志:查看平台端的设备日志,确认数据接收情况
5.2 优化建议
- 断线重连:实现 MQTT 断线自动重连机制
- 数据缓存:网络异常时缓存数据,恢复后补发
- 功耗优化:休眠模式下关闭 WiFi,定时唤醒上传数据
- 安全加固:使用 TLS/SSL 加密连接,保护数据安全
5.3 常见问题解决
| 问题 | 原因 | 解决方案 |
|---|---|---|
| 连接超时 | WiFi 信号弱 | 检查 WiFi 密码,增强信号 |
| MQTT 认证失败 | 设备密钥错误 | 核对 OneNET 平台设备信息 |
| 数据上传失败 | JSON 格式错误 | 验证 JSON 语法,使用 cJSON 库 |
| 命令无响应 | 订阅主题错误 | 检查订阅主题格式是否正确 |