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

如何捕获MQTT.connect()的ConnectionRefusedError并实现自动重连?

MQTT连接失败自动重试方案

问题核心

调用G_mqtt_client.connect()时,若MQTT Broker未就绪(比如设备重启顺序导致Broker后启动),会直接抛出ConnectionRefusedError导致程序崩溃,根本无法触发on_connect或on_disconnect回调,需要实现静默重试直到连接成功。

解决思路

  • 捕获连接异常:用try-except包裹connect()调用,阻止程序直接崩溃
  • 实现智能重试机制:失败后延迟重试,推荐用指数退避策略(重试间隔逐步拉长),避免频繁请求占用资源
  • 结合内置重连逻辑:成功连接后,利用paho-mqtt的loop机制处理后续断开场景,在on_disconnect回调中触发重连
  • 控制loop启动时机:确保连接成功后再启动loop_start(),避免无效的循环运行

修改后的代码示例

import time

def mqtt_initialisation():
    mod = '(messaging.initialisation) '
    # 绑定回调函数
    G_mqtt_client.on_connect = on_mqtt_connect
    G_mqtt_client.on_disconnect = on_mqtt_disconnect
    G_mqtt_client.on_message = message_received
    G_mqtt_client.on_log = on_mqtt_log
    G_mqtt_client.username_pw_set(params.G_broker_user, params.G_broker_pwd)
    
    # 重试参数配置
    retry_delay = 5
    max_retry_delay = 60
    
    while True:
        try:
            main_log.debug(mod + "尝试连接MQTT Broker...")
            G_mqtt_client.connect(params.G_broker_url, params.G_broker_port)
            # 连接成功后启动消息循环
            G_mqtt_client.loop_start()
            main_log.debug(mod + "MQTT连接成功,消息循环已启动")
            break
        except ConnectionRefusedError:
            main_log.warning(mod + f"连接被拒绝,{retry_delay}秒后重试")
            time.sleep(retry_delay)
            # 指数退避,避免频繁重试
            retry_delay = min(retry_delay * 2, max_retry_delay)
        except Exception as e:
            main_log.error(mod + f"连接出现未知错误: {str(e)},{retry_delay}秒后重试")
            time.sleep(retry_delay)
            retry_delay = min(retry_delay * 2, max_retry_delay)

def on_mqtt_connect(client, userdata, flags, rc):
    main_log.info(f"MQTT连接成功,结果码: {str(rc)}")
    do_subscribe()

def on_mqtt_disconnect(client, userdata, rc):
    main_log.info("MQTT客户端已断开连接")
    # 非主动断开时自动重试
    if rc != 0:
        main_log.info("尝试重新连接MQTT Broker...")
        client.reconnect()

def on_mqtt_log(client, userdata, level, buff):
    main_log.debug(f"MQTT日志: {buff}")

补充说明

  • 指数退避机制:第一次重试等待5秒,之后每次翻倍,直到最大60秒,避免短时间内频繁请求给Broker造成压力
  • on_disconnect回调处理后续断开:连接成功后若Broker重启,客户端会触发断开回调,此时调用client.reconnect()会让消息循环自动发起重连
  • 全局异常捕获:除了ConnectionRefusedError,还兼容DNS解析失败、网络中断等其他异常场景,确保程序稳定性

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 02:05:20