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

Python socketserver.BaseRequestHandler接收文件触发事件实现咨询

解决方法:用线程安全队列实现任务解耦

核心思路是把耗时的预处理/分类逻辑从TCPHandler线程中剥离,用独立工作线程从线程安全队列中取任务执行。TCPHandler只负责接收文件,完成后把文件名扔进队列即可,不用关心后续处理。

具体实现步骤:

1. 创建线程安全任务队列

用Python标准库的queue.Queue,它自带线程锁,无需额外处理同步问题。

2. 启动工作线程处理队列任务

编写专门的工作函数,循环从队列取出文件名,执行预处理和分类逻辑。这个线程可以在服务器启动前就运行。

3. 让TCPHandler访问队列

ThreadedTCPServer实例化TCPHandler时,会把自身作为self.server传给Handler,因此可以把队列绑定到服务器实例上,Handler通过self.server.task_queue就能拿到队列(或者直接用全局队列)。

完整代码示例

import socketserver
import queue
import threading
import os

# 线程安全任务队列
task_queue = queue.Queue()

def worker():
    """工作线程:从队列取文件名,执行预处理和分类"""
    while True:
        filename = task_queue.get()
        # 收到终止信号则退出
        if filename is None:
            break
        try:
            print(f"开始处理文件: {filename}")
            # 这里替换为你的预处理逻辑
            # 如:读取图片、缩放、归一化
            # 再送入分类器执行分类
            # process_image(filename)
            # classify_image(processed_data)
            print(f"文件处理完成: {filename}")
        except Exception as e:
            print(f"处理文件{filename}出错: {str(e)}")
        finally:
            # 标记任务完成,队列可处理下一个任务
            task_queue.task_done()

class ThreadedTCPHandler(socketserver.BaseRequestHandler):
    def handle(self):
        # 现有逻辑:判断请求类型
        data = self.request.recv(1024).decode().strip()
        if data == "PING":
            self.request.sendall(b"PONG")
            return
        
        # 接收文件逻辑(示例)
        # 接收文件名
        filename = self.request.recv(1024).decode().strip()
        # 接收文件大小
        file_size = int(self.request.recv(1024).decode().strip())
        # 接收文件内容
        with open(filename, "wb") as f:
            received = 0
            while received < file_size:
                chunk = self.request.recv(min(1024, file_size - received))
                if not chunk:
                    break
                f.write(chunk)
                received += len(chunk)
        
        # 文件接收完成,将文件名放入任务队列
        task_queue.put(filename)
        self.request.sendall(b"FILE_RECEIVED")

class ThreadedTCPServer(socketserver.ThreadingMixIn, socketserver.TCPServer):
    pass

if __name__ == "__main__":
    HOST, PORT = "0.0.0.0", 9999

    # 启动单工作线程,若需并发处理可启动多个
    threading.Thread(target=worker, daemon=True).start()
    # 多线程示例:启动3个工作线程
    # for _ in range(3):
    #     threading.Thread(target=worker, daemon=True).start()

    # 启动服务器
    with ThreadedTCPServer((HOST, PORT), ThreadedTCPHandler) as server:
        print(f"服务器启动在 {HOST}:{PORT}")
        server.serve_forever()

关键细节说明

  • 线程安全保障:queue.Queue的put()和get()方法自带线程锁,无需担心多Handler线程同时写队列的冲突问题。
  • 优雅退出:若需关闭服务器时终止工作线程,可在服务器关闭前向队列传入None,工作线程收到后会退出循环。
  • 并发能力优化:把耗时操作剥离到工作线程,TCPHandler可快速释放,继续处理新的客户端连接,提升服务器并发响应能力。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 07:00:54