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

Python多线程:如何让已有Processor对象共享参数队列?

实现自定义多进程任务分发(基于队列)

你可以借助multiprocessing模块的Queue和Process来实现需求:让已创建的5个Processor实例作为独立工作进程,从共享队列中动态获取任务执行,无需预分配参数。以下是完整实现代码:

import time
import numpy as np
from multiprocessing import Process, Queue

class Processor :
    def __init__(self, name) :
        self.name = name

    def process(self, arg) :
        print(f'{self.name} : processing {arg}...')
        time.sleep(arg)
        print(f'{self.name} : processing {arg}... DONE')

def worker(processor, task_queue):
    # 循环从队列取任务,直到收到结束标记None
    while True:
        arg = task_queue.get()
        if arg is None:
            # 收到结束信号,退出循环
            break
        processor.process(arg)

if __name__ == '__main__':
    l_processors = [Processor(f'Processor_{i}') for i in range(5)]
    l_arguments = list(range(10))
    np.random.shuffle(l_arguments)

    # 1. 创建任务队列,放入所有待处理参数
    task_queue = Queue()
    for arg in l_arguments:
        task_queue.put(arg)
    
    # 2. 为每个Worker添加结束标记(5个None,对应5个Processor)
    for _ in range(len(l_processors)):
        task_queue.put(None)
    
    # 3. 创建并启动子进程
    processes = []
    for processor in l_processors:
        p = Process(target=worker, args=(processor, task_queue))
        processes.append(p)
        p.start()
    
    # 4. 等待所有子进程执行完成
    for p in processes:
        p.join()

关键逻辑说明

  • 任务队列:所有待处理参数先放入Queue,队列会自动处理多进程间的同步,确保同一任务不会被多个Worker重复获取。
  • Worker函数:每个子进程运行此函数,持续从队列拉取任务,执行processor.process(arg);当取到None时,判定为任务结束,退出进程。
  • 结束标记:添加与Worker数量相同的None,确保每个Worker都能收到结束信号,避免进程无限阻塞。
  • 进程管理:为每个已创建的Processor实例绑定一个子进程,启动后通过join()等待所有任务完成。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 13:23:23