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
相关产品推荐
相关产品推荐

