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

Python Socket多线程项目BrokenPipeError及EOFError问题求助

树莓派土壤湿度传感器Socket通信异常问题

我正在开发一个Python项目:树莓派通过传感器读取土壤湿度数据,再通过Socket发送到笔记本电脑,笔记本电脑用UI展示数据并控制流程。服务器启动正常,但客户端连接后出现以下错误:

Exception in thread Thread-2:
Traceback (most recent call last):
  File "/usr/lib/python3.5/threading.py", line 914, in _bootstrap_inner
    self.run()
  File "/usr/lib/python3.5/threading.py", line 862, in run
    self._target(*self._args, **self._kwargs)
  File "/home/pi/Py_Scripts/MoistureSensor/Server.py", line 77, in worker_recv
    self.JSON_Control = pickle.loads(self.data_recv)
EOFError: Ran out of input

Exception in thread Thread-1:
Traceback (most recent call last):
  File "/usr/lib/python3.5/threading.py", line 914, in _bootstrap_inner
    self.run()
  File "/usr/lib/python3.5/threading.py", line 862, in run
    self._target(*self._args, **self._kwargs)
  File "/home/pi/Py_Scripts/MoistureSensor/Server.py", line 70, in worker_send
    self.conn.send(self.data_send)
BrokenPipeError: [Errno 32] Broken pipe

客户端代码

import sys
from socket import *
from threading import Thread
from time import sleep
import pickle
from PySide6 import QtWidgets, QtGui


class MyClient: 
    def __init__(self, server_port, buf_size, host):
        self.server_port = server_port
        self.buf_size = buf_size
        self.host = host
        self.data_send = None
        self.data_recv = None
        self.exit = False

        self.JSON_Control = {
            "measure": False,
            "exit": False,
        }

        self.JSON_Moisture = {
            "moisture_level": 0
        }

        self.socket_connection = socket(AF_INET, SOCK_STREAM)
        self.socket_connection.connect((self.host, self.server_port))
        print("Connected with Server: %s: " % self.host)

        # thread for sending
        self.thread_send = Thread(target=self.worker_send)
        # thread for receiving
        self.thread_recv = Thread(target=self.worker_recv)

        # starting Threads
        self.thread_send.start()
        self.thread_recv.start()

    def worker_send(self):
        while not self.exit:
            self.data_send = pickle.dumps(self.JSON_Control)
            self.socket_connection.send(self.data_send)
            sleep(0.5)

    def worker_recv(self):
        while not self.exit:
            self.data_recv = self.socket_connection.recv(self.buf_size)
            if self.data_recv is not None:
                self.JSON_Moisture = pickle.loads(self.data_recv)


class UiDisplay(QtWidgets.QWidget):
    def __init__(self):
        super().__init__()

        self.moisture_label = None
        self.moisture_level = None
        self.start_button = None
        self.refresh_button = None
        self.stop_button = None
        self.reset_button = None
        self.quit_button = None
        self.create_widgets()

    def create_widgets(self):
        # create a label to display the current time
        self.moisture_label = QtWidgets.QLabel(self)
        self.moisture_label.setFont(QtGui.QFont("Helvetica", 16))
        self.moisture_level = client.JSON_Moisture["moisture_level"]
        self.moisture_label.setText('Bodenfeuchtigkeit: ' + str(self.moisture_level)[:3] + '%')

        # button to start the weight measuring
        self.start_button = QtWidgets.QPushButton("Start measuring", self)
        self.start_button.clicked.connect(self.measure_moisture)

        # button to refresh the Weight label
        self.refresh_button = QtWidgets.QPushButton("Refresh", self)
        self.refresh_button.clicked.connect(self.refresh)

        # button to stop measuring
        self.stop_button = QtWidgets.QPushButton("Stop measuring", self)
        self.stop_button.clicked.connect(self.stop_measuring)

        # button to quit the program
        self.quit_button = QtWidgets.QPushButton("Quit", self)
        self.quit_button.clicked.connect(self.quit)

        # create a layout to hold the widgets
        layout = QtWidgets.QVBoxLayout()
        layout.addWidget(self.moisture_label)
        layout.addWidget(self.start_button)
        layout.addWidget(self.refresh_button)
        layout.addWidget(self.stop_button)
        layout.addWidget(self.reset_button)
        layout.addWidget(self.quit_button)
        self.setLayout(layout)

    def measure_moisture(self):
        client.JSON_Control["measure"] = True
    def refresh(self):
        # get current weight ( niederlag auf ein qm rechnen fehlt noch)
        self.moisture_level = round(client.JSON_Moisture["moisture_level"] - 2.7) / (1 - 2.7)
        print(self.moisture_level)
        # umrechnen von analogem Wert zu Prozentanteil
        print(self.moisture_level)
        # update the weight label with the current time
        self.moisture_label.setText('Bodenfeuchtigkeit:  ' + str(self.moisture_level)[:5] + '%')

    def stop_measuring(self):
        if client.JSON_Control["measure"]:
            client.JSON_Control["measure"] = False
        else:
            pass

    def quit(self):
        QtWidgets.QApplication.instance().quit()
        client.JSON_Control["exit"] = True
        sleep(2)
        client.exit = True
        client.thread_recv.join()
        client.thread_send.join()
        client.socket_connection.close()
        print("Server connection is closed")
        print('exiting...')
        sleep(1)
        sys.exit()


client = MyClient(1957, 1024, "192.168.86.121")

app = QtWidgets.QApplication()
window = UiDisplay()
window.show()
app.exec()

服务端代码

import sys
from socket import *
from threading import Thread
import pickle
from time import sleep
import Adafruit_ADS1x15


adc_channel_0 = 0


class SoilMoistureSensor:
    def __init__(self, gain, sps, dry_voltage, saturated_voltage):
        self.adc = Adafruit_ADS1x15.ADS1115()
        self.raw_data = None
        self.moisture_level = None
        # self.voltage = None
        self.gain = gain
        self.sps = sps
        self.dry_voltage = dry_voltage
        self.saturated_voltage = saturated_voltage


class MyServer:
    def __init__(self, echo_port, buf_size):
        self.buf_size = buf_size
        self.echo_port = echo_port
        self.data_send = None
        self.data_recv = None
        self.exit = False
        self.data_json = None
        self.moisture_sensor = SoilMoistureSensor(2, 32, 1, 2.7)  # gain, sps, saturated_voltage, dry_voltage

        self.JSON_Control = {
            "measure": False,
            "exit": False
        }

        self.JSON_Moisture = {
            "moisture_level": 0,
        }

        self.socket_connection = socket(AF_INET, SOCK_STREAM)
        self.socket_connection.bind(("", self.echo_port))
        self.socket_connection.listen(1)

        print("Server gestartet")
        print("Name des Hosts: ", gethostname())
        print("IP des Hosts: ", gethostbyname(gethostname()))

        self.conn, (self.remotehost, self.remoteport) = self.socket_connection.accept()
        print("Verbunden mit %s %s " % (self.remotehost, self.remoteport))

        # thread for sending
        self.thread_send = Thread(target=self.worker_send)
        # thread for receiving
        self.thread_recv = Thread(target=self.worker_recv)
        # thread for checking Json Control
        self.thread_check_json_control = Thread(target=self.check_json_control)

        # starting Threads
        self.thread_send.start()
        self.thread_recv.start()
        self.thread_check_json_control.start()

    def worker_send(self):
        while not self.exit:
            self.data_send = pickle.dumps(self.JSON_Moisture)
            self.conn.send(self.data_send)
            sleep(0.5)

    def worker_recv(self):
        while not self.exit:
            self.data_recv = self.conn.recv(self.buf_size)
            if self.data_recv is not None:
                self.JSON_Control = pickle.loads(self.data_recv)

    def measure_moisture(self, channel):
        channel = adc_channel_0
        self.moisture_sensor.raw_data = self.moisture_sensor.adc.read_adc(
            channel, self.moisture_sensor.gain, self.moisture_sensor.sps)
        self.JSON_Moisture["moisture_level"] = self.moisture_sensor.raw_data
        print(self.moisture_sensor.raw_data)

    def stop_connection(self):
        self.thread_recv.join()
        self.thread_send.join()
        self.thread_check_json_control.join()
        self.socket_connection.close()
        print("Server connection is closed")
        print('exiting...')
        sys.exit()

    def check_json_control(self):
        while not self.exit:
            if self.JSON_Control["measure"]:
                self.measure_moisture(0)
            if self.JSON_Control["exit"]:
                self.stop_connection()
            sleep(0.5)


server = MyServer(1957, 1024)

问题分析与修复方案

1. 核心问题:Socket流数据边界处理错误

TCP是流式协议,recv(buf_size)不一定能一次性接收完整的pickle数据;同时客户端断开连接时,recv会返回空字节串b''而非None,此时直接调用pickle.loads会触发EOFError。

2. 具体错误点及修复

(1)服务端worker_recv函数修复

原代码判断if self.data_recv is not None无法识别客户端断开的情况,需修改为判断空字节串,并添加异常处理:

def worker_recv(self):
    while not self.exit:
        self.data_recv = self.conn.recv(self.buf_size)
        # 客户端断开时recv返回空字节串
        if not self.data_recv:
            self.exit = True
            break
        try:
            self.JSON_Control = pickle.loads(self.data_recv)
        except pickle.UnpicklingError:
            print("Received incomplete pickle data")

(2)客户端worker_recv函数修复

同服务端逻辑,处理空字节串和不完整数据:

def worker_recv(self):
    while not self.exit:
        self.data_recv = self.socket_connection.recv(self.buf_size)
        if not self.data_recv:
            self.exit = True
            break
        try:
            self.JSON_Moisture = pickle.loads(self.data_recv)
        except pickle.UnpicklingError:
            print("Received incomplete pickle data")

(3)服务端worker_send捕获BrokenPipeError

客户端断开后服务端仍尝试发送数据会触发该错误,需捕获并终止线程:

def worker_send(self):
    while not self.exit:
        try:
            self.data_send = pickle.dumps(self.JSON_Moisture)
            self.conn.send(self.data_send)
            sleep(0.5)
        except BrokenPipeError:
            self.exit = True
            break

(4)客户端UI未定义变量修复

UiDisplay.create_widgets中添加了未初始化的self.reset_button,需删除该行或补充按钮初始化:

# 删除该行
# layout.addWidget(self.reset_button)

(5)线程安全优化

多线程同时读写字典会引发数据竞争,需添加线程锁:

  • 在客户端和服务端类中添加锁:self.lock = threading.Lock()
  • 读写字典时加锁,例如服务端check_json_control:
def check_json_control(self):
    while not self.exit:
        with self.lock:
            measure_flag = self.JSON_Control["measure"]
            exit_flag = self.JSON_Control["exit"]
        if measure_flag:
            self.measure_moisture(0)
        if exit_flag:
            self.stop_connection()
        sleep(0.5)

3. 额外优化建议

  • 用固定长度前缀标记pickle数据长度:发送时先发送4字节的大端整数表示数据长度,接收时先读取长度再接收对应字节数,确保完整接收数据。
  • 客户端退出时先关闭Socket再终止线程,避免线程在Socket关闭后仍尝试读写。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 08:25:32