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

Python聊天客户端线程接收消息时触发Broken Pipe错误

Python聊天客户端退出时Broken Pipe错误的解决方案

问题根源

  1. 客户端后台接收线程未实现优雅退出机制,当主程序关闭socket后,线程仍尝试对已关闭的socket执行send/recv操作,触发Broken Pipe错误。
  2. 服务器处理!quit命令时直接关闭连接,未同步客户端的收尾流程;广播逻辑中强制调用recv会导致无数据时阻塞,影响消息推送。

修复后的客户端代码

import socket
import threading
import sys

class Client:
    DATA_BUFFER_SIZE = 1024

    def __init__(self, host, port):
        self.client = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
        self.client.connect((host, port))
        self.running = True  # 线程运行标志
        self.receive_thread = threading.Thread(target=self.receive_messages, daemon=True)
        self.receive_thread.start()

    def close(self, send_message=True):
        self.running = False  # 通知后台线程停止
        if send_message:
            try:
                self.client.send('!quit'.encode())
                # 等待服务器确认
                ack = self.client.recv(Client.DATA_BUFFER_SIZE)
                if ack:
                    print(ack.decode())
            except Exception:
                pass
        print('closing client')
        self.client.close()
        sys.exit(0)

    def receive_messages(self):
        while self.running:
            try:
                msg = self.client.recv(Client.DATA_BUFFER_SIZE)
                if not msg:
                    break  # 连接已关闭
                print(f"\n收到消息: {msg.decode()}")
                print('message -> ', end='', flush=True)  # 回到输入提示符
            except Exception:
                if self.running:
                    print("\n连接已断开")
                break

    def loop(self):
        try:
            while self.running:
                msg = input('message -> ')
                if not self.running:
                    break
                self.client.send(msg.encode())
                if msg == '!quit':
                    self.close(send_message=False)
                # 普通消息无需等待ack,交给后台线程处理推送
        except KeyboardInterrupt:
            self.close()

def main():
    client = Client('172.31.196.42', 7777)
    client.loop()

if __name__ == '__main__':
    main()

修复后的服务器代码

import threading
import socket
import sys

class Server:
    DATA_BUFFER_SIZE = 1024
  
    class Connection:
        def __init__(self, server):
            self.conn, self.addr = server.accept()
            print(f'accepted new connection from {self.addr}')
            self.user = 'Anonymous'
            self.running = True
            # 为每个客户端启动单独的消息处理线程
            self.handle_thread = threading.Thread(target=self.handle_messages, args=(server,), daemon=True)
            self.handle_thread.start()

        def close(self):
            self.running = False
            print(f'closing connection to user {self.user}, addr: {self.addr}')
            try:
                self.conn.close()
            except Exception:
                pass
        
        def handle_messages(self, server):
            while self.running:
                try:
                    msg = self.conn.recv(Server.DATA_BUFFER_SIZE)
                    if not msg:
                        # 客户端主动断开
                        server.close_client(server.clients.index(self))
                        break
                    msg_str = msg.decode()
                    print(f'received message from {self.user}: {msg_str}')
                    if msg_str == '!quit':
                        self.conn.send('ack'.encode())
                        server.close_client(server.clients.index(self))
                        break
                    else:
                        # 发送ack确认收到消息
                        self.conn.send('ack'.encode())
                        # 广播消息给其他客户端
                        server.broadcast(f'{self.user}: {msg_str}'.encode(), exclude=self)
                except Exception:
                    server.close_client(server.clients.index(self))
                    break
  
    def __init__(self, port):
        host = socket.gethostbyname(socket.gethostname())
        self.server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
        self.server.bind((host, port))
        self.server.listen()
        print(f'server binded to {host}:{port}')
        self.running = True
        self.clients = []
        self.accept_thread = threading.Thread(target=self.accept_connections, daemon=True)
        self.accept_thread.start()

    def close(self):
        print('closing server')
        for client in self.clients:
            client.close()
        self.server.close()

    def close_client(self, client_id):
        if 0 <= client_id < len(self.clients):
            self.clients[client_id].close()
            self.clients.pop(client_id)

    def accept_connections(self):
        while self.running:
            try:
                self.clients.append(Server.Connection(self.server))
            except Exception:
                break

    def broadcast(self, msg, exclude=None):
        for client in self.clients:
            if client != exclude:
                try:
                    client.conn.send(msg)
                except Exception:
                    self.close_client(self.clients.index(client))

    def loop(self):
        try:
            while self.running:
                input()  # 保持主进程运行,可按回车退出
        except KeyboardInterrupt:
            self.close()
            sys.exit(0)

def main():
    server = Server(7777)
    server.loop()

if __name__ == '__main__':
    main()

关键修复点

  1. 客户端:

    • 添加running标志控制后台线程生命周期,避免连接关闭后仍操作socket。
    • 后台线程专注于接收服务器推送的消息,移除无效的定时ack发送逻辑。
    • 关闭时先通知线程停止,再处理socket操作,确保优雅退出。
  2. 服务器:

    • 为每个客户端单独启动消息处理线程,避免遍历客户端时的阻塞问题。
    • 广播逻辑移除强制recv,直接推送消息,同时处理发送异常(客户端已断开)。
    • 客户端断开时通过recv返回空值检测,实现自动清理连接。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 07:05:05