切换pathos规避multiprocessing的PicklingError时报Pool not running怎么办
报错原因
pathos.multiprocessing.ProcessingPool的uimap方法返回惰性迭代器,调用该方法时不会立刻执行任务,只有当你实际迭代返回的responses对象时,任务才会被提交到进程池执行。- 你的代码在调用
uimap后直接执行了pool.close()和pool.join()销毁了进程池,后续执行reduce操作遍历responses提交任务时,进程池已经停止运行,因此抛出Pool not running错误。
修复方案
有两种可行的修改方式:
方案1:调整操作顺序,先迭代完所有结果再关闭池
你可以先强制触发responses的完整迭代,等所有任务执行完成后再销毁进程池,示例代码如下:
from pathos.multiprocessing import ProcessingPool as Pool from functools import reduce pool = Pool(processes=9) responses = pool.uimap(execute_queries, [session]*32, [query]*32,[parse_query]*32, [filter_values]*32, range(32)) # 先强制迭代所有结果,触发所有任务执行完成 response_list = list(responses) # 再关闭、回收进程池 pool.close() pool.join() response = reduce(reduce_dic, response_list)
注:你代码中定义的tuples变量未实际使用,可以直接删除。
方案2:替换为阻塞式的map方法
如果你不需要惰性迭代的特性,可以直接用map方法替换uimap,map会阻塞到所有任务执行完成并返回结果列表,后续再关闭进程池不会报错:
from pathos.multiprocessing import ProcessingPool as Pool from functools import reduce pool = Pool(processes=9) # map为阻塞调用,所有任务执行完成后才会返回结果列表 responses = pool.map(execute_queries, [session]*32, [query]*32,[parse_query]*32, [filter_values]*32, range(32)) pool.close() pool.join() response = reduce(reduce_dic, responses)
内容的提问来源于stack exchange,提问作者Shivangi Singh
相关产品推荐
相关产品推荐

