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

如何在Python中实现具备任务队列与取消功能的可控HTTPServer?

解决方案:自主控制连接与任务队列的Python Web服务器

我完全理解你的痛点——ThreadingMixIn的锁机制确实繁琐,而直接在do_POST里开线程又会因为Handler自动关闭连接导致客户端拿不到结果。下面这个方案可以让你完全掌控连接关闭时机,同时独立管理任务队列,不需要适配ThreadingMixIn的行为:

核心思路

  1. 独立任务管理器:用一个单独的TaskManager类来处理任务队列、执行逻辑和取消信号,线程安全且与Web服务解耦。
  2. Handler阻塞等待任务:在do_POST方法中,针对/work请求,我们不直接返回响应,而是阻塞等待任务完成或被取消,确保连接在我们准备好后才关闭。
  3. 手动控制响应流程:完全接管请求体读取、响应发送的时机,避免Handler自动触发连接关闭。

完整代码示例

import threading
import queue
import json
from http.server import BaseHTTPRequestHandler, HTTPServer

class TaskManager:
    def __init__(self):
        self.task_queue = queue.Queue()
        self.cancel_event = threading.Event()
        self.current_task = None
        # 启动后台工作线程
        self.worker_thread = threading.Thread(target=self._worker, daemon=True)
        self.worker_thread.start()

    def _worker(self):
        while True:
            # 等待队列中的任务
            task_func, task_data, result_event, result_dict = self.task_queue.get()
            self.current_task = (task_func, task_data)
            self.cancel_event.clear()  # 重置取消信号

            try:
                if not self.cancel_event.is_set():
                    # 执行指定任务函数
                    result = task_func(task_data)
                    result_dict['value'] = result
                    result_dict['status'] = 'completed'
                else:
                    result_dict['status'] = 'cancelled_before_start'
            except Exception as e:
                result_dict['value'] = str(e)
                result_dict['status'] = 'failed'
            finally:
                self.current_task = None
                result_event.set()  # 通知Handler任务已完成
                self.task_queue.task_done()

    def add_task(self, task_func, task_data):
        """添加任务并返回用于等待的事件和结果存储字典"""
        result_event = threading.Event()
        result_dict = {}
        self.task_queue.put((task_func, task_data, result_event, result_dict))
        return result_event, result_dict

    def cancel_current_task(self):
        """取消当前正在执行的任务"""
        if self.current_task is not None:
            self.cancel_event.set()
            return True
        return False

# 示例业务函数:你可以替换成自己需要执行的逻辑
def sample_business_task(data):
    import time
    # 模拟耗时任务,定期检查取消信号
    for i in range(5):
        if threading.current_thread().cancel_event.is_set():
            raise RuntimeError("Task cancelled by user")
        time.sleep(1)
    return f"Successfully processed data: {data}"

class CustomRequestHandler(BaseHTTPRequestHandler):
    # 类级别共享TaskManager实例,确保全局唯一
    task_manager = TaskManager()

    def do_POST(self):
        if self.path == '/work':
            self._handle_work_request()
        elif self.path == '/cancel':
            self._handle_cancel_request()
        else:
            self.send_error(404, "Requested path not found")

    def _handle_work_request(self):
        # 读取并解析请求体
        try:
            content_length = int(self.headers.get('Content-Length', 0))
            body = self.rfile.read(content_length).decode('utf-8')
            task_data = json.loads(body)
        except Exception as e:
            self.send_response(400)
            self.send_header('Content-Type', 'application/json')
            self.end_headers()
            self.wfile.write(json.dumps({"error": f"Invalid request: {str(e)}"}).encode('utf-8'))
            return

        # 将任务加入队列,获取等待事件和结果容器
        result_event, result_dict = self.task_manager.add_task(sample_business_task, task_data)

        # 阻塞等待任务完成/取消,此时连接保持打开状态
        result_event.wait()

        # 任务完成后发送响应
        self.send_response(200)
        self.send_header('Content-Type', 'application/json')
        self.end_headers()
        self.wfile.write(json.dumps(result_dict).encode('utf-8'))

    def _handle_cancel_request(self):
        success = self.task_manager.cancel_current_task()
        self.send_response(200)
        self.send_header('Content-Type', 'application/json')
        self.end_headers()
        response = {
            "success": success,
            "message": "Current task cancelled" if success else "No running task to cancel"
        }
        self.wfile.write(json.dumps(response).encode('utf-8'))

    # 可选:覆盖默认日志输出,避免控制台冗余信息
    def log_message(self, format, *args):
        pass

if __name__ == '__main__':
    server_address = ('', 8000)
    httpd = HTTPServer(server_address, CustomRequestHandler)
    print(f"Server started on http://localhost:8000")
    httpd.serve_forever()

关键细节解释

  1. TaskManager的线程安全设计:

    • 用queue.Queue实现线程安全的任务队列,无需手动加锁;工作线程持续从队列取任务执行。
    • cancel_event是线程安全的信号量,触发后会中断当前正在执行的任务(需要业务函数定期检查这个信号,比如示例中的循环)。
  2. 连接控制逻辑:

    • 在/work请求处理中,result_event.wait()会阻塞当前Handler线程,直到任务完成或被取消——这样Handler不会提前结束,连接也就不会被自动关闭。
    • 只有当任务处理完成、响应发送完毕后,do_POST方法才会结束,此时连接才会正常关闭。
  3. 取消机制的灵活性:

    • /cancel请求会触发取消信号,正在执行的任务会在下次检查信号时终止;如果任务还未开始执行,会直接标记为取消状态。
    • 你可以根据业务需求调整检查取消信号的频率,比如在耗时操作的关键节点添加检查。

这个方案完全贴合你的需求:自主掌控连接关闭时机,独立管理任务队列,不需要依赖ThreadingMixIn的复杂锁机制,同时完美支持任务的添加与取消功能。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 04:18:45