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

Python Multiprocessing:实现同时运行x个进程并自定义任务队列管理的技术咨询

自定义多进程任务调度方案(支持权重控制)

我完全理解你的需求——既要并行固定数量的任务,又要后续基于任务权重控制总负载,自己管理队列确实是更灵活的选择。下面一步步解决你的问题:

核心问题拆解与解决思路

1. 复杂类对象无法放入队列的问题

multiprocessing的队列依赖Python的pickle序列化机制,复杂类对象只要满足可pickle条件就能传递。如果你的类包含不可pickle的属性(比如文件句柄、网络连接),可以通过自定义序列化方法解决:

class MyComplexTask:
    def __init__(self, task_id, weight, file_path):
        self.task_id = task_id
        self.weight = weight
        self.file = open(file_path, "r")  # 不可pickle的资源
    
    # 自定义序列化逻辑:只保存可序列化的属性
    def __getstate__(self):
        state = self.__dict__.copy()
        del state["file"]  # 移除不可序列化的对象
        state["file_path"] = self.file.name  # 保存重建资源需要的信息
        return state
    
    # 自定义反序列化逻辑:重建不可序列化的资源
    def __setstate__(self, state):
        self.__dict__.update(state)
        self.file = open(state["file_path"], "r")  # 重新打开文件
    
    def main_func(self):
        print(f"Executing task {self.task_id} (weight: {self.weight})")
        # 这里写你的任务逻辑
        time.sleep(2)
        print(f"Task {self.task_id} completed")
        return self.weight

2. 子进程完成通知与任务分配逻辑

我们用两个队列实现通信:

  • 任务队列:主进程放待执行任务,子进程取任务执行
  • 状态队列:子进程完成任务后发送状态,主进程据此更新负载并分配新任务

子进程工作逻辑

子进程持续监听任务队列,收到None时退出,完成任务后向主进程发送完成信号:

def worker(task_queue, status_queue):
    while True:
        task = task_queue.get()
        if task is None:  # 主进程发送的退出标记
            break
        # 执行任务并获取权重(用于后续负载计算)
        completed_weight = task.main_func()
        # 向主进程发送完成状态(任务ID、权重)
        status_queue.put(("completed", task.task_id, completed_weight))

主进程调度逻辑

主进程负责初始化队列、启动子进程、动态分配任务(严格控制权重总和不超过阈值):

import multiprocessing
import time

def main():
    # 配置参数
    MAX_WORKERS = 3  # 最大并行任务数
    MAX_TOTAL_WEIGHT = 10  # 并行任务权重总和阈值
    # 准备所有任务
    tasks = [
        MyComplexTask(1, 4, "task1.txt"),
        MyComplexTask(2, 3, "task2.txt"),
        MyComplexTask(3, 5, "task3.txt"),
        MyComplexTask(4, 2, "task4.txt"),
        MyComplexTask(5, 6, "task5.txt"),
        MyComplexTask(6, 1, "task6.txt"),
    ]
    
    # 初始化队列
    task_queue = multiprocessing.Queue()
    status_queue = multiprocessing.Queue()
    
    # 启动子进程
    workers = []
    for _ in range(MAX_WORKERS):
        p = multiprocessing.Process(target=worker, args=(task_queue, status_queue))
        p.start()
        workers.append(p)
    
    # 任务调度变量
    current_total_weight = 0
    task_index = 0
    completed_task_count = 0
    total_tasks = len(tasks)
    
    # 初始填充任务(不超过权重阈值)
    while task_index < total_tasks and current_total_weight + tasks[task_index].weight <= MAX_TOTAL_WEIGHT:
        task = tasks[task_index]
        task_queue.put(task)
        current_total_weight += task.weight
        task_index += 1
    
    # 监听子进程状态,动态分配新任务
    while completed_task_count < total_tasks:
        status = status_queue.get()
        if status[0] == "completed":
            completed_task_count += 1
            current_total_weight -= status[2]
            # 尝试分配下一个任务(只要权重总和不超限)
            while task_index < total_tasks and current_total_weight + tasks[task_index].weight <= MAX_TOTAL_WEIGHT:
                task = tasks[task_index]
                task_queue.put(task)
                current_total_weight += task.weight
                task_index += 1
    
    # 发送退出信号给所有子进程
    for _ in range(MAX_WORKERS):
        task_queue.put(None)
    
    # 等待所有子进程结束
    for p in workers:
        p.join()
    
    print("All tasks completed successfully!")

if __name__ == "__main__":
    main()

关于Pool的补充说明

你提到的Pool确实可以处理不同函数的并行任务,比如用apply_async提交不同的可调用对象:

from multiprocessing import Pool

class TaskA:
    def run(self):
        print("Running Task A")
        time.sleep(2)

class TaskB:
    def execute(self):
        print("Running Task B")
        time.sleep(1)

def main():
    tasks = [TaskA().run, TaskB().execute, TaskA().run]
    with Pool(2) as pool:
        results = [pool.apply_async(task) for task in tasks]
        for res in results:
            res.get()  # 等待任务完成

但这种方式无法实时跟踪任务权重和负载,所以你选择自定义队列管理的方案完全正确,更适合后续的权重控制需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 12:02:30