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

Python多线程并发与多服务器通信的客户端问题排查

Python多线程客户端与多服务器持续通信问题修复

核心问题分析

原代码存在两个关键问题:

  • 输入bye后线程终止,socket连接被关闭,无法保留连接用于后续功能
  • 多线程共用input()导致阻塞,同一时间仅能与一台服务器通信

问题根源

  1. 连接关闭问题:原客户端线程函数中,当输入bye时退出循环,函数执行完毕后socket对象被销毁,连接自动关闭。
  2. 通信阻塞问题:每个线程都调用input(),而标准输入是全局共享资源,同一时间只有一个线程能获取输入,导致其他线程被阻塞。

修复方案

客户端修改思路

  • 用全局消息队列统一处理输入,主线程负责获取输入,分发给各个通信线程
  • 输入bye时不终止线程,仅标记暂停发送,保持socket连接
  • 每个通信线程独立监听队列,发送消息并接收服务器响应

修复后的Client.py

import sys
import threading
import socket
from queue import Queue

# 服务器配置
servers = [
    ('127.0.0.1', 6000),
    ('127.0.0.2', 7000)
]

# 全局消息队列,主线程放消息,通信线程取消息
message_queue = Queue()
# 控制线程是否运行的标志,避免输入bye后关闭线程
running = True

def connect_to_server(host, port):
    client_socket = socket.socket()
    client_socket.connect((host, port))
    print(f"已连接到服务器 {host}:{port}")

    while running:
        try:
            # 从队列获取消息,超时1秒避免一直阻塞
            message = message_queue.get(timeout=1)
            if message.lower().strip() == 'bye':
                # 输入bye时不关闭连接,仅跳过发送(可根据需求调整逻辑)
                print(f"服务器 {host}:{port} 连接已保留,暂停发送消息")
                continue
            # 发送消息并接收响应
            client_socket.send(message.encode())
            data = client_socket.recv(1024).decode()
            print(f"从 {host}:{port} 收到响应: {data}")
        except:
            # 队列超时无消息时,继续循环保持连接
            continue
    # 当running变为False时关闭连接(后续功能可控制此标志)
    client_socket.close()
    print(f"已断开与服务器 {host}:{port} 的连接")

def input_handler():
    global running
    while running:
        message = input(" -> ")
        if message.lower().strip() == 'bye':
            # 把bye放入队列,通知所有线程暂停发送
            for _ in servers:
                message_queue.put(message)
            print("已发送暂停指令,所有服务器连接保留")
            # 若需要完全退出,可设置running=False,否则保持线程运行
            # running = False
        else:
            # 把消息放入队列,所有线程都会发送该消息
            for _ in servers:
                message_queue.put(message)

if __name__ == "__main__":
    # 启动通信线程
    threads = []
    for host, port in servers:
        th = threading.Thread(target=connect_to_server, args=(host, port))
        th.start()
        threads.append(th)
    
    # 启动输入处理线程
    input_th = threading.Thread(target=input_handler)
    input_th.start()

    # 等待所有线程结束(若running一直为True,线程会持续运行)
    input_th.join()
    for th in threads:
        th.join()

服务器端优化(可选)

原服务器仅能处理一个客户端连接,若需要支持多客户端,可修改为多线程服务器:

import socket
import sys
import threading

def handle_client(conn, address):
    print(f"新连接: {address}")
    while True:
        data = conn.recv(1024).decode()
        if not data:
            break
        print(f"来自 {address} 的消息: {data}")
        response = input(" -> ")
        conn.send(response.encode())
    conn.close()
    print(f"连接 {address} 已关闭")

def server_program():
    host = sys.argv[1]
    port = int(sys.argv[2])

    server_socket = socket.socket()
    server_socket.bind((host, port))
    server_socket.listen(5)
    print(f"服务器 {host}:{port} 启动,等待连接...")

    while True:
        conn, address = server_socket.accept()
        # 为每个客户端启动新线程处理
        client_th = threading.Thread(target=handle_client, args=(conn, address))
        client_th.start()

if __name__ == '__main__':
    server_program()

修复效果

  • 输入bye后,所有服务器连接保留,仅暂停发送消息,后续可通过输入其他消息恢复通信
  • 主线程统一处理输入,消息会同时发送给所有服务器,实现多服务器并行通信

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 02:40:39