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

Python MQTT订阅程序无报错停止运行问题排查求助

问题描述

我有一个Python程序,用于订阅MQTT Broker并将匹配消息写入数据库。程序运行一段时间(数分钟到数小时不等)后无报错停止/崩溃,VS Code调试器未捕获到任何异常。我是Python新手,不清楚是否遗漏了关键细节。调试器显示脚本仍在运行,MQTT Broker也显示向客户端发送PUBLISH消息,但程序不再输出日志,也不更新数据库。用systemd作为服务运行时,状态显示为"active"且无错误/退出码。想请教哪里处理有误,导致无法看到堆栈跟踪?


程序代码

main.py

#!/usr/bin/python3

from queue import Queue
import paho.mqtt.client as mqtt
import sci_message_pb2
from google.protobuf.json_format import MessageToDict
import ruckus_parser_functions as rks
import datetime as dt


def run_client(q):
    def on_message(client, userdata, message):
        scim = sci_message_pb2.SciMessage().FromString(message.payload)
        q.put(scim)

    broker_address = "xxx.xxx.xxx.xxx"
    client = mqtt.Client("xxx-client")
    client.connect(broker_address)
    client.subscribe("xxx-topic")
    client.on_message = on_message
    client.loop_start()


def run():
    q = Queue()
    run_client(q)
    t = dt.datetime.now()
    print(t.strftime("%Y-%m-%d %H:%M:%S"), "Started...")

    while True:
        try:
            now = dt.datetime.now()
            delta = now-t
            if delta.seconds >= 60:
                print(now.strftime("%Y-%m-%d %H:%M:%S"), "Alive...")
                t = dt.datetime.now()
            msg = q.get()

            # Processing
            msgDict = MessageToDict(msg)
            switchStatus = False

            try:
                switchStatus = msgDict["switchMessage"]["switchStatus"]
            except KeyError:
                pass


            if (switchStatus):
                switchStatus = rks.process_switch_status(msg)

        except Exception as e:
            print("Caught exception:" + str(e))


if __name__ == "__main__":
    run()

ruckus_parser_functions.py

import mysql.connector
from google.protobuf.json_format import MessageToDict
import datetime

class SwitchStatusObj(object):
    switch_id = ""
    uptime = ""
    cpu_percent = 0
    memory_percent = 0
    poe_percent = 0
    status = ""

def process_switch_status (message):
    #print(message.switchMessage.switchStatus)
    statusData = message.switchMessage.switchStatus
    switchStatusData = SwitchStatusObj()
    switchStatusData.switch_id = statusData.id
    switchStatusData.uptime = statusData.uptime
    switchStatusData.cpu_percent = statusData.cpu
    switchStatusData.memory_percent = statusData.memory
    switchStatusData.poe_percent = int((statusData.poeUtilization / statusData.poeTotal) * 100)
    switchStatusData.status = statusData.status
    
    sql = """
    INSERT INTO vsz_switch_status (switch_id, uptime, cpu_percent, memory_percent, poe_percent, status) 
    VALUES (%s,%s,%s,%s,%s,%s) 
    ON DUPLICATE KEY UPDATE 
    uptime = VALUES(uptime), 
    cpu_percent = VALUES(cpu_percent), 
    memory_percent = VALUES(memory_percent),
    poe_percent = VALUES(poe_percent),
    status = VALUES(status)
    ;"""
    vals = (switchStatusData.switch_id, switchStatusData.uptime, switchStatusData.cpu_percent, switchStatusData.memory_percent, switchStatusData.poe_percent, switchStatusData.status)

    insert_to_mysql(sql, vals)

    now = datetime.datetime.now()
    print(now.strftime("%Y-%m-%d %H:%M:%S"), " SwitchStatus updated in MySQL:", switchStatusData.switch_id)
    
    return switchStatusData


def insert_to_mysql (sql, vals):
    try:
        db = mysql.connector.connect(
            host="xxxxxxxxxx",
            user="xxx",
            password="xxxxxxx",
            database="xxxx",
        )
        cursor = db.cursor()
        cursor.execute(sql, vals)
        db.commit()
        db.close()
    except mysql.connector.Error as err:
        print("MySQL error: {}".format(err))

问题排查与修复方案

1. MQTT连接无重连机制

当前代码仅调用client.connect(),未处理连接断开后的重连逻辑。若Broker重启或网络波动导致连接中断,后台loop_start()线程会静默失效,主线程会卡在q.get()处无限等待。

修复:添加断开重连回调:

def run_client(q):
    def on_message(client, userdata, message):
        try:
            scim = sci_message_pb2.SciMessage().FromString(message.payload)
            q.put(scim)
        except Exception as e:
            print(f"解析MQTT消息失败: {str(e)}")
    
    def on_disconnect(client, userdata, rc):
        if rc != 0:
            print("意外断开连接,尝试重连...")
            client.reconnect()

    broker_address = "xxx.xxx.xxx.xxx"
    client = mqtt.Client("xxx-client")
    client.on_disconnect = on_disconnect
    client.connect(broker_address)
    client.subscribe("xxx-topic")
    client.on_message = on_message
    client.loop_start()

2. 阻塞式队列导致程序假死

q.get()默认是无限阻塞的,如果MQTT无消息传入,主线程会卡在该调用,连"Alive..."日志都无法输出(日志检查逻辑在q.get()之后)。

修复:给队列获取添加超时,超时后回到循环执行状态检查:

# 替换原msg = q.get()
try:
    msg = q.get(timeout=1)  # 1秒超时
except Queue.Empty:
    continue  # 超时后跳过,执行Alive检查

3. 未捕获后台线程异常

loop_start()启动的MQTT线程若发生异常(比如Protobuf解析失败),不会被主线程的try-except捕获,会导致线程静默死亡,无法接收消息。

修复:在on_message内部添加异常捕获,避免线程崩溃。

4. 除零错误风险

process_switch_status中statusData.poeTotal可能为0,导致除零错误,该异常未被捕获会终止后续处理。

修复:添加除零检查:

# 替换原poe_percent计算代码
if statusData.poeTotal > 0:
    switchStatusData.poe_percent = int((statusData.poeUtilization / statusData.poeTotal) * 100)
else:
    switchStatusData.poe_percent = 0
    print(f"警告:交换机{statusData.id}的poeTotal为0")

5. 数据库连接资源泄漏

每次调用insert_to_mysql都新建连接,若异常发生可能导致连接未正确关闭,长期运行会耗尽数据库连接池。

优化:使用连接池复用连接:

# 在ruckus_parser_functions.py开头初始化连接池
from mysql.connector.pooling import MySQLConnectionPool

db_pool = MySQLConnectionPool(
    pool_name="mypool",
    pool_size=5,
    host="xxxxxxxxxx",
    user="xxx",
    password="xxxxxxx",
    database="xxxx"
)

def insert_to_mysql(sql, vals):
    try:
        db = db_pool.get_connection()
        cursor = db.cursor()
        cursor.execute(sql, vals)
        db.commit()
    except mysql.connector.Error as err:
        print(f"MySQL错误: {err}")
    finally:
        if 'db' in locals():
            db.close()  # 归还连接到池

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 11:00:44