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

Python MQTT脚本connect_mqtt函数陷入无限循环问题求助

解决MQTT脚本发完开机状态后陷入循环的问题

我写了个Python脚本,用来收发MQTT消息,收到指令后通过SSH给服务器发命令,还得每3分钟发一次MQTT消息报服务器开机状态。但现在出问题了:发完第一条开机状态通知后,脚本就卡进无限循环里了。代码如下:

import time
from paho.mqtt import client as mqtt_client
from paramiko import SSHClient, AutoAddPolicy

broker = 'localhost'
port = 1883
user = "user"
passw = "pass"
topic = "LV-Automation/Server"
client_id = 'ServerPower'

msgState = f"ON"
msgON = f"Server_Encendido"
msgOFF = f"Server_Apagado"

rState = True


def connect_mqtt():
    def on_connect(client, userdata, flags, rc):
        if rc == 0:
            print("Connected to MQTT Broker!")
        else:
            print("Failed to connect, return code %d\n", rc)

    client = mqtt_client.Client(client_id)
    client.username_pw_set("user", "pass")
    client.on_connect = on_connect
    client.connect(broker, port)
    return client

def subscribe(client: mqtt_client):
    global rState
    def on_message(client, userdata, msg):
        global rState
        print(f"Received `{msg.payload.decode()}` from `{msg.topic}` topic")
        message = msg.payload.decode()
        if message == "OFF":
            rState = False
            client.publish(topic, msgOFF)
            client.publish(topic, "APAGADO")
            time.sleep(2)
            client1 = SSHClient()
            client1.load_system_host_keys()
            client1.set_missing_host_key_policy(AutoAddPolicy())
            client1.connect('12.168.1.22', username= 'user', password= 'pass2w')
            stdin, stdout, stderr = client1.exec_command('powerdown')

        elif message == "Win10":
            client.publish(topic, "Iniciando Win 10 VM, resultado...")
            time.sleep(3)
            client2 = SSHClient()
            client2.load_system_host_keys()
            client2.set_missing_host_key_policy(AutoAddPolicy())
            client2.connect('12.16.1.20', username= 'user', password= 'pass2w')
            stdin, stdout, stderr = client2.exec_command('virsh start Windows10')
            output = stdout.read()
            client.publish(topic, output)

        elif message == "Win10OFF":
            client.publish(topic, "Apagando Win 10 VM, resultado...")
            time.sleep(3)
            client3 = SSHClient()
            client3.load_system_host_keys()
            client3.set_missing_host_key_policy(AutoAddPolicy())
            client3.connect('12.16.1.20', username= 'user', password= 'pass2w')
            stdin, stdout, stderr = client3.exec_command('virsh shutdown Windows10')
            output1 = stdout.read()
            client.publish(topic, output1)

    client.subscribe(topic)
    client.on_message = on_message



def prinRun(client):
    global rState
    client.publish(topic, msgON)
    while rState:
        client.publish(topic, msgState)
        time.sleep(250) # Pause for 3 minutes


def run():
    client = connect_mqtt()
    client.loop_start() # Start network loop in separate thread
    subscribe(client)
    prinRun(client)


if __name__ == '__main__':
    run()

我知道代码写得不够规范,但求能解决问题。


问题排查

你说的“connect_mqtt函数陷入无限循环”其实是误解,实际是主线程被prinRun里的while rState循环占死了,同时代码里还有几个坑会导致异常:

  1. MQTT回调里搞阻塞操作:on_message里的time.sleep和SSH操作都是阻塞的,这些代码跑在MQTT的网络线程里,会卡得线程没法处理其他MQTT消息,甚至假死。
  2. 全局变量rState不安全:这个变量被主线程和MQTT线程同时读写,没有同步机制,可能导致状态更新不及时,循环没法正常退出。
  3. SSH客户端用完不关闭:每次执行SSH命令都新建SSHClient但不关闭,时间长了会占满资源。
  4. 时间算错了:time.sleep(250)不是3分钟,3分钟是180秒。
  5. MQTT认证写死了:client.username_pw_set("user", "pass")硬编码了账号密码,没用到你定义的user和passw变量。

修复后的代码

import time
import threading
from paho.mqtt import client as mqtt_client
from paramiko import SSHClient, AutoAddPolicy, SSHException

broker = 'localhost'
port = 1883
user = "user"
passw = "pass"
topic = "LV-Automation/Server"
client_id = 'ServerPower'

msgState = "ON"
msgON = "Server_Encendido"
msgOFF = "Server_Apagado"

# 用线程安全的事件替代全局变量,控制循环启停
running = threading.Event()
running.set()


def connect_mqtt():
    def on_connect(client, userdata, flags, rc):
        if rc == 0:
            print("Connected to MQTT Broker!")
        else:
            print(f"Failed to connect, return code {rc}")

    client = mqtt_client.Client(client_id)
    client.username_pw_set(user, passw)  # 用定义好的变量,别写死
    client.on_connect = on_connect
    client.connect(broker, port)
    return client


def execute_ssh(host, username, password, cmd):
    """封装SSH操作,统一处理连接、命令执行和资源释放"""
    output = b""
    error = b""
    ssh_client = None
    try:
        ssh_client = SSHClient()
        ssh_client.load_system_host_keys()
        ssh_client.set_missing_host_key_policy(AutoAddPolicy())
        ssh_client.connect(host, username=username, password=password)
        stdin, stdout, stderr = ssh_client.exec_command(cmd)
        output = stdout.read()
        error = stderr.read()
    except SSHException as e:
        error = str(e).encode()
    finally:
        if ssh_client:
            ssh_client.close()
    return output, error


def subscribe(client: mqtt_client):
    def on_message(client, userdata, msg):
        print(f"Received `{msg.payload.decode()}` from `{msg.topic}` topic")
        cmd = msg.payload.decode()

        # 把SSH操作放到单独线程里,别阻塞MQTT线程
        def handle_cmd():
            if cmd == "OFF":
                running.clear()  # 停止状态上报循环
                client.publish(topic, msgOFF)
                client.publish(topic, "APAGADO")
                _, err = execute_ssh('12.168.1.22', 'user', 'pass2w', 'powerdown')
                if err:
                    client.publish(topic, f"关机失败:{err.decode()}")

            elif cmd == "Win10":
                client.publish(topic, "正在启动Win10虚拟机...")
                out, err = execute_ssh('12.16.1.20', 'user', 'pass2w', 'virsh start Windows10')
                client.publish(topic, err.decode() if err else out.decode())

            elif cmd == "Win10OFF":
                client.publish(topic, "正在关闭Win10虚拟机...")
                out, err = execute_ssh('12.16.1.20', 'user', 'pass2w', 'virsh shutdown Windows10')
                client.publish(topic, err.decode() if err else out.decode())

        threading.Thread(target=handle_cmd, daemon=True).start()

    client.subscribe(topic)
    client.on_message = on_message


def publish_status(client):
    client.publish(topic, msgON)
    while running.is_set():
        client.publish(topic, msgState)
        time.sleep(180)  # 修正为3分钟
    # 退出前发最后一条离线通知
    client.publish(topic, msgOFF)


def run():
    client = connect_mqtt()
    client.loop_start()
    subscribe(client)
    publish_status(client)
    client.loop_stop()
    client.disconnect()


if __name__ == '__main__':
    try:
        run()
    except KeyboardInterrupt:
        print("脚本被中断")

核心修复点

  • 替换全局变量为线程安全事件:threading.Event能保证多线程下状态更新及时生效,避免循环没法退出的问题。
  • 封装SSH操作:把重复的SSH代码抽成函数,用完自动关闭客户端,防止资源泄漏。
  • SSH操作放单独线程:不让阻塞操作卡MQTT的消息处理线程,保证MQTT通信正常。
  • 修正时间和认证的小错误:把250秒改成180秒,用变量做MQTT认证。
  • 加异常处理:捕获SSH错误并通过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.03 05:30:53