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

Touch pHat+MQTT发布脚本出现Broken pipe错误及自动重连需求

解决Touch pHat MQTT闲置断开及自动重连问题

Hey there, let's break down what's going on with your script and fix that frustrating broken pipe error once and for all.

问题原因分析

Your hunch about signal.pause() is spot-on, but there's a bit more to it:

  • Most MQTT brokers (like CloudMQTT) have an idle timeout (usually around 5 minutes by default). If a client doesn't send a heartbeat (PINGREQ) within that window, the broker will drop the connection to free up resources.
  • When you use signal.pause(), it completely blocks the main thread. This means the Paho MQTT client's internal network loop (which handles heartbeats, connection checks, and reconnections) can't run. After 5 minutes of inactivity, the broker cuts the connection. When you press a button again, mqttc.publish() tries to send data over a dead socket, causing the Broken pipe [Errno 32] error.
  • Without any reconnection logic, the client stays disconnected, so subsequent button presses won't send any messages even though there's no error.

解决方法 & 修改后的代码

We need to do three key things:

  1. Replace signal.pause() with a non-blocking loop that lets both Touch pHat events and MQTT network operations run.
  2. Add MQTT connection callbacks to handle automatic reconnections when the link drops.
  3. Add error handling for message publishing to catch issues early.

修改后的MQTT发布脚本

#!/usr/bin/env python
import paho.mqtt.client as mqtt
import signal
import time
import touchphat

# MQTT Configuration
MQTT_BROKER = "m14.cloudmqtt.com"
MQTT_PORT = 21543
MQTT_USER = "user"
MQTT_PASS = "mypassword"
MQTT_TOPIC = "touchphat"

def on_connect(client, userdata, flags, rc):
    print(f"MQTT Connected successfully, return code: {rc}")

def on_disconnect(client, userdata, rc):
    print(f"MQTT Disconnected, return code: {rc}")
    # Auto-reconnect loop
    while True:
        try:
            client.reconnect()
            print("Reconnected to MQTT broker!")
            break
        except Exception as e:
            print(f"Reconnection failed, retrying in 5s: {str(e)}")
            time.sleep(5)

# Initialize MQTT client
mqttc = mqtt.Client("client1", clean_session=False)
mqttc.username_pw_set(MQTT_USER, MQTT_PASS)
mqttc.on_connect = on_connect
mqttc.on_disconnect = on_disconnect

# Initial connection attempt
try:
    # Set keepalive to 60s (shorter than broker's 5min timeout)
    mqttc.connect(MQTT_BROKER, MQTT_PORT, 60)
except Exception as e:
    print(f"Initial connection failed: {str(e)}")
    on_disconnect(mqttc, None, 1)

# Touch pHat Button Callbacks with error handling
@touchphat.on_touch(['A'])
def handle_touch_A(event):
    try:
        publish_result = mqttc.publish(MQTT_TOPIC, payload="A", qos=0)
        # Wait to confirm publish success
        publish_result.wait_for_publish()
        print("Button A pressed - message published")
    except Exception as e:
        print(f"Failed to publish A: {str(e)}")

@touchphat.on_touch(['B'])
def handle_touch_B(event):
    try:
        publish_result = mqttc.publish(MQTT_TOPIC, payload="B", qos=0)
        publish_result.wait_for_publish()
        print("Button B pressed - message published")
    except Exception as e:
        print(f"Failed to publish B: {str(e)}")

@touchphat.on_touch(['Back'])
def handle_touch_Back(event):
    try:
        publish_result = mqttc.publish(MQTT_TOPIC, payload="Back", qos=0)
        publish_result.wait_for_publish()
        print("Button Back pressed - message published")
    except Exception as e:
        print(f"Failed to publish Back: {str(e)}")

@touchphat.on_touch(['Enter'])
def handle_touch_Enter(event):
    try:
        publish_result = mqttc.publish(MQTT_TOPIC, payload="Enter", qos=0)
        publish_result.wait_for_publish()
        print("Button Enter pressed - message published")
    except Exception as e:
        print(f"Failed to publish Enter: {str(e)}")

# Graceful exit handler
def signal_handler(signal, frame):
    print("\nShutting down...")
    mqttc.disconnect()
    exit(0)

signal.signal(signal.SIGINT, signal_handler)

# Start MQTT network loop in background thread
mqttc.loop_start()

# Keep main thread alive to handle Touch pHat events
try:
    while True:
        time.sleep(0.1)  # Reduce CPU usage
except KeyboardInterrupt:
    signal_handler(signal.SIGINT, None)

优化后的MQTT客户端脚本

We should also add auto-reconnection logic to the client, and clean up the message handling:

#!/usr/bin/env python
import signal
import time
from subprocess import call
import paho.mqtt.client as mqtt

def on_connect(client, userdata, flags, rc):
    print("Connected with result code "+str(rc))
    client.subscribe("touchphat")

def on_message(client, userdata, msg):
    payload = msg.payload.decode('utf-8')  # Convert bytes to string
    print(f"{msg.topic} {payload}")
    if msg.topic == "touchphat":
        if payload == "A":
            print("Button A is pressed")
            call(["aplay", "/home/pi/projectfolder/music.wav"])
        elif payload == "B":
            print("Button B is pressed")
        elif payload == "Back":
            print("Button Back is pressed")
        elif payload == "Enter":
            print("Button Enter is pressed")

def on_disconnect(client, userdata, rc):
    print(f"Disconnected from MQTT broker, return code: {rc}")
    while True:
        try:
            client.reconnect()
            print("Reconnected successfully!")
            break
        except Exception as e:
            print(f"Reconnection failed, retrying in 5s: {str(e)}")
            time.sleep(5)

# Initialize client
client = mqtt.Client()
client.on_connect = on_connect
client.on_message = on_message
client.on_disconnect = on_disconnect
client.username_pw_set("user", "mypassword")

# Initial connection
try:
    client.connect("m14.cloudmqtt.com", 21543, 60)
except Exception as e:
    print(f"Initial connection failed: {str(e)}")
    on_disconnect(client, None, 1)

# Graceful exit
def signal_handler(signal, frame):
    print("\nShutting down client...")
    client.disconnect()
    exit(0)

signal.signal(signal.SIGINT, signal_handler)

client.loop_forever()

关键修改说明

  1. Replace signal.pause(): We use mqttc.loop_start() to run the MQTT network loop in a background thread, and a simple while True loop with time.sleep(0.1) to keep the main thread alive for Touch pHat events. This lets both systems run without blocking each other.
  2. Auto-reconnection: The on_disconnect callback triggers an infinite loop that tries to reconnect every 5 seconds until it succeeds. This ensures the script automatically recovers from dropped connections.
  3. Keepalive setting: We set the MQTT keepalive to 60 seconds (the third parameter in connect()), which tells the client to send a heartbeat every 60 seconds. This is shorter than the typical 5-minute broker timeout, so the connection stays active.
  4. Error handling: We wrap publish calls in try/except blocks and use wait_for_publish() to confirm messages are sent, so you get clear feedback if something goes wrong.
  5. Payload decoding: In the client, we decode the payload from bytes to a string to avoid comparison issues (e.g., b"A" vs "A").

测试建议

  • Run both scripts, let them idle for 5+ minutes, then press a button. You should see the reconnection log (if needed) and the message get published successfully.
  • Test manually disconnecting the network briefly to verify the auto-reconnect works.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:02:25