基于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!
异常原因
- 重连后订阅丢失:原代码仅在初始化时执行一次订阅,MQTT连接断开重连后,之前的订阅会失效,导致脚本虽然显示连接成功,但无法接收GPIO控制消息,表现为无法控制GPIO。
- 缺乏断开原因排查机制:未设置
on_disconnect回调,无法记录连接断开的具体原因(如网络波动、Broker重启、心跳超时等),只能看到重连成功日志,难以定位问题根源。 - 连接稳定性配置不足:未设置心跳间隔、重连间隔等参数,默认配置可能在网络不稳定时频繁断开重连。
解决方法
修改要点
- 在
on_connect回调中执行订阅:确保每次连接(包括重连)都重新订阅主题,避免订阅丢失。 - 添加
on_disconnect回调:记录断开原因,辅助排查问题。 - 配置客户端重连与心跳参数:设置
keepalive心跳间隔,开启自动重连。 - 添加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
相关产品推荐
相关产品推荐

