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

如何在同一脚本中同时实现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 12:01:14