Python ThreadPool的map方法传入多可迭代对象与任务均分问题
解决方案
核心问题原因
- 原生
map方法仅支持向目标函数传递单个参数,你传入的zip生成的元组会被当成单个参数传入需要2个参数的f,会触发参数数量不匹配报错 - 直接用
zip(workers, work)会因为workers只有2个元素,仅生成2个任务,剩余8个任务会被直接丢弃,无法完成10个任务的分配需求
方案1:使用starmap方法(无需修改原函数)
starmap会自动将可迭代对象中的每个元组拆分为多个位置参数传递给目标函数,同时使用itertools.cycle循环生成worker标识,匹配任务总数量:
from multiprocessing.pool import ThreadPool import itertools def f(worker,work): print(f"{worker} | {work}") workers = ["worker1","worker2"] work = [1,2,3,4,5,6,7,8,9,10] if __name__ == "__main__": # 循环生成worker序列,和任务数量对齐 worker_sequence = itertools.cycle(workers) task_params = zip(worker_sequence, work) with ThreadPool(len(workers)) as p: # starmap自动拆包元组为多个参数 p.starmap(f, task_params)
运行后两个worker会各分配到5个任务,执行顺序由线程调度决定,符合你的预期要求。
方案2:修改函数适配map方法
如果你坚持要用map,可以修改目标函数接收单个元组参数,自行拆包即可:
from multiprocessing.pool import ThreadPool import itertools # 修改函数接收单个参数,自行拆包 def f(args): worker, work = args print(f"{worker} | {work}") workers = ["worker1","worker2"] work = [1,2,3,4,5,6,7,8,9,10] if __name__ == "__main__": worker_sequence = itertools.cycle(workers) task_params = zip(worker_sequence, work) with ThreadPool(len(workers)) as p: p.map(f, task_params)
补充说明
如果你不需要自定义worker标识,只是想确认任务是两个线程均分的,也可以直接在函数内部打印线程ID,不需要手动传worker参数:
import threading def f(work): print(f"{threading.current_thread().name} | {work}")
这种情况下直接p.map(f, work)即可,pool大小设为2会自动把10个任务均分给两个线程执行。
内容的提问来源于stack exchange,提问作者Minh Doan
相关产品推荐
相关产品推荐

