Python MQTT 发布订阅实战:paho-mqtt、QoS 与断线重连
2026/9/20 7:47:48 网站建设 项目流程

简介:这份围绕 Python 实现 MQTT 发布与订阅的参考资料,面向物联网开发初学者、嵌入式与后端工程师,以及需要快速搭建消息通信链路的自学者。内容以 paho-mqtt 库为主线,梳理 Broker 连接、消息发布、主题订阅与回调处理等关键环节,并给出可对照的函数参数注解,帮助读者理解 QoS 等级、retain 保留标志、心跳间隔等概念在实际场景中的取舍。压缩包共 1 个 pdf 文件,约 45KB,篇幅精简,便于离线查阅与打印标注,适合边编码边对照。目前已有 3400 人学习下载,可见其对动手实践者颇具参考价值。读者可从中获得发布与订阅两端的完整示例代码、connect、publish、subscribe 三个核心函数的参数含义说明,以及消息到达后回调打印的排错思路,为构建 IoT 数据采集、设备指令下发等轻量级实时通信方案提供直接参考。

1. 从一盏灯到一条消息:Python 里跑通 MQTT 发布订阅的现实场景

一个不大的实验室里,温湿度传感器每 5 秒上报一次数据,空调控制器要根据这些数据决定是否启动。如果让传感器直接调用控制器的 HTTP 接口,两边要互相知道地址、处理超时、还要考虑防火墙策略;换成 MQTT,传感器只管往一个 topic 发消息,控制器只管订阅这个 topic,中间靠 broker 解耦,任何一端重启都不影响另一端。Python 在这个链路里通常承担两类角色:跑在工控机或树莓派上的采集程序,以及后端把设备数据桥接进数据库的服务。这篇内容从 mqtt协议详解 里最容易被跳过的几个字段讲起,落到 paho-mqtt 的发布端和订阅端代码、QoS 参数怎么设、断线重连怎么写,最后给一套把消息链路切两半来定位问题的验证手法。刚过完 python安装教程 的人能跟着跑通,已经在做设备接入的人可以对照参数和坑位查漏。

2. MQTT 协议里必须先搞清的四个概念与 Python 客户端选型

2.1 Broker、Topic、Client ID 与 Session 在 Python 侧的含义

MQTT 走的是发布订阅模型,三个角色里 broker 是唯一有状态的组件,它维护一张订阅关系表:谁订阅了哪些 topic 过滤规则,消息到达时按规则分发。Python 客户端代码里所有配置错误,最终都表现为这张表没建对,或者 session 状态被清掉了,所以先把四个概念对齐再写代码,能省掉大量抓包时间。

Topic 是分层字符串,用/分隔,比如lab/room1/temp。层级本身没有语义,broker 只做前缀匹配,但实际项目里建议按「业务域/位置/设备/指标」固定层级,后续做权限控制和数据落库时不用再改。Topic 区分大小写,Lab/room1lab/room1是两条完全不同的路由,这一点在跨语言协作时最容易踩,比如 Java 侧写的常量是首字母大写,Python 侧全小写,结果订阅端一直没消息。

Client ID 是 broker 识别客户端的唯一标识。同一个 broker 上出现两个相同 Client ID 的连接,broker 会按协议把先到的那条连接踢掉,而且是静默踢,Python 侧的on_disconnect里只能看到一个非零的 reason code。常见做法是在 ID 后面拼进程号或主机名后缀,比如sub-room1-01,避免容器扩缩容时实例撞车。

Session 决定 broker 是否替你保存订阅关系和未确认消息。这个开关就是clean_session(在较新的 paho-mqtt 里参数名改成了clean_start)。设成 True 时每次连上都是白纸一张,设成 False 时 broker 会记住你之前的订阅,QoS 1/2 的离线消息在重连后补发。想做「设备断网半小时,回来还能收到期间的告警」,就必须关掉 clean session,同时订阅时用 QoS 1 以上。

2.2 paho-mqtt、gmqtt、aiomqtt 三种 Python 客户端的选型对比

Python 生态里能用的 MQTT 客户端不止一个,选错了会在并发和回调模型上反复绕路。下面这张表是按实际项目里最常见的三种场景整理的判断依据。

客户端库编程模型适合的场景需要留意的地方
paho-mqtt同步 + 后台线程事件循环采集脚本、上位机、快速验证回调在独立线程执行,共享变量要加锁
aiomqtt基于 asyncio 的异步迭代已用 FastAPI/NoneBot 等异步栈的服务必须在线程内,不能用同步阻塞调用
gmqtt原生 asyncio 实现需要精细控制重连、多 broker 切换生态资料相对少,调试要靠日志

如果是第一次写,直接用 paho-mqtt,它是官方维护的参考实现,API 稳定,出问题搜到的答案也最多。已经有一套 asyncio 服务在跑,再往里塞一个带后台线程的 paho-mqtt 会很难受,这种情况选 aiomqtt,订阅端写成async for message in client.messages就能和现有协程共存。gmqtt 更适合做网关类中间件,普通业务代码用不上它的复杂度。

paho-mqtt 在 2.x 之后改了回调签名,创建客户端时要显式声明回调 API 版本,不声明会直接报错,这是从旧教程复制代码时最常见的失败点:

import paho.mqtt.client as mqtt # paho-mqtt 2.x 要求显式指定回调 API 版本,否则构造时抛 ValueError client = mqtt.Client( callback_api_version=mqtt.CallbackAPIVersion.VERSION2, client_id="pub-room1-01", clean_session=True, )

callback_api_version决定on_connecton_disconnect等回调收到的参数个数;client_id要保证同一 broker 上唯一;clean_session=True表示不保留会话,如果后面要做离线补发,这里改成 False 并同步在connect时保持一致。

2.3 用 Docker 起一个本地 MQTT 服务器

本地验证不要直接用公网测试 broker,消息内容会泄露,而且网络抖动会让你误判代码有问题。mosquitto 是最轻的选择,它的 2.x 版本默认只监听本地回环并且禁止匿名连接,所以必须挂一份配置文件进去,否则容器起来了也连不上。

mkdir -p /tmp/mqtt && cd /tmp/mqtt cat > mosquitto.conf <<'EOF' listener 1883 allow_anonymous true persistence true persistence_location /mosquitto/data/ log_dest stdout EOF docker run -d --name mqtt-broker \ -p 1883:1883 \ -v /tmp/mqtt/mosquitto.conf:/mosquitto/config/mosquitto.conf \ eclipse-mosquitto docker logs --tail 20 mqtt-broker

listener 1883指定监听端口;allow_anonymous true只用于本地开发,生产环境要换成密码文件或接入认证;persistence true让 broker 把会话和保留消息写到磁盘,重启不丢,验证离线补发时必须有这一行。启动后用docker logs确认没有Error: Unable to open config file之类的报错,再往下写代码。

2.4 Python 环境准备:venv、pip 与解释器选择

依赖尽量装在虚拟环境里,避免系统 Python 被污染。装完之后跑一条导入语句确认路径,比在编辑器里猜解释器要可靠得多。

python -m venv .venv source .venv/bin/activate # Windows 用 .venv\Scripts\activate pip install paho-mqtt python -c "import paho.mqtt.client as m; print(m.__file__)"

最后一行输出的路径应该落在.venv目录下。如果在 VSCode 里写代码,按Ctrl+Shift+P执行Python: Select Interpreter选中这个虚拟环境,否则编辑器会提示找不到paho,但终端里其实跑得好好的,这种不一致会浪费很多排查时间。aiomqtt 用pip install aiomqtt单独装,两者可以共存。

3. 用 paho-mqtt 写出第一个可复现的发布者与订阅者

3.1 订阅端:on_connect 里 subscribe 才是正确位置

新手最容易犯的错是把subscribe写在connect之前或者之后直接裸调,网络断一次重连上来,订阅就丢了,消息再也不来。正确做法是把订阅动作放进on_connect回调,每次连接建立都会重新执行一遍。

import paho.mqtt.client as mqtt BROKER, PORT = "127.0.0.1", 1883 TOPIC = "lab/room1/temp" def on_connect(client, userdata, flags, reason_code, properties): # reason_code 为 0 才代表连接成功,非 0 时订阅是无效操作 print(f"connected rc={reason_code}") if reason_code == 0: client.subscribe(TOPIC, qos=1) def on_message(client, userdata, msg): # payload 是 bytes,按发布端的编码方式解回来 text = msg.payload.decode("utf-8") print(f"{msg.topic} qos={msg.qos} payload={text}") client = mqtt.Client( callback_api_version=mqtt.CallbackAPIVersion.VERSION2, client_id="sub-room1-01", ) client.on_connect = on_connect client.on_message = on_message client.connect(BROKER, PORT, keepalive=60) client.loop_forever()

keepalive=60表示客户端承诺 60 秒内至少发一次心跳,broker 超过 1.5 倍时间没收到就判定掉线并触发遗嘱;subscribeqos参数决定 broker 转发这条订阅下消息时使用的最大服务质量,实际生效值取发布端 QoS 和订阅端 QoS 的较小者;loop_forever()是阻塞调用,会一直处理网络收发,按Ctrl+C才会退出,脚本类程序用它最省心。

3.2 发布端:publish 的四个参数怎么填

发布端比订阅端简单,但publish的四个参数每一个都有默认值陷阱,尤其是retain。下面这段模拟十次温湿度上报,带发送确认。

import json, time, random import paho.mqtt.client as mqtt client = mqtt.Client( callback_api_version=mqtt.CallbackAPIVersion.VERSION2, client_id="pub-room1-01", ) client.connect("127.0.0.1", 1883, keepalive=60) client.loop_start() # 起后台线程跑网络循环,主线程继续发 for seq in range(10): payload = json.dumps({"seq": seq, "temp": round(20 + random.random() * 5, 2)}) info = client.publish( topic="lab/room1/temp", payload=payload, qos=1, retain=False, ) info.wait_for_publish(timeout=2) # 阻塞到 PUBACK 回来或超时 print(f"rc={info.rc} published={info.is_published()}") time.sleep(1) client.loop_stop() client.disconnect()

topic必须和订阅端完全一致,包括大小写;payload传字符串时会按 UTF-8 编码,传bytes则原样发送,二进制协议建议提前序列化好;qos=1表示至少送达一次,代价是可能重复,接收端要做幂等;retain=True会让 broker 保存这条消息,之后任何新订阅者一连上就立刻收到它,适合发设备当前状态,不适合发高频采样数据,否则 broker 里会残留一份过期值。wait_for_publish只在 QoS 1/2 下有意义,QoS 0 时不能用来判断对方是否收到。

3.3 loop_forever、loop_start、手动 loop 的取舍

循环方式行为适用场景
loop_forever()阻塞当前线程,内部自动重连纯订阅脚本、单职责进程
loop_start()起后台线程跑循环主线程还要做采集、计算、写库
loop(timeout)手动驱动,只处理一轮需要嵌入已有事件循环,比如 PyQt

loop_start()时要注意回调函数是在后台线程执行的,在里面直接改主线程的列表或数据库连接会出并发问题,稳妥做法是用queue.Queue把消息丢给主线程消费。loop_stop()之后再disconnect(),顺序反了会留下未关闭的 socket。

3.4 通配符订阅 + 与 # 的边界规则

+匹配单层,#匹配多层且只能出现在末尾。lab/+/temp能匹配lab/room1/templab/room2/temp,但匹配不到lab/room1/floor2/templab/#能匹配lab下的所有层级。要留意的是#不会匹配以$开头的系统主题,比如$SYS/broker/uptime,想拿 broker 自身的运行指标必须显式订阅$SYS/#。通配符订阅在 broker 侧是按订阅规则逐条比对的,一个客户端订阅几百条细粒度规则会明显增加 broker 负担,能用一条#覆盖就别拆成几十条。

4. 把示例改成能长期运行:QoS、重连、遗嘱与 TLS

4.1 QoS 0/1/2 的真实差别与选择

QoS交互过程可靠性典型用途
0发出去就不管可能丢,最多一次高频传感器采样、可丢的监控指标
1PUBLISH / PUBACK至少一次,可能重复指令下发、状态上报、告警
2四次握手恰好一次,开销最大计费、开关动作等不能重复执行的场景

选 QoS 的本质是问自己两个问题:这条消息丢了会不会出事,重复执行会不会出事。丢了没事、重复也没事的,用 0;丢了不行、重复能靠业务侧去重的,用 1;两个都不行的,用 2。很多项目一上来全用 2,结果在弱网设备上握手包来回四次,延迟直接翻倍,吞吐掉一半。真正需要 2 的场景比想象中少得多。

4.2 断线重连:reconnect_delay_set 与 clean_start

loop_forever()loop_start()内部都带自动重连,但默认是固定间隔,网络长时间不通时会疯狂重试。用reconnect_delay_set改成指数退避,同时把重连事件打进日志,方便事后判断设备在线率。

def on_disconnect(client, userdata, flags, reason_code, properties): # reason_code 非 0 表示非正常断开,需要关注 if reason_code != 0: print(f"unexpected disconnect: {reason_code}") client.on_disconnect = on_disconnect client.reconnect_delay_set(min_delay=1, max_delay=60) # 1s 起步,最长退到 60s

min_delay是首次重连等待秒数,max_delay是封顶值,中间按倍数递增,避免断网时把 broker 的连接数打满。另一个坑是会话清理:如果之前用clean_session=False建了持久会话,重连时又传了True,broker 会立刻删掉旧的订阅关系和排队消息,表现为「重连成功了但消息少了」。这类问题要在连接参数上一以贯之,改配置时同步检查发布端和订阅端。

4.3 遗嘱消息与保留消息:离线告警怎么写

设备掉线这件事,靠订阅端轮询判断既慢又费资源,MQTT 自带的遗嘱机制更直接:客户端在连接时预先登记一条消息,broker 检测到它异常断开(没发 DISCONNECT 就跑掉)时替它发出来。

client.will_set( topic="lab/room1/status", payload="offline", qos=1, retain=True, ) def on_connect(client, userdata, flags, reason_code, properties): if reason_code == 0: client.publish("lab/room1/status", "online", qos=1, retain=True) client.subscribe("lab/room1/cmd", qos=1)

配合retain=True,broker 会保存最后一条状态,监控端一连上就能立刻知道设备当前是在线还是离线,不用等一个心跳周期。注意遗嘱的触发条件是「异常断开」,主动调用disconnect()属于正常断开,遗嘱不会发,所以正常关机时应该由业务代码自己补一条离线状态。

4.4 用户名密码与 TLS 的连接参数

生产环境的 broker 通常开在 8883 端口走 TLS,同时要求账号认证。这两步在 paho-mqtt 里分别在connect之前调用。

client.username_pw_set("device-01", "your-password") client.tls_set( ca_certs="/etc/mqtt/ca.crt", certfile="/etc/mqtt/client.crt", keyfile="/etc/mqtt/client.key", ) client.connect("mqtt.example.com", 8883, keepalive=60)

ca_certs是服务端证书链的根证书,用自签证书时最容易出错的地方就是漏了它,或者把服务端证书当成 CA 传进去;certfilekeyfile只在服务端要求双向认证时才需要,单向认证的场景留空即可。参数顺序和端口要对上,把 TLS 参数配好却仍然连 1883,握手会直接失败并抛出ssl.SSLError,这类错误在日志里通常只有一行,很容易被忽略。

5. 进阶:asyncio 并发订阅与消息链路的验证手法

5.1 用 aiomqtt 把订阅并进 asyncio

后端服务已经是异步栈时,再引入带后台线程的 paho-mqtt 会让上下文切换和异常传播变得混乱。aiomqtt 把订阅写成异步迭代器,写法上和async for读队列几乎一样。

import asyncio import aiomqtt async def consume(): async with aiomqtt.Client("127.0.0.1", 1883) as client: await client.subscribe("lab/+/temp", qos=1) async for message in client.messages: topic = str(message.topic) data = message.payload.decode("utf-8") print(topic, data) asyncio.run(consume())

async with负责连接和断开,退出代码块时自动发 DISCONNECT;client.messages是一个无限异步生成器,断线时 aiomqtt 会按内置策略重连并恢复订阅,不用自己写on_connect。要注意回调式的写法在这里不适用,所有处理逻辑必须在async for循环体内完成,里面别放time.sleep这类阻塞调用,否则整个事件循环会被卡住。

5.2 消息没收到时,先用命令行把链路切两半

订阅端没输出时,不要急着改 Python 代码,先用命令行工具确认 broker 和 topic 是否正常。这一步能把问题范围缩小一半:命令行收到了说明 broker 和发布端没问题,问题在 Python 订阅端;命令行也收不到,那就是发布端或 broker 配置的问题。

# 终端 A:命令行订阅,观察是否有消息 mosquitto_sub -h 127.0.0.1 -p 1883 -t 'lab/#' -q 1 -v # 终端 B:命令行手动发一条 mosquitto_pub -h 127.0.0.1 -p 1883 -t 'lab/room1/temp' -q 1 -m '{"seq":0,"temp":23.1}' # 查看 broker 自身指标,确认连接数和消息计数在变 mosquitto_sub -h 127.0.0.1 -t '$SYS/#' -v

-v会同时打印 topic 和 payload,方便核对路由是否符合预期。图形化工具 MQTTX 连本地 broker 时,Host 填127.0.0.1、端口填 1883、Client ID 随便换一个唯一值,用它的消息列表能看到每条的 QoS 和时间戳,比看控制台输出直观。

5.3 常见报错与参数对照表

现象常见原因处理方向
ConnectionRefusedErrorbroker 没起或端口未映射docker psss -lntp | grep 1883
rc=5 Not authorized服务端要求认证或 ACL 不匹配username_pw_set或检查 topic 权限
订阅成功但收不到历史消息发布时retain=False且会话未持久化retain=True或关闭 clean session
连接反复掉线两个进程用了同一 Client ID 互踢给 ID 拼进程号或主机名后缀
回调函数完全不触发忘了loop_start或没进入事件循环检查循环调用位置

排查时把on_connect里的 reason code、docker logs mqtt-broker里的连接记录、以及$SYS主题下的消息计数放在一起看,三个位置的信息能直接指向问题出在发布端、broker 还是订阅端,比逐行读代码快得多。

本文还有配套的精品资源,点击获取

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

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

立即咨询