STM32F103 基于 MQTT 协议与 OneNET 平台通信

STM32F103 基于 MQTT 协议与 OneNET 平台通信

一、系统架构设计

┌─────────────────────────────────────────────────────────────┐
│                    STM32F103 MQTT 物联网系统                 │
├─────────────────────────────────────────────────────────────┤
│  STM32F103     │  ESP8266/   │  OneNET      │  应用端      │
│  (主控制器)    │  Ethernet   │  云平台      │  (手机/网页) │
│                │  模块       │              │              │
│  • 数据采集    │  • WiFi/    │  • 设备管理   │  • 远程监控  │
│  • 数据处理    │    TCP/IP   │  • 数据存储   │  • 设备控制  │
│  • MQTT客户端  │  • MQTT桥接 │  • 规则引擎   │  • 数据分析  │
│  • 本地控制    │             │  • 告警通知   │              │
└─────────────────────────────────────────────────────────────┘

二、OneNET 平台配置

2.1 创建产品和设备

  1. 登录 OneNET 平台 (https://open.iot.10086.cn/)
  2. 创建产品:选择 MQTT 协议,填写产品信息
  3. 创建设备:获取设备 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 调试技巧

  1. 串口调试:使用串口打印 MQTT 报文,验证协议正确性
  2. 网络抓包:使用 Wireshark 抓包分析 MQTT 通信过程
  3. OneNET 日志:查看平台端的设备日志,确认数据接收情况

5.2 优化建议

  1. 断线重连:实现 MQTT 断线自动重连机制
  2. 数据缓存:网络异常时缓存数据,恢复后补发
  3. 功耗优化:休眠模式下关闭 WiFi,定时唤醒上传数据
  4. 安全加固:使用 TLS/SSL 加密连接,保护数据安全

5.3 常见问题解决

问题 原因 解决方案
连接超时 WiFi 信号弱 检查 WiFi 密码,增强信号
MQTT 认证失败 设备密钥错误 核对 OneNET 平台设备信息
数据上传失败 JSON 格式错误 验证 JSON 语法,使用 cJSON 库
命令无响应 订阅主题错误 检查订阅主题格式是否正确

专注于matlab/simulink,电子电路,编程