如何在同一脚本中同时实现MQTT定时发消息与实时收消息
问题原因
你的代码无法实现预期效果,核心是3个逻辑错误:
- 定时发送的判断逻辑写在了
subscribe()函数内,这个函数仅在脚本启动时执行1次。当你调用client.loop_forever()后,程序会阻塞在MQTT的网络事件循环中,仅会在收到消息、连接状态变更时触发对应回调,不会重复执行subscribe()里的时间差判断,因此只会启动时发送1次消息。 - 消息接收回调里存在变量覆盖问题:处理关机指令时,你用
client = SSHClient()将传入的MQTT客户端实例覆盖为SSH连接对象,后续触发MQTT操作时会直接报错。 - 定时逻辑没有和MQTT持续运行的事件循环绑定,无法被周期触发。
- 原代码把订阅操作写在了
subscribe()函数中,一旦MQTT连接断开重连,不会自动重新订阅主题,会导致收不到关机指令。
修复方案
不需要额外引入多线程,直接用paho-mqtt自带的事件循环调度能力实现周期发送,逻辑和事件循环完全兼容,不会阻塞消息接收,修复后完整可运行代码如下:
import time from paho.mqtt import client as mqtt_client from paramiko import SSHClient, AutoAddPolicy broker = '192.168.15.210' port = 1883 topic = "LV-Automation/Server" topicState = "LV-Automation/Server/State" client_id = 'ServerPower' msgState = "ON" # 定时上报间隔,单位:秒 SEND_INTERVAL = 120 def connect_mqtt(): def on_connect(client, userdata, flags, rc): if rc == 0: print("Connected to MQTT Broker!") # 连接/重连成功后自动订阅主题,避免重连后收不到消息 client.subscribe(topic) else: print(f"Failed to connect, return code {rc}\n") def on_message(client, userdata, msg): print(f"Received `{msg.payload.decode()}` from `{msg.topic}` topic") message = msg.payload.decode() if message == "OFF": client.publish(topicState, "OFF") # 替换变量名,避免覆盖MQTT客户端实例 ssh_client = SSHClient() ssh_client.load_system_host_keys() ssh_client.set_missing_host_key_policy(AutoAddPolicy()) ssh_client.connect('192.168.15.220', username='root', password='Caro1992') stdin, stdout, stderr = ssh_client.exec_command('powerdown') ssh_client.close() client = mqtt_client.Client(client_id) client.username_pw_set("user", "password") client.on_connect = on_connect client.on_message = on_message client.connect(broker, port, keepalive=60) return client def publish_state(client): """上报服务器运行状态""" msg0 = "ENCENDIDO" client.publish(topic, msg0) result = client.publish(topicState, msgState) status = result[0] if status == 0: print(f"Send `{msg0}` to topic `{topic}`") else: print(f"Failed to send message to topic {topic}") def run(): client = connect_mqtt() # 启动后立即上报一次状态 publish_state(client) # 循环处理消息+定时上报:loop方法会阻塞处理网络事件,到达间隔时间后退出执行上报 while True: client.loop(timeout=SEND_INTERVAL) publish_state(client) if __name__ == '__main__': run()
效果说明
- 脚本启动后立刻发送一次在线状态,之后每2分钟自动重复上报,不会阻塞消息监听
- 收到
OFF指令时会正常执行关机操作,不会出现变量覆盖导致的运行报错 - MQTT连接意外断开后,客户端会自动重连,重连成功后自动恢复订阅和定时上报逻辑
- 不需要额外维护时间戳变量,逻辑更稳定,不存在多线程竞争问题
内容的提问来源于stack exchange,提问作者Apha25
相关产品推荐
相关产品推荐

