自定义ThreadPool与请求队列:Python玩具店服务器并发请求故障排查
问题分析与解决方案
核心问题诊断
- 无法接收后续客户端请求:原服务器代码仅调用一次
accept(),之后循环仅处理第一个客户端的请求,未持续监听新连接。 - 线程池效率低下:原线程通过轮询队列判断任务,造成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
相关产品推荐
相关产品推荐

