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

如何用Python multiprocessing Pool永久消费队列并重启Worker进程

问题:使用multiprocessing Pool实现动态任务队列+单任务重启进程的Worker系统

我需要用Python的multiprocessing模块实现一个Worker系统,满足以下需求:

  • 监听HTTP请求,动态向任务队列添加任务ID
  • 用固定数量(比如2个)的进程池处理队列任务
  • 每个Worker进程处理仅一个任务后就重启(避免内存泄漏)
  • 进程池永久运行,持续消费动态添加的任务

当前代码的问题是:Pool在处理完初始的2个任务后不会重启进程处理后续任务(比如示例中的任务"c"),使用initializer和maxtasksperchild=1没有达到预期效果。

示例代码:

from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from multiprocessing import Pool, SimpleQueue, current_process

queue = SimpleQueue()

def do_something(q):
    worker_id = current_process().pid
    print(f"Worker {worker_id} spawned")
    item_id = q.get()
    print(f"Worker {worker_id} received id: {item_id}")
    # long_term_operation_that_leaks_memory(item_id)
    # print(f"Worker {worker_id} completed id: {item_id}")

def main():
    with Pool(
        processes=2, initializer=do_something, initargs=(queue,), maxtasksperchild=1
    ):
        queue.put("a")
        queue.put("b")
        queue.put("c")
        server_address = ("", 8000)
        httpd = ThreadingHTTPServer(server_address, BaseHTTPRequestHandler)
        try:
            httpd.serve_forever()
        except (KeyboardInterrupt, SystemExit):
            pass

if __name__ == "__main__":
    main()

解决方案

核心问题在于误用了initializer参数——initializer是进程启动时执行的初始化函数,不是任务处理循环。要实现动态任务消费+单任务重启进程,需要结合Pool.apply_async的异步任务提交机制,配合maxtasksperchild=1让进程处理完一个任务就退出重启。

修改思路

  • 把任务处理函数改为单次任务处理逻辑,而非无限循环的队列消费
  • 启动独立线程,负责从队列中取出任务并提交给Pool
  • 利用maxtasksperchild=1确保每个进程只处理一个任务就被替换
  • 扩展HTTP请求处理器,实现动态添加任务的逻辑

完整代码

from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from multiprocessing import Pool, SimpleQueue, current_process
import threading

queue = SimpleQueue()

def process_single_task(item_id):
    """单次任务处理函数:处理一个任务后退出,由Pool自动重启新进程"""
    worker_id = current_process().pid
    print(f"Worker {worker_id} started processing task: {item_id}")
    # 替换为你的实际任务逻辑,比如可能泄漏内存的操作
    # long_term_operation_that_leaks_memory(item_id)
    print(f"Worker {worker_id} finished task: {item_id}")
    return f"Task {item_id} processed by {worker_id}"

def task_consumer(pool):
    """独立线程:持续从队列取任务并提交给Pool"""
    while True:
        item_id = queue.get()
        # 异步提交任务,不阻塞主线程
        pool.apply_async(process_single_task, args=(item_id,))

class TaskRequestHandler(BaseHTTPRequestHandler):
    """自定义HTTP处理器:接收POST请求添加任务"""
    def do_POST(self):
        content_length = int(self.headers['Content-Length'])
        task_id = self.rfile.read(content_length).decode('utf-8').strip()
        queue.put(task_id)
        self.send_response(200)
        self.send_header('Content-Type', 'text/plain')
        self.end_headers()
        self.wfile.write(f"Added task: {task_id}".encode('utf-8'))
        print(f"HTTP server added task: {task_id}")

def main():
    # 初始化Pool,maxtasksperchild=1确保每个进程只处理一个任务
    with Pool(processes=2, maxtasksperchild=1) as pool:
        # 启动任务消费线程
        consumer_thread = threading.Thread(target=task_consumer, args=(pool,), daemon=True)
        consumer_thread.start()

        # 添加初始测试任务
        queue.put("a")
        queue.put("b")
        queue.put("c")

        # 启动HTTP服务器
        server_address = ("", 8000)
        httpd = ThreadingHTTPServer(server_address, TaskRequestHandler)
        print("HTTP server running on port 8000")
        try:
            httpd.serve_forever()
        except (KeyboardInterrupt, SystemExit):
            print("Shutting down...")
            httpd.shutdown()

if __name__ == "__main__":
    main()

关键说明

  • maxtasksperchild=1:该参数让Pool中的每个进程处理完1个任务后自动退出,Pool会立即启动新进程补充,彻底解决内存泄漏问题
  • 任务消费线程:独立于主线程的线程持续从队列取任务,用apply_async异步提交给Pool,不会阻塞HTTP服务
  • HTTP处理器:自定义TaskRequestHandler处理POST请求,可通过curl -X POST -d "task_d" http://localhost:8000动态添加任务
  • 进程自动重启:每个任务完成后,对应的Worker进程被销毁,Pool自动创建新进程等待下一个任务,完全符合需求

测试验证

运行代码后,输出示例:

HTTP server running on port 8000
Worker 1234 started processing task: a
Worker 5678 started processing task: b
Worker 1234 finished task: a
Worker 9012 started processing task: c  # 新进程替代了1234
Worker 5678 finished task: b
Worker 3456 started processing task: task_d  # 新进程替代了5678,处理HTTP添加的任务

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 17:40:46