如何用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,这样只有收到新消息时才会再次转发。
验证方法
- 先启动你的Mosquitto Broker(如果还没启动的话,树莓派上一般是
sudo systemctl start mosquitto) - 运行上面的Python脚本:
python3 mqtt_forwarder.py
- 打开另一个终端,订阅主主题:
mosquitto_sub -t master_topic
- 再开几个终端,给各个来源主题发消息:
# 给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
相关产品推荐
相关产品推荐

