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

如何用Python+Mosquitto MQTT实现多主题订阅并轮流转发至单主题

嘿,我刚帮朋友搞定过几乎一模一样的需求,用Python结合paho-mqtt库就能完美实现!下面给你一步步拆解怎么做:

Python实现MQTT多主题轮流转发到主主题

先搞定依赖

首先得装Python的MQTT客户端库,树莓派上直接用pip装就行:

pip install paho-mqtt

核心思路

我们要做的是:

  • 连接到本地的Mosquitto Broker
  • 订阅你指定的3个主题(rpi_gateway/nrf、rpi_gateway/test_topic,还有你没说全的第三个,我这里用rpi_gateway/third_topic占位,你自己替换)
  • 缓存每个主题的最新消息,确保我们有内容可以转发
  • 用轮询的方式,依次把各个主题的消息推送到master_topic
  • 单独开个线程处理转发逻辑,避免阻塞MQTT的订阅循环

完整代码示例

import paho.mqtt.client as mqtt
import threading
import time

# 配置参数
MQTT_BROKER = "localhost"  # 本地Mosquitto地址,默认是localhost
MQTT_PORT = 1883
SOURCE_TOPICS = [
    "rpi_gateway/nrf",
    "rpi_gateway/test_topic",
    "rpi_gateway/third_topic"  # 替换成你的第三个主题
]
MASTER_TOPIC = "master_topic"

# 缓存每个主题的最新消息,键是主题,值是消息内容
message_cache = {topic: None for topic in SOURCE_TOPICS}
# 轮询索引,记录当前要转发的主题位置
current_index = 0
# 线程锁,避免多线程操作缓存时出现冲突
cache_lock = threading.Lock()

def on_connect(client, userdata, flags, rc):
    print(f"Connected with result code {rc}")
    # 订阅所有来源主题
    for topic in SOURCE_TOPICS:
        client.subscribe(topic)
        print(f"Subscribed to {topic}")

def on_message(client, userdata, msg):
    # 收到消息时更新缓存
    with cache_lock:
        message_cache[msg.topic] = msg.payload.decode()
    print(f"Received message from {msg.topic}: {message_cache[msg.topic]}")

def forward_messages(client):
    global current_index
    while True:
        with cache_lock:
            # 只转发有内容的主题
            if message_cache[SOURCE_TOPICS[current_index]] is not None:
                # 发布到主主题,这里可以加上主题标识,方便订阅者知道来源
                payload = f"[{SOURCE_TOPICS[current_index]}] {message_cache[SOURCE_TOPICS[current_index]]}"
                client.publish(MASTER_TOPIC, payload)
                print(f"Forwarded to {MASTER_TOPIC}: {payload}")
                # 可选:转发后清空该主题的缓存,避免重复转发相同内容
                # message_cache[SOURCE_TOPICS[current_index]] = None
        
        # 切换到下一个主题,循环往复
        current_index = (current_index + 1) % len(SOURCE_TOPICS)
        # 这里可以调整转发间隔,比如每1秒轮询一次,根据你的需求改
        time.sleep(1)

if __name__ == "__main__":
    # 创建MQTT客户端实例
    client = mqtt.Client()
    client.on_connect = on_connect
    client.on_message = on_message

    # 连接到Broker
    client.connect(MQTT_BROKER, MQTT_PORT, 60)

    # 启动转发线程
    forward_thread = threading.Thread(target=forward_messages, args=(client,))
    forward_thread.daemon = True
    forward_thread.start()

    # 启动MQTT循环,保持连接和消息监听
    client.loop_forever()

代码关键点解释

  • 消息缓存:用字典message_cache保存每个主题的最新消息,这样即使某个主题暂时没新消息,我们也不会空转发。
  • 轮询逻辑:通过current_index变量循环切换主题,确保每个主题的消息都能被依次转发到主主题。
  • 线程安全:用cache_lock线程锁保护缓存的读写操作,避免多线程同时操作导致的数据混乱。
  • 可选优化:如果你不想重复转发相同的消息,可以在转发后把对应主题的缓存设为None,这样只有收到新消息时才会再次转发。

验证方法

  1. 先启动你的Mosquitto Broker(如果还没启动的话,树莓派上一般是sudo systemctl start mosquitto)
  2. 运行上面的Python脚本:
python3 mqtt_forwarder.py
  1. 打开另一个终端,订阅主主题:
mosquitto_sub -t master_topic
  1. 再开几个终端,给各个来源主题发消息:
# 给rpi_gateway/nrf发消息
mosquitto_pub -t rpi_gateway/nrf -m "Temperature: 25°C"
# 给rpi_gateway/test_topic发消息
mosquitto_pub -t rpi_gateway/test_topic -m "Test message 123"
# 给第三个主题发消息
mosquitto_pub -t rpi_gateway/third_topic -m "Humidity: 60%"

你会看到订阅master_topic的终端里,依次收到各个主题的消息,循环往复。

自定义调整

  • 如果你的Broker需要用户名密码,在client.connect之前加上client.username_pw_set("your_username", "your_password")
  • 可以调整time.sleep(1)的间隔,控制轮询转发的速度
  • 如果不需要显示来源主题标识,可以把payload直接设为message_cache[SOURCE_TOPICS[current_index]]

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 07:39:49