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

基于Paho MQTT的树莓派GPIO控制脚本异常重连问题排查

MQTT控制GPIO脚本异常排查与修复

问题背景

树莓派上运行Python 3.6脚本,通过MQTT控制板载GPIO,以系统服务形式启动(状态active)。脚本可正常运行数天,但突然出现无法控制GPIO的情况,日志持续输出Connected to MQTT Broker!,推测脚本一直在尝试重新连接MQTT Broker。

原代码

# python 3.6
import RPi.GPIO as GPIO
import time
import json

from paho.mqtt import client as mqtt_client
from socket import timeout

# MQTT broker settings
broker = '192.*.*.*'
port = 1883
topic = "RaspberryMK/GPIO"
client_id = '2223'
username = '****'
password = '****'

# Relays GPIO numbers
RELAIS_1_GPIO = 22 # K1 - Relay 1
RELAIS_2_GPIO = 27 # K2 - Spare

GPIO.setmode(GPIO.BCM) # GPIO Numbers instead of board numbers
GPIO.setwarnings(False)
GPIO.setup(RELAIS_1_GPIO, GPIO.OUT) # GPIO Assign mode
GPIO.setup(RELAIS_2_GPIO, GPIO.OUT) # GPIO Assign mode

GPIO.output(RELAIS_1_GPIO, GPIO.HIGH) # Switch off RELAIS_1_GPIO
GPIO.output(RELAIS_2_GPIO, GPIO.HIGH) # Switch off RELAIS_2_GPIO

def relay_on(Relay):
    GPIO.output(Relay, GPIO.LOW) # out
    print("GPIO LOW")
    
def relay_off(Relay):
    GPIO.output(Relay, GPIO.HIGH) # out
    print("GPIO HIGH")
    
def getInfo():
    if GPIO.input(RELAIS_1_GPIO) == 1:
        relaisstatus = "Uit"
    else:
        relaisstatus = "Aan"

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(username, password)
    client.on_connect = on_connect
    client.connect(broker, port)
    return client
    
def publish(client):
    msg_count = 0
    while True:
        time.sleep(1)
        msg = f"messages: {msg_count}"
        result = client.publish(topic, msg)
        # result: [0, 1]
        status = result[0]
        if status == 0:
            print(f"Send `{msg}` to topic `{topic}`")
        else:
            print(f"Failed to send message to topic {topic}")
        msg_count += 1

def subscribe(client: mqtt_client):
    def on_message(client, userdata, msg):
        print(f"Received `{msg.payload.decode()}` from `{msg.topic}` topic")
        
        if msg.topic == "RaspberryMK/GPIO/GPIO22":
            if msg.payload.decode() == "ON":
                relay_on(RELAIS_1_GPIO)
            elif msg.payload.decode() == "OFF":
                relay_off(RELAIS_1_GPIO)
            else:
                print ("Failed")
        elif msg.topic == "RaspberryMK/GPIO/GPIO27":
            if msg.payload.decode() == "ON":
                relay_on(RELAIS_2_GPIO)
            elif msg.payload.decode() == "OFF":
                relay_off(RELAIS_2_GPIO)
            else:
                print ("Failed")
        else:
            print ("Failed")
   
    client.subscribe(topic+"/GPIO22")
    client.subscribe(topic+"/GPIO27")
    client.on_message = on_message
    
def run():
    client = connect_mqtt()
    
    subscribe(client)
    client.loop_forever()


if __name__ == '__main__':
    run()

异常日志

Connected to MQTT Broker!
Connected to MQTT Broker!
Connected to MQTT Broker!

异常原因

  1. 重连后订阅丢失:原代码仅在初始化时执行一次订阅,MQTT连接断开重连后,之前的订阅会失效,导致脚本虽然显示连接成功,但无法接收GPIO控制消息,表现为无法控制GPIO。
  2. 缺乏断开原因排查机制:未设置on_disconnect回调,无法记录连接断开的具体原因(如网络波动、Broker重启、心跳超时等),只能看到重连成功日志,难以定位问题根源。
  3. 连接稳定性配置不足:未设置心跳间隔、重连间隔等参数,默认配置可能在网络不稳定时频繁断开重连。

解决方法

修改要点

  1. 在on_connect回调中执行订阅:确保每次连接(包括重连)都重新订阅主题,避免订阅丢失。
  2. 添加on_disconnect回调:记录断开原因,辅助排查问题。
  3. 配置客户端重连与心跳参数:设置keepalive心跳间隔,开启自动重连。
  4. 添加GPIO操作异常捕获:避免GPIO操作的隐性错误阻塞脚本逻辑。

修改后的完整代码

# python 3.6
import RPi.GPIO as GPIO
import time
import json

from paho.mqtt import client as mqtt_client
from socket import timeout

# MQTT broker settings
broker = '192.*.*.*'
port = 1883
topic = "RaspberryMK/GPIO"
client_id = '2223'
username = '****'
password = '****'

# Relays GPIO numbers
RELAIS_1_GPIO = 22 # K1 - Relay 1
RELAIS_2_GPIO = 27 # K2 - Spare

# GPIO初始化
GPIO.setmode(GPIO.BCM)
GPIO.setwarnings(False)
GPIO.setup(RELAIS_1_GPIO, GPIO.OUT)
GPIO.setup(RELAIS_2_GPIO, GPIO.OUT)

# 默认关闭继电器
GPIO.output(RELAIS_1_GPIO, GPIO.HIGH)
GPIO.output(RELAIS_2_GPIO, GPIO.HIGH)

def relay_on(Relay):
    try:
        GPIO.output(Relay, GPIO.LOW)
        print(f"Relay {Relay} turned ON (GPIO LOW)")
    except Exception as e:
        print(f"Failed to turn on relay {Relay}: {str(e)}")
    
def relay_off(Relay):
    try:
        GPIO.output(Relay, GPIO.HIGH)
        print(f"Relay {Relay} turned OFF (GPIO HIGH)")
    except Exception as e:
        print(f"Failed to turn off relay {Relay}: {str(e)}")
    
def getInfo():
    if GPIO.input(RELAIS_1_GPIO) == 1:
        relaisstatus = "Uit"
    else:
        relaisstatus = "Aan"
    return relaisstatus  # 补充返回值,修复原函数无返回问题

def connect_mqtt():
    def on_connect(client, userdata, flags, rc):
        if rc == 0:
            print("Connected to MQTT Broker!")
            # 每次连接成功后重新订阅主题
            client.subscribe([(topic+"/GPIO22", 0), (topic+"/GPIO27", 0)])
            print("Re-subscribed to control topics")
        else:
            print(f"Failed to connect, return code {rc}")
            
    def on_disconnect(client, userdata, rc):
        print(f"Disconnected from MQTT Broker, return code {rc}")
        # rc != 0表示异常断开,触发自动重连
        if rc != 0:
            print("Attempting to reconnect...")
    
    client = mqtt_client.Client(client_id)
    client.username_pw_set(username, password)
    client.on_connect = on_connect
    client.on_disconnect = on_disconnect
    # 设置心跳间隔(秒),配置重连延迟范围
    client.keepalive = 60
    client.reconnect_delay_set(min_delay=1, max_delay=60)
    client.connect(broker, port)
    return client

def on_message(client, userdata, msg):
    try:
        payload = msg.payload.decode().strip()
        print(f"Received `{payload}` from `{msg.topic}` topic")
        
        if msg.topic == f"{topic}/GPIO22":
            if payload == "ON":
                relay_on(RELAIS_1_GPIO)
            elif payload == "OFF":
                relay_off(RELAIS_1_GPIO)
            else:
                print(f"Invalid payload for {msg.topic}: {payload}")
        elif msg.topic == f"{topic}/GPIO27":
            if payload == "ON":
                relay_on(RELAIS_2_GPIO)
            elif payload == "OFF":
                relay_off(RELAIS_2_GPIO)
            else:
                print(f"Invalid payload for {msg.topic}: {payload}")
        else:
            print(f"Unknown topic: {msg.topic}")
    except Exception as e:
        print(f"Error processing message: {str(e)}")
    
def run():
    client = connect_mqtt()
    client.on_message = on_message
    client.loop_forever()

if __name__ == '__main__':
    try:
        run()
    except KeyboardInterrupt:
        print("Script stopped by user")
    finally:
        GPIO.cleanup()  # 退出时清理GPIO资源,避免引脚占用

额外建议

  • 查看MQTT Broker的日志,确认是否有客户端断开的记录,排查是否是Broker端问题。
  • 监控树莓派的网络状态,确认是否存在频繁断网、IP地址变化等情况。
  • 考虑使用loop_start()替代loop_forever(),配合主线程的异常处理,提升脚本的健壮性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 08:25:58