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

自定义ThreadPool与请求队列:Python玩具店服务器并发请求故障排查

问题分析与解决方案

核心问题诊断

  1. 无法接收后续客户端请求:原服务器代码仅调用一次accept(),之后循环仅处理第一个客户端的请求,未持续监听新连接。
  2. 线程池效率低下:原线程通过轮询队列判断任务,造成CPU空转;且请求解析、库存查询逻辑在主线程执行,易导致主线程阻塞。

修复后的服务器代码

import socket
import json
from threading import Thread, Lock, Condition
from collections import deque

class Server():
    def __init__(self, host, port, num_threads):
        self.items = {
            "tux": {
                "qty": 100,
                "cost": 25.99
            },
            "whale": {
                "qty": 100,
                "cost": 19.99
            }
        }
        self.lock = Lock()
        self.request_queue = deque([])
        self.cond = Condition(self.lock)  # 用条件变量实现线程等待/唤醒,避免轮询空转
        self.thread_pool = [Thread(target=self.serve_request) for _ in range(num_threads)]
        self.s = socket.socket()
        self.s.bind((host, port))
        print("socket binded to port", port)
        self.s.listen(5)
        print("socket is listening")

    def get_item(self, item):
        with self.lock:
            if item not in self.items:
                return -1
            if self.items[item]['qty'] == 0:
                return 0
            self.items[item]['qty'] -= 1
            return self.items[item]['cost']

    def serve_request(self):
        while True:
            with self.cond:
                # 队列为空时线程进入等待,有新任务时被唤醒
                while not self.request_queue:
                    self.cond.wait()
                client_socket = self.request_queue.popleft()
            
            # 完整处理客户端请求
            try:
                data = client_socket.recv(4096)
                if not data:  # 客户端断开连接
                    client_socket.close()
                    continue
                data = json.loads(data.decode('utf-8'))
                cost = self.get_item(data['query'])
                client_socket.send(str(cost).encode('utf-8'))
            except Exception as e:
                print(f"请求处理出错: {e}")
            finally:
                client_socket.close()

    def run(self):
        # 启动线程池,设置守护线程随主进程退出
        for thread in self.thread_pool:
            thread.daemon = True
            thread.start()
        
        # 持续接受新客户端连接
        while True:
            client_socket, addr = self.s.accept()
            print(f"新客户端连接: {addr}")
            with self.cond:
                self.request_queue.append(client_socket)
                self.cond.notify()  # 唤醒等待的线程

host = "127.0.0.1"
port = 12345
server = Server(host, port, 10)  # 线程数建议设为CPU核心数1-2倍,避免过多上下文切换
server.run()

关键修改说明

  • 持续监听新连接:run()方法中循环调用accept(),每次接收新客户端后将socket加入队列,保证能处理多个并发连接。
  • 线程池效率优化:使用Condition替代轮询,线程在队列空时进入等待状态,有新任务时被唤醒,大幅降低CPU资源浪费。
  • 解耦请求处理:主线程仅负责接受连接,完整的请求解析、库存查询、响应逻辑交给线程池处理,避免主线程阻塞。
  • 线程数合理设置:线程数无需设为100,建议根据CPU核心数调整(如1-2倍),过多线程会增加上下文切换开销。

客户端代码(无需修改)

import socket
import json
import random
def main():
    host = "127.0.0.1"
    port = 12345
    s = socket.socket()
    s.connect((host, port))
    while True:
        toys = ["tux", "whale"]
        choice = random.choice(toys)
        message = {"query": str(choice)}
        serialzed_message = json.dumps(message)
        print("requesting: ", choice)
        s.send(serialzed_message.encode('utf-8'))
        data = s.recv(4096)
        print("Server replied: {}".format(str(data.decode('utf-8'))))

if __name__ == "__main__":
    main()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 05:25:39