
简介这份资源围绕 Python 实现 MQTT 消息发布与订阅展开面向物联网开发初学者、嵌入式与后端工程师以及需要快速搭建设备间实时通信链路的开发者。内容以 paho-mqtt 库为核心讲解客户端连接、发布publish与订阅subscribe三类接口的参数含义并给出可直接运行的示例程序帮助读者理解发布/订阅模型、主题topic、QoS 服务质量等级与 retain 保留消息等关键概念。资源包为 1 个 PDF 文件大小约 45KB篇幅精炼适合作为随手查阅的速查手册或课堂演示讲义。该资源已有 3400 余人学习说明其在 MQTT 入门场景中具有一定的参考认可度。通过阅读读者可以掌握连接 Broker 的常用写法、消息回调函数的绑定方式以及在本机与远程服务器上分别测试发布和订阅的完整思路为后续构建物联网数据采集与实时推送系统打下基础。1. 从网关每秒几十条读数说起一台边缘网关每秒要往平台推几十条传感器读数。用 HTTP 短连接的做法是每条数据建一次 TCP、带一整套请求头、再等一个响应换成 MQTT网关只和服务端保持一条长连接每帧消息的固定头只有 2 字节载荷按需拼装带宽和电量的差距是数量级的。Python 加 MQTT 最常见的落地面就在这里设备侧采集上报服务端订阅落库控制指令反向下发。反直觉的地方在于很多人第一次用 paho-mqtt 写完 publish就默认消息一定到了服务端。QoS 0 下 publish 只是把数据交给本地 socket 缓冲区网络断开的瞬间写进缓冲区的消息会直接消失on_publish 回调触发的含义也不是「服务端收到了」而是「本地发出去了」。把发布订阅模型、QoS 选择、回调写法和断线重连一条条拆开才是能扛住现场网络的那套代码。2. 先把 MQTT 的 5 个概念站稳再装 Broker 和 paho-mqtt2.1 Broker、Topic、QoS、Retain、Keep Alive 的边界Broker 是消息枢纽负责把发布者投递到某个 Topic 的消息路由给所有订阅了该 Topic 的客户端。它不关心里面装的是 JSON 还是二进制只看 Topic 字符串。常见的服务端有 Mosquitto、EMQXRabbitMQ 则需要额外开启 mqtt 插件才能接 MQTT 客户端。本地开发用 Mosquitto 就够了配置简单、资源占用低。Topic 是斜杠分层的字符串比如sensor/a1/temp大小写敏感不需要提前创建。发布者往一个不存在的 Topic 发消息不会报错只要有人订阅了匹配的过滤器就能收到。这一点和消息队列的「先建队列再发」完全不同也是新手最容易困惑的地方。QoS 决定投递保证级别Retain 决定消息要不要留在服务端Keep Alive 决定连接多久没心跳就判定断开。三者的行为差异直接影响你的重发逻辑和流量成本QoS投递语义是否可能重复是否可能丢失额外开销0最多一次否是最小只发一次1至少一次是否需要 PUBACK 确认2恰好一次否否两次往返握手开销最大Retain 的语义是「服务端保留该 Topic 上最后一条带保留标记的消息」。新订阅者一连上来立刻收到这条消息不用等下一次发布。它适合放设备当前状态、配置项这类「随时来问都要有答案」的数据不适合放秒级变化的时序读数否则每个新订阅者都会被历史数据糊一脸。Keep Alive 是在 connect 时协商的秒数客户端在这个周期内没有其他报文就会发一个 PINGREQ。服务端在 1.5 倍 Keep Alive 时间内没收到任何报文就认为客户端掉线触发遗嘱消息。设成 60 秒是常见起点移动网络下可以压到 30 秒。2.2 用 Docker 起一个 Mosquitto用 mosquitto_sub/mosquitto_pub 验通先起服务端用官方的 eclipse-mosquitto 镜像把 1883 端口映射出来。这条命令带上了--name方便后续 stop/rm配置文件用镜像内置的无认证版本本地开发足够docker run -d --name mosquitto \ -p 1883:1883 -p 9001:9001 \ eclipse-mosquitto:2 \ mosquitto -c /mosquitto-no-auth.conf参数说明-p 1883:1883是标准 MQTT 端口9001是 WebSocket 端口浏览器端和部分前端库会用到。-c /mosquitto-no-auth.conf指定镜像内自带的允许匿名连接的配置生产环境要换成带password_file的配置并禁掉匿名。容器起来后用官方命令行工具验通一条订阅一条发布分两个终端执行# 终端 A订阅通配符主题-v 打印 topic 前缀 mosquitto_sub -h 127.0.0.1 -p 1883 -t sensor//temp -v -q 1 # 终端 B发布一条 QoS 1 消息 mosquitto_pub -h 127.0.0.1 -p 1883 -t sensor/a1/temp -m 23.5 -q 1终端 A 应该输出sensor/a1/temp 23.5。如果终端 B 没有任何报错、终端 A 也收不到优先查两件事容器是否真的在跑docker ps看状态以及主题过滤器是否写错只匹配一层sensor/a1/temp这种三层路径用sensor//temp才能命中。这个「先命令行验通、再上 Python」的顺序很重要它把服务端问题和客户端代码问题隔离开了。2.3 pip 安装 paho-mqtt2.x 与 1.x 的回调签名差异客户端库用 paho-mqttPython 生态里用得最广的一个。安装指定 2.x 大版本pip install paho-mqtt2.0,3.0paho-mqtt 2.x 引入了 CallbackAPIVersion回调函数签名和 1.x 不一样这是从旧代码迁移时最容易踩的坑。1.x 的on_connect(client, userdata, flags, rc)里 rc 是整数2.x 的 VERSION2 回调变成on_connect(client, userdata, connect_flags, reason_code, properties)其中 reason_code 是 ReasonCode 对象判断成功要用reason_code 0或reason_code.is_failure。on_publish同样从三参数变成五参数mid 的位置没变。回调1.x 签名2.x VERSION2 签名on_connect(client, userdata, flags, rc)(client, userdata, flags, reason_code, properties)on_publish(client, userdata, mid)(client, userdata, mid, reason_code, properties)on_message(client, userdata, msg)(client, userdata, msg)on_disconnect(client, userdata, rc)(client, userdata, disconnect_flags, reason_code, properties)提示创建客户端时显式传入mqtt.CallbackAPIVersion.VERSION2否则 2.x 会按 1.x 的签名回退并打印弃用警告一旦你按新签名写了回调就会因为参数错位收到一堆莫名其妙的 TypeError。3. 发布端paho-mqtt 的 connect、publish、loop 调用顺序3.1 二十行跑通一次 QoS 1 上报先看最小可运行的发布脚本。它连本机 Broker循环发 5 条温度值每条都等发布完成再发下一条import time import paho.mqtt.client as mqtt BROKER, PORT 127.0.0.1, 1883 TOPIC sensor/a1/temp def on_connect(client, userdata, flags, reason_code, propertiesNone): # 2.x 里 reason_code 为 0 表示连接成功 print(connected:, reason_code) def on_publish(client, userdata, mid, reason_codeNone, propertiesNone): # mid 是本次 publish 的报文标识可用于对账 print(published mid , mid) client mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, client_idpub-a1) client.on_connect on_connect client.on_publish on_publish client.connect(BROKER, PORT, keepalive60) client.loop_start() # 后台线程跑网络循环主线程继续发数据 for i in range(5): info client.publish(TOPIC, payloadf23.{i}, qos1, retainFalse) print(rc , info.rc, mid , info.mid) info.wait_for_publish(timeout2) # 阻塞直到该条消息发出 time.sleep(1) client.loop_stop() client.disconnect()逻辑说明connect只是发起 TCP 连接并发送 CONNECT 报文真正的收发包由网络循环驱动。loop_start()在后台起一个线程跑 loop这样主线程可以同步调用 publish如果不用后台线程就得在发布后手动loop()或loop_forever()否则 QoS 1 的 PUBACK 收不到、重发也触发不了。参数说明qos1表示至少一次retainFalse表示这条不留存wait_for_publish(timeout2)在 2 秒内没有完成就返回配合info.rc判断失败原因。keepalive60是心跳周期现场网络差可以降到 30。3.2 QoS 与 retain 的组合怎么选这四个组合基本覆盖了实际场景选错的表现通常在联调后期才暴露组合典型场景踩坑点QoS 0 retainFalse高频遥测、丢一条无所谓断网期间产生的数据全丢QoS 1 retainFalse告警、计费上报网络抖动时可能收到重复消息消费端要幂等QoS 1 retainTrue设备在线状态、最新配置每次状态变更都会覆盖新订阅者只拿到最后一条QoS 2 retainTrue指令下发、资金相关两次往返握手吞吐明显下降别拿它做高频通道QoS 1 的重复投递是协议层面的正常行为不是 bug。消费端按msg.topic加一个业务序列号做去重比在发布端强行抬到 QoS 2 划算得多。3.3 client_id、clean_session 与 topic 命名client_id 在同一个 Broker 上必须唯一。两个进程用了同一个 client_id后连上的会把先连上的踢下线现场表现为「发布端时不时断一下又自己好了」很难查。clean_session2.x 里叫 clean_start决定会话要不要保留。设为 False 时Broker 会替你缓存离线期间订阅的 QoS 1/2 消息重连后补发。做断点续传的采集端一般会这么设但前提是 client_id 固定否则每次重连都是新会话缓存等于没有。Topic 命名上我一般按{业务域}/{设备ID}/{指标}三层走设备 ID 里不放斜杠指标名前缀统一。这样订阅端用sensor//temp就能一次拿到所有设备的温度。3.4 发布异常定位on_publish 的 mid 与 on_disconnect 的 rc发布端的排错入口就两个回调。on_publish里的 mid 是本次 publish 的报文标识把它和client.publish()返回的info.mid对上就能确认哪条消息真的走完了握手。on_disconnect的 reason_code 会告诉你断开的原因常见的是 16正常断开、7连接被服务端拒绝以及网络层的直接超时。一个高频误用是在publish之后立刻disconnect。QoS 1 的 PUBACK 还没回来就断开这条消息就被丢了而且不会有任何报错。正确顺序是wait_for_publish之后再loop_stop()和disconnect()。4. 订阅端on_message 回调、通配符与断线重连4.1 subscribe 为什么必须写在 on_connect 里订阅脚本最常见的错误是在connect()之后直接调用client.subscribe()。连接还没建立时调用可能返回失败码或者订阅在会话重置后失效。正确做法是把 subscribe 放进 on_connect 回调这样每次重连都会重新订阅一遍import paho.mqtt.client as mqtt BROKER, PORT 127.0.0.1, 1883 def on_connect(client, userdata, flags, reason_code, propertiesNone): if reason_code 0: # 重连时 on_connect 会再次触发订阅在这里天然幂等 client.subscribe([(sensor//temp, 1), (cmd/#, 2)]) print(subscribed) def on_message(client, userdata, msg): print(f[{msg.topic}] qos{msg.qos} retain{msg.retain} fpayload{msg.payload.decode(utf-8)}) client mqtt.Client(mqtt.CallbackAPIVersion.VERSION2, client_idsub-01) client.on_connect on_connect client.on_message on_message client.reconnect_delay_set(min_delay1, max_delay60) client.connect(BROKER, PORT, keepalive60) client.loop_forever() # 阻塞主线程异常断开后自动重连参数说明subscribe接受一个 (topic, qos) 元组列表一次订阅多个过滤器。msg.qos是 Broker 实际投递时用的级别可能低于订阅时申请的因为发布端 QoS 才是上限。msg.retain为 True 表示这是一条保留消息不是实时产生的。4.2 通配符 与 # 的边界和多主题订阅匹配单层#匹配剩余所有层且只能出现在过滤器末尾。sensor//temp能匹配sensor/a1/temp匹配不了sensor/a1/room/tempsensor/#两者都能匹配。#单独一个主题过滤器的含义是订阅全部消息调试时可以用线上千万不要一旦有其他业务共用 Broker你的消费端会被灌爆。QoS 的最终生效值取发布端和订阅端申请值的较小者。订阅时申请 2、发布端用 0实际仍然是 0不会因为订阅端要求高就变可靠。这一点经常被人误解为「订阅端设 2 就安全了」。4.3 loop_forever 与 loop_start 的差别loop_forever()是阻塞式的内部自己处理重连适合纯消费进程是订阅端的默认选择。loop_start()起后台线程主线程可以继续做别的事适合把订阅集成进一个已经有主循环的程序比如同时跑一个 Web 服务或用例调度器。用loop_start()时要注意回调运行在后台线程里on_message中如果操作了主线程的共享数据结构得自己加锁。没有共享状态的话直接在里面写库、落文件都没问题。4.4 reconnect_delay_set 参数与网络抖动验证reconnect_delay_set(min_delay1, max_delay60)控制重连退避首次断开后等 1 秒重连失败则等待时间翻倍上限 60 秒。设得太短会在服务端挂掉时造成重连风暴设得太长又会让恢复时间变慢1 到 60 秒是较通用的配置。验证断线重连的方法很直接把容器停掉再启回来参数建议值说明keepalive30~60 秒移动网络取小值min_delay1 秒首连失败后的等待max_delay60 秒退避上限clean_start按需需要补发离线消息时置 Falsedocker stop mosquitto sleep 10 docker start mosquitto观察订阅端日志里 on_connect 是否再次触发、订阅是否重新建立。如果容器重启后一直没恢复检查 on_disconnect 打印的 reason_code以及重连次数是否已经触发了 max_delay 上限。5. 交叉验证用 mosquitto_sub、MQTTX 和 $SYS 主题核对 Python 端行为Python 端行为异常时最快的定位方式不是改代码而是换一个客户端去对照。命令行订阅常驻再用 Python 发一条就能判断问题出在哪一侧mosquitto_sub -h 127.0.0.1 -p 1883 -t sensor/# -v -q 1收到说明 Broker 和发布端都没问题问题在订阅端代码收不到说明要去查发布端的info.rc和 Broker 状态。MQTTX 这类桌面客户端连同一个 Broker也能起到同样的对照作用记得把 client_id 改成和 Python 端不同否则会互相挤下线。验证保留消息有个小技巧先发一条带 retain 的消息再启动订阅如果立刻收到说明 retain 生效mosquitto_pub -h 127.0.0.1 -t device/a1/status -m online -r -q 1 mosquitto_sub -h 127.0.0.1 -t device/a1/status -v # 应立即输出 online想确认连接数、消息吞吐这些服务端指标订阅$SYS主题即可Mosquitto 默认每 10 秒更新一次mosquitto_sub -h 127.0.0.1 -t $SYS/broker/clients/connected -v最后是遗嘱消息的验证这是发布端最容易漏配的一项。给客户端设好遗嘱后用kill -9强杀进程模拟真实宕机client.will_set(device/a1/status, payloadoffline, qos1, retainTrue)订阅端应在 Keep Alive 的 1.5 倍时间内收到offline。如果收不到先确认 will_set 是在 connect 之前调用再检查 Keep Alive 是否设得过长导致判定延迟。把这一条和前面的保留消息配合起来用设备上下线的状态面板就能做到新订阅者一连上就看到当前全量状态而不是等下一次心跳才补齐。本文还有配套的精品资源点击获取