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
相关产品推荐
相关产品推荐

