☰
从零上手MQTT:用阿里云物联网平台实现两个设备互发消息
2026/10/9 12:36:40 网站建设 项目流程

1. 为什么我建议新手从"两个设备互发消息"开始学MQTT

如果你刚开始接触物联网开发,面对一堆协议名词和云平台控制台,很容易陷入"每个按钮都认识,但连起来就不知道在干嘛"的状态。我见过太多人卡在第一步:设备接进来了,但不知道数据到底有没有发出去、发给了谁、对方收没收到。这种"黑盒感"是学习物联网最大的障碍。

MQTT协议恰好是打破这个障碍的最佳切入点。它的核心模型极其简单——发布者把消息丢给一个中间人(Broker),订阅者从中间人那里取走消息。你不需要关心对方设备的IP地址,不需要处理复杂的网络穿透,只需要约定好一个"话题"(Topic),双方就能通信。阿里云物联网平台提供的MQTT Broker服务,把这个中间人角色托管到了云端,你只需要在控制台点几下、在代码里填几个参数,就能让两个虚拟设备互相收发消息。

这篇内容就是带你走完这个完整流程:从阿里云控制台创建产品、注册设备、获取三元组,到本地用Python写两个独立进程模拟两个设备,一个持续上报温度数据,另一个订阅并接收。全程不需要真实硬件,一台能跑Python的电脑就够。我会把每一步的意图讲清楚,把容易踩的坑提前标出来,让你在动手的时候心里有底。

适合谁看?如果你写过一点Python,知道什么是变量和函数,但对物联网云平台完全陌生,这篇就是为你准备的。如果你已经用过其他物联网平台,想快速了解阿里云物联网平台的MQTT接入方式,也可以直接跳到第3节开始看操作细节。

2. 动手之前,先把MQTT和阿里云物联网平台的关系理清楚

2.1 MQTT协议到底解决了什么问题

MQTT的全称是Message Queuing Telemetry Transport,翻译过来叫"消息队列遥测传输"。名字很长,但核心思想可以用一个生活场景解释:你住在一个小区里,想给邻居送东西,但不知道邻居具体住哪栋哪户。于是你把东西放到小区门口的快递柜,告诉邻居"东西在3号柜,取件码是123"。邻居凭取件码去取,你们不需要直接见面。

在这个类比里,快递柜就是Broker,3号柜就是Topic,取件码就是消息内容。MQTT的设计目标就是让设备之间不需要知道彼此的网络位置,只通过一个公共的Broker和约定的Topic来交换数据。这种"发布-订阅"模式比传统的"请求-响应"模式更适合物联网场景,因为设备数量可能成千上万,网络环境也不稳定,Broker可以帮它们解耦。

MQTT协议有几个关键特性你需要知道:第一,它基于TCP长连接,设备连上Broker后保持在线,有新消息时Broker主动推送,不需要设备轮询;第二,它支持三种QoS等级,QoS 0是"发出去就不管",QoS 1是"至少送达一次",QoS 2是"确保只送达一次",你可以根据数据重要性选择;第三,它有一个"遗嘱消息"机制,设备异常断线时Broker可以自动发布一条预设消息,通知其他设备"我掉线了"。

2.2 阿里云物联网平台在MQTT之上加了什么

阿里云物联网平台本质上是一个托管的MQTT Broker,但它不是裸的MQTT,而是在上面加了一层设备身份认证和管理体系。裸MQTT的Broker通常允许匿名连接,或者用简单的用户名密码。阿里云物联网平台要求每个设备必须用"三元组"来认证:ProductKey(产品密钥)、DeviceName(设备名称)、DeviceSecret(设备密钥)。这三个东西组合起来,通过一套签名算法生成MQTT连接所需的用户名和密码。

为什么要这么设计?因为物联网设备数量大、分布广,如果只靠简单的用户名密码,一旦泄露整个产品下的设备都可能被冒用。三元组机制让每个设备有独立的身份,即使某个设备的密钥泄露,也只影响那一个设备。另外,阿里云物联网平台还提供了Topic权限管理,你可以精确控制某个设备只能发布或订阅哪些Topic,防止设备越权访问。

对于我们要做的"两个设备通信"实验,阿里云物联网平台的角色就是:提供Broker地址、验证设备身份、转发消息。两个设备各自用三元组登录,一个往Topic A发消息,另一个订阅Topic A,消息就会经过阿里云Broker转发到订阅方。

2.3 为什么用Python而不是其他语言

做这个实验,语言选择其实很多:C语言适合嵌入式设备,JavaScript适合Web端,Java适合企业级应用。但我推荐Python,原因有三个:第一,Python的MQTT客户端库paho-mqtt非常成熟,安装简单,API直观,几行代码就能跑起来;第二,Python在电脑上运行不需要交叉编译,不需要烧录固件,改一行代码就能重新运行,调试成本极低;第三,Python的语法接近自然语言,你可以把精力集中在理解MQTT的通信逻辑上,而不是被语言本身的复杂性分散注意力。

当然,如果你后续要把代码移植到真实硬件上,比如ESP32或树莓派,Python也是支持的。所以这个实验的代码不是"玩具",它的核心逻辑可以直接复用到真实项目中。

3. 阿里云控制台操作:从零创建产品和两个设备

3.1 开通物联网平台并创建产品

首先你需要有一个阿里云账号,并且完成实名认证。登录阿里云控制台后,搜索"物联网平台",进入后如果是第一次使用,会提示你开通服务。阿里云物联网平台有免费额度,对于我们的实验来说完全够用,不用担心费用问题。

开通后进入"公共实例"(新版控制台可能叫"实例概览"),点击"创建产品"。这里有几个关键字段需要填写:

  • 产品名称:随便起一个,比如"温湿度监测实验",这个名称只是给你自己看的,不影响通信。
  • 品类:可以选择"自定义品类",也可以选"温湿度传感器",选什么不影响功能,只是影响控制台展示的图标。
  • 节点类型:选择"直连设备"。如果你的设备需要通过网关中转,才选"网关子设备",我们这里两个设备都是直接连阿里云的,所以选直连。
  • 连接方式:选择"Wi-Fi"或"蜂窝"都可以,这只是标记用途。
  • 数据格式:这一步非常关键,选择"Alink JSON"。Alink是阿里云定义的一套JSON格式规范,虽然我们也可以选"透传/自定义",但Alink JSON对新手更友好,因为控制台可以直接解析展示,调试时能看到结构化的数据。

创建产品后,你会看到产品列表里多了一条记录,记下它的ProductKey,后面代码里要用。

3.2 为两个设备分别注册身份

进入刚才创建的产品,点击"设备列表",然后"添加设备"。我们需要添加两个设备,分别命名为device_1和device_2。DeviceName在同一个产品下必须唯一,你可以用任何字符串,但建议用有意义的名称,方便后面区分。

添加完成后,每个设备都会生成一个DeviceSecret。这个密钥只会在创建时显示一次,关闭页面后就看不到了,所以一定要立即复制保存。如果你不小心关了页面,可以在设备详情里点击"重置"来生成新的DeviceSecret,但重置后旧密钥立即失效,如果设备已经烧录了旧密钥就需要重新烧录。

现在你手上有两组三元组:

设备ProductKeyDeviceNameDeviceSecret
设备1你的ProductKeydevice_1设备1的密钥
设备2你的ProductKeydevice_2设备2的密钥

ProductKey是产品级别的,两个设备共用同一个ProductKey,但DeviceName和DeviceSecret不同。

3.3 配置Topic和权限

阿里云物联网平台默认会为每个产品预置一些系统Topic,比如用于设备属性上报的/sys/{ProductKey}/{DeviceName}/thing/event/property/post。但我们要做的是两个设备之间的自定义通信,所以需要自己创建自定义Topic。

在产品的"Topic类列表"里,点击"自定义Topic",定义一个Topic。比如我们可以定义:

  • 上行Topic(设备发布消息用):/{ProductKey}/device_1/user/update
  • 下行Topic(设备订阅消息用):/{ProductKey}/device_2/user/get

但这样定义的话,设备1只能发布,设备2只能订阅,不够灵活。为了实验方便,我们可以定义一个两个设备都能发布和订阅的Topic,比如:

/{ProductKey}/experiment/pubsub

然后在"Topic类列表"里设置这个Topic的操作为"发布和订阅"。接着进入每个设备的"Topic列表",点击"订阅Topic",把/experiment/pubsub添加进去。这样两个设备就都有权限往这个Topic发消息,也都能从这个Topic收消息。

注意:阿里云物联网平台的Topic权限控制比较严格,如果设备没有订阅某个Topic,即使Broker收到了消息也不会转发给它。所以这一步不能省。

4. 本地环境搭建与Python客户端代码编写

4.1 安装paho-mqtt库

在命令行里执行:

pip install paho-mqtt

如果你用的是Python 3,pip会自动安装最新版。目前paho-mqtt已经更新到2.x版本,API和1.x有一些差异,网上很多老教程用的是1.x的写法,直接复制可能会报错。我下面给出的代码是基于2.x版本的,如果你安装的是1.x,需要把mqtt.Client()的调用方式改一下。

安装完成后,可以用pip show paho-mqtt确认版本号。

4.2 理解阿里云MQTT连接参数的计算方式

阿里云物联网平台不接受直接用三元组作为MQTT用户名密码,而是要求用三元组计算出一组签名。具体来说:

  • ClientId:格式为{DeviceName}|securemode=3,signmethod=hmacsha256|,其中securemode=3表示TCP直连,signmethod=hmacsha256表示用HMAC-SHA256算法签名。
  • Username:格式为{DeviceName}&{ProductKey}。
  • Password:用DeviceSecret作为密钥,对特定字符串做HMAC-SHA256运算,然后转成十六进制字符串。

需要签名的内容格式是:clientId{ClientId}deviceName{DeviceName}productKey{ProductKey},注意这里没有分隔符,就是按这个顺序拼接。

举个例子,假设:

  • ProductKey =pk123456
  • DeviceName =dev1
  • DeviceSecret =secret123

那么:

  • ClientId =dev1|securemode=3,signmethod=hmacsha256|
  • Username =dev1&pk123456
  • 签名内容 =clientIddev1|securemode=3,signmethod=hmacsha256|deviceNamedev1productKeypk123456
  • Password = HMAC-SHA256(签名内容,secret123) 的十六进制字符串

这个计算过程不需要你手动做,Python代码里用hmac和hashlib库几行就能搞定。但理解这个机制很重要,因为如果连接失败,你需要知道是哪个参数出了问题。

4.3 编写设备1的发布端代码

设备1的角色是"发布者",它每隔几秒往Topic发一条模拟的温度数据。代码如下:

import paho.mqtt.client as mqtt import hmac import hashlib import time import json import random # 三元组信息,替换成你自己的 PRODUCT_KEY = "你的ProductKey" DEVICE_NAME = "device_1" DEVICE_SECRET = "设备1的DeviceSecret" # 阿里云MQTT接入点,公共实例的地址通常是这个格式 MQTT_HOST = f"{PRODUCT_KEY}.iot-as-mqtt.cn-shanghai.aliyuncs.com" MQTT_PORT = 1883 # 自定义Topic TOPIC = f"/{PRODUCT_KEY}/experiment/pubsub" def calculate_sign(client_id, device_name, product_key, device_secret): """计算阿里云MQTT密码签名""" content = f"clientId{client_id}deviceName{device_name}productKey{product_key}" sign = hmac.new( device_secret.encode('utf-8'), content.encode('utf-8'), hashlib.sha256 ).hexdigest() return sign def on_connect(client, userdata, flags, rc): if rc == 0: print("[设备1] 连接成功") else: print(f"[设备1] 连接失败,返回码:{rc}") def on_publish(client, userdata, mid): print(f"[设备1] 消息已发布,mid={mid}") # 构造连接参数 client_id = f"{DEVICE_NAME}|securemode=3,signmethod=hmacsha256|" username = f"{DEVICE_NAME}&{PRODUCT_KEY}" password = calculate_sign(client_id, DEVICE_NAME, PRODUCT_KEY, DEVICE_SECRET) # 创建客户端并连接 client = mqtt.Client(client_id=client_id, protocol=mqtt.MQTTv311) client.username_pw_set(username, password) client.on_connect = on_connect client.on_publish = on_publish print(f"[设备1] 正在连接 {MQTT_HOST}:{MQTT_PORT} ...") client.connect(MQTT_HOST, MQTT_PORT, keepalive=60) client.loop_start() # 持续发布模拟数据 try: while True: temperature = round(random.uniform(20.0, 30.0), 2) payload = { "device": "device_1", "temperature": temperature, "timestamp": int(time.time()) } result = client.publish(TOPIC, json.dumps(payload), qos=1) print(f"[设备1] 发布温度:{temperature}°C") time.sleep(5) except KeyboardInterrupt: print("[设备1] 停止发布") client.loop_stop() client.disconnect()

这段代码的关键点:client.loop_start()会启动一个后台线程处理网络收发,这样主线程可以继续执行发布逻辑。如果你用client.loop_forever(),它会阻塞主线程,你就没法在循环里发布消息了。

4.4 编写设备2的订阅端代码

设备2的角色是"订阅者",它连上Broker后订阅同一个Topic,然后等待消息到达。代码如下:

import paho.mqtt.client as mqtt import hmac import hashlib import json PRODUCT_KEY = "你的ProductKey" DEVICE_NAME = "device_2" DEVICE_SECRET = "设备2的DeviceSecret" MQTT_HOST = f"{PRODUCT_KEY}.iot-as-mqtt.cn-shanghai.aliyuncs.com" MQTT_PORT = 1883 TOPIC = f"/{PRODUCT_KEY}/experiment/pubsub" def calculate_sign(client_id, device_name, product_key, device_secret): content = f"clientId{client_id}deviceName{device_name}productKey{product_key}" sign = hmac.new( device_secret.encode('utf-8'), content.encode('utf-8'), hashlib.sha256 ).hexdigest() return sign def on_connect(client, userdata, flags, rc): if rc == 0: print("[设备2] 连接成功") client.subscribe(TOPIC, qos=1) print(f"[设备2] 已订阅:{TOPIC}") else: print(f"[设备2] 连接失败,返回码:{rc}") def on_message(client, userdata, msg): try: data = json.loads(msg.payload.decode('utf-8')) print(f"[设备2] 收到消息 - 来自:{data.get('device')}," f"温度:{data.get('temperature')}°C," f"时间戳:{data.get('timestamp')}") except json.JSONDecodeError: print(f"[设备2] 收到非JSON消息:{msg.payload}") client_id = f"{DEVICE_NAME}|securemode=3,signmethod=hmacsha256|" username = f"{DEVICE_NAME}&{PRODUCT_KEY}" password = calculate_sign(client_id, DEVICE_NAME, PRODUCT_KEY, DEVICE_SECRET) client = mqtt.Client(client_id=client_id, protocol=mqtt.MQTTv311) client.username_pw_set(username, password) client.on_connect = on_connect client.on_message = on_message print(f"[设备2] 正在连接 {MQTT_HOST}:{MQTT_PORT} ...") client.connect(MQTT_HOST, MQTT_PORT, keepalive=60) client.loop_forever()

设备2的代码用loop_forever(),因为它不需要主动做其他事情,只需要被动等待消息。on_message回调会在每次收到消息时被触发。

4.5 运行两个脚本并观察通信效果

打开两个命令行窗口,一个运行设备1的脚本,另一个运行设备2的脚本。你应该会看到设备1每隔5秒打印一条"发布温度"的日志,设备2几乎同时打印"收到消息"的日志。

如果一切正常,设备2收到的温度值应该和设备1发布的一模一样。这说明消息经过了阿里云Broker的转发,从设备1到达了设备2。

你可以试着把设备2的脚本停掉,观察设备1是否还在正常发布。然后再启动设备2,看它是否能继续收到后续消息。这个过程中,设备1完全不知道设备2的存在,它只是往Topic发消息,这就是发布-订阅模式的解耦特性。

5. 调试过程中最容易卡住的几个地方

5.1 连接返回码不是0怎么办

on_connect回调里的rc参数是MQTT协议的连接返回码。如果rc不等于0,说明连接被Broker拒绝了。常见的返回码和原因:

返回码含义可能原因
1协议版本不支持客户端协议版本和Broker不匹配
2ClientId无效ClientId格式不对,或包含非法字符
3服务器不可用Broker地址或端口写错
4用户名或密码错误签名计算错误,或三元组填错
5未授权设备没有权限连接,检查产品状态

最常见的是返回码4,也就是密码错误。这时候你需要检查:DeviceSecret是否复制完整(有没有多余空格)、ProductKey和DeviceName是否对应、签名内容的拼接顺序是否正确。我建议在代码里把计算出的password打印出来,和阿里云控制台"设备详情"里的"MQTT连接参数"对比一下。控制台会直接显示正确的ClientId、Username和Password,你可以逐字符比对。

5.2 消息发出去了但对方收不到

如果设备1显示发布成功,但设备2没有收到消息,按以下顺序排查:

第一,检查设备2是否真的订阅成功了。on_connect里调用subscribe后,Broker会返回一个SUBACK报文,但paho-mqtt默认不会打印这个确认。你可以在on_subscribe回调里打印确认信息,确保订阅生效。

第二,检查Topic是否完全一致。阿里云的Topic是大小写敏感的,/experiment/pubsub和/Experiment/PubSub是两个不同的Topic。另外,Topic前后的斜杠也不能多也不能少。

第三,检查设备的Topic权限。在阿里云控制台的设备详情里,查看"Topic列表",确认/experiment/pubsub这个Topic的权限是"发布和订阅"。如果只勾了"发布",设备2订阅时会失败,但paho-mqtt可能不会报错,只是收不到消息。

第四,检查QoS等级。如果设备1用QoS 0发布,设备2用QoS 1订阅,消息仍然能到达,但QoS 0不保证送达。建议两边都用QoS 1,这样至少能保证消息到达一次。

5.3 设备频繁掉线重连

如果你看到设备反复打印"连接成功",说明它在不断掉线重连。可能的原因有:

  • KeepAlive设置太短:connect时的keepalive参数默认是60秒,如果网络延迟大,Broker可能在60秒内没收到心跳就断开连接。可以适当调大到120秒。
  • ClientId冲突:如果两个设备用了相同的ClientId,Broker会踢掉先连接的那个。确保每个设备的ClientId包含唯一的DeviceName。
  • 网络不稳定:如果你在本地运行,检查网络是否正常。阿里云公共实例的接入点在国内,一般延迟很低。

paho-mqtt有自动重连机制,但需要你手动调用client.reconnect_delay_set()来设置重连间隔。默认情况下,连接断开后不会自动重连,你需要在on_disconnect回调里调用client.reconnect()。

5.4 阿里云控制台看不到消息内容

阿里云物联网平台的"日志服务"可以查看设备的上行和下行消息,但默认可能没有开启。你需要在控制台里找到"日志服务"或"消息轨迹",开启后就能看到每条消息的Topic、Payload和时间戳。这对于调试非常有用,因为你可以确认消息到底有没有到达Broker。

另外,控制台的"设备详情"里有一个"在线调试"功能,你可以直接在网页上给设备发消息,或者查看设备最近上报的数据。这个功能在验证设备是否正常在线时很方便。

6. 从实验到实用:几个可以立刻上手的扩展方向

6.1 让两个设备互相"对话"而不是单向广播

目前的实验是设备1发、设备2收,属于单向通信。你可以很容易地把它改成双向:让设备2也发布消息到另一个Topic,设备1订阅那个Topic。这样两个设备就能互相收发,形成一个简单的对话系统。

具体做法是定义两个Topic:/experiment/device1_to_device2和/experiment/device2_to_device1。设备1订阅后者、发布到前者,设备2订阅前者、发布到后者。然后在各自的on_message回调里根据消息内容决定是否回复。

6.2 用设备影子实现离线消息

MQTT的发布-订阅模式有一个天然限制:如果订阅者不在线,消息就丢了。阿里云物联网平台提供了"设备影子"功能来解决这个问题。设备影子本质上是一个JSON文档,存储了设备的期望状态和最新状态。即使设备离线,你也可以更新影子,设备上线后会自动同步。

对于我们的实验,你可以把设备1的温度数据写入设备2的影子,设备2上线后读取影子就能拿到最新的温度值。这个功能在实际项目中非常有用,因为物联网设备经常因为网络问题离线,不能假设对方永远在线。

6.3 用规则引擎把消息转发到其他服务

阿里云物联网平台的"规则引擎"可以把MQTT消息转发到其他阿里云服务,比如数据库、消息队列、函数计算等。你可以写一条SQL规则,把/experiment/pubsub的消息筛选出来,转发到一个数据库表中。这样你就能在数据库里看到所有历史温度数据,而不是只在命令行里看实时输出。

规则引擎的配置在控制台里是可视化的,不需要写代码。你只需要定义数据来源(Topic)、筛选条件(SQL语句)和目标(数据库表),阿里云会自动帮你完成转发。这个功能对于做数据分析和持久化存储非常方便。

6.4 把代码移植到真实硬件上

如果你手头有ESP32或树莓派,可以把设备1的代码移植过去。ESP32可以用MicroPython运行类似的paho-mqtt代码,树莓派可以直接跑Python脚本。移植时需要注意:硬件上的网络环境可能和电脑不同,Wi-Fi信号强度、路由器防火墙设置都可能影响MQTT连接。建议先在电脑上跑通,再移植到硬件,这样出问题时容易定位是代码问题还是网络问题。

另外,真实硬件上的DeviceSecret需要妥善保存,不要硬编码在代码里然后上传到公开仓库。可以用配置文件或环境变量来存储敏感信息,避免泄露。

7. 我在这个实验里踩过的坑和总结的经验

第一次做这个实验的时候,我卡在密码计算上整整一个下午。当时我按照网上的教程拼接签名内容,但那个教程用的是旧版签名算法,拼接顺序和现在不一样。后来我直接在阿里云控制台的设备详情里找到了"MQTT连接参数"一栏,里面直接显示了正确的ClientId、Username和Password。我把控制台显示的值复制到代码里,连接立刻就成功了。所以我的建议是:先用控制台生成的标准参数跑通连接,再自己写签名计算代码。这样你可以先确认网络和权限没问题,再把精力集中在签名算法上。

另一个坑是Topic权限。我一开始只给设备1配置了发布权限,给设备2配置了订阅权限,但忘了给设备2也配置发布权限。后来我想让设备2回复消息时,发现发布失败。阿里云的Topic权限是分开控制的,发布和订阅是独立的权限位,你需要根据实际需求勾选。对于实验来说,直接给两个设备都勾上"发布和订阅"最省事。

还有一个经验是关于QoS的选择。在局域网环境里,QoS 0几乎不会丢消息,因为网络很稳定。但在真实的物联网场景里,设备可能通过蜂窝网络连接,信号时好时坏,这时候QoS 1就很有必要。QoS 1的代价是消息可能重复,所以你的应用层需要做去重处理,比如用消息里的timestamp或msgId来判断是否已经处理过。

最后说一个调试技巧:如果你不确定消息有没有发出去,可以在阿里云控制台的"日志服务"里查看消息轨迹。每条消息都会有记录,包括发送时间、Topic、Payload大小和QoS等级。这个功能比在代码里打印日志更可靠,因为它记录的是Broker实际收到的消息,不受客户端代码影响。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询