You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python使用paho-mqtt时如何在同一MQTT主题下同时实现订阅与发布?

MQTT同主题收发区分及发布逻辑优化

同主题收发消息区分方案

因为设备仅监听test/{imei}单个主题,无法通过主题拆分区分上下行,所以直接约定统一的消息 payload 结构即可:

  • 所有消息统一用JSON格式封装,必须携带msg_type字段标识消息方向
    • 设备上行(设备发给平台):msg_type = "report",携带设备上报数据
    • 平台下行(平台发给设备):msg_type = "command",携带平台下发指令
  • 平台收到消息后先解析msg_type,如果是command类型直接跳过即可,避免消费到自己发布的下行消息,也不会出现循环处理的问题。

发布逻辑优化说明

你现有代码的问题是每次发布都新建MQTT连接、用完就销毁,频繁建连会严重增加性能开销和出错概率,完全没必要,直接复用已经初始化的全局长连接client即可:

  • 建议单独封装发布方法,实现逻辑解耦,后续修改下发规则不需要修改消息接收的处理逻辑
  • 方法入参直接传从上报消息中提取的imei、要下发的指令内容即可,不需要额外传其他参数
  • 不需要在on_message里写复杂的发布逻辑,只需要提取必要参数后调用发布方法即可

注意避坑

你当前代码订阅的是/test/#,如果设备发布的主题是test/123(无前导斜杠),两者是完全不同的MQTT主题,会导致订阅不到设备消息,需要统一主题前缀,建议统一去掉前导斜杠避免匹配错误。

修正后可运行代码

import paho.mqtt.client as mqtt
import json

# 全局MQTT客户端复用
client = mqtt.Client()
user = 'test'
passw = '1111'
client.username_pw_set(user, passw)

# 客户端收到服务端CONNACK响应时的回调
def on_connect(client, userdata, flags, rc):
    print("Connected with result code "+str(rc))
    # 统一主题规则,去掉前导斜杠,和设备发布的主题匹配
    client.subscribe("test/#")

# 单独封装发布方法,解耦逻辑
def publish_command(imei, command_data):
    topic = f"test/{imei}"
    # 下行消息封装固定格式
    payload = json.dumps({
        "msg_type": "command",
        "data": command_data
    })
    client.publish(topic, payload)
    print(f"SEND command to {topic}: {payload}")

# 收到服务端PUBLISH消息时的回调
def on_message(client, userdata, msg):
    print(f"Received message from {msg.topic}: {str(msg.payload)}")
    imei = msg.topic.split('test/')[1]
    try:
        payload = json.loads(msg.payload.decode())
    except json.JSONDecodeError:
        print(f"Invalid message format from {imei}, skip")
        return
    # 过滤自己发布的下行消息,只处理设备上行的上报消息
    if payload.get("msg_type") != "report":
        return
    # 提取设备上报数据处理
    report_data = payload.get("data")
    print(f"Get report from {imei}: {report_data}")
    # 处理完成后下发指令,这里示例下发hello指令,可按业务修改
    publish_command(imei, "hello")

if __name__ == "__main__":
    client.on_connect = on_connect
    client.on_message = on_message
    client.connect("localhost", 1883, 60)
    client.loop_forever()

内容的提问来源于stack exchange,提问作者Kate12345

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.24 18:15:05