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

