如何让Python multiprocessing.pool.map_async在函数调用间设置等待间隔
解决方案
你无法通过map_async的原生参数实现任务提交间隔,它的设计逻辑是一次性将所有任务推送进进程池的待执行队列。你可以通过以下两种方案实现提交延时,规避黑盒程序的并发文件访问冲突:
方案1:改用apply_async逐个提交任务并添加间隔
这是最灵活可控的方案,你可以在每提交1个/一组任务后等待指定时长,确保先提交的任务已经完成文件初始化读取,再提交下一个针对同文件的任务:
import multiprocessing import time def 黑盒分析函数(文件路径): # 此处替换为你调用黑盒程序的逻辑 pass if __name__ == "__main__": 任务列表 = ["文件1路径", "文件1路径", "文件2路径", "文件2路径"] 进程池大小 = multiprocessing.cpu_count() 提交间隔 = 4 # 可按测试结果在3-5秒区间调整 pool = multiprocessing.Pool(processes=进程池大小) 结果对象列表 = [] for 任务参数 in 任务列表: res = pool.apply_async(黑盒分析函数, args=(任务参数,)) 结果对象列表.append(res) # 每次提交后等待指定时长 time.sleep(提交间隔) # 等待所有任务执行完成 pool.close() pool.join() # 获取所有任务返回结果 所有结果 = [res.get() for res in 结果对象列表]
方案2:针对同文件任务做分组提交
如果你有大量不同文件的任务,不需要给所有任务都加间隔,只需要保证针对同一个文件的两次提交之间有3-5秒间隔即可,避免不必要的时间浪费:
import multiprocessing import time from collections import defaultdict def 黑盒分析函数(文件路径): # 此处替换为你调用黑盒程序的逻辑 pass if __name__ == "__main__": 任务列表 = ["文件1路径", "文件1路径", "文件2路径", "文件2路径", "文件3路径"] 提交间隔 = 4 # 按文件路径对任务做分组 路径分组 = defaultdict(list) for idx, 路径 in enumerate(任务列表): 路径分组[路径].append(idx) pool = multiprocessing.Pool() 结果存储 = [None] * len(任务列表) # 先提交每个文件的第一次分析任务 for 路径, 索引列表 in 路径分组.items(): first_idx = 索引列表[0] res = pool.apply_async(黑盒分析函数, args=(路径,)) 结果存储[first_idx] = res # 等待指定间隔后提交每个文件的第二次分析任务 time.sleep(提交间隔) for 路径, 索引列表 in 路径分组.items(): if len(索引列表) >= 2: second_idx = 索引列表[1] res = pool.apply_async(黑盒分析函数, args=(路径,)) 结果存储[second_idx] = res pool.close() pool.join() 所有结果 = [r.get() for r in 结果存储]
注意:如果你的任务量远大于进程池大小,需要注意避免任务在队列里积压导致同文件的两个任务仍被同时调度执行。这种情况可以进一步调整逻辑,每轮只提交和进程池大小相同的任务,执行完一轮再提交下一轮。
内容的提问来源于stack exchange,提问作者Acceptable Meanderings
相关产品推荐
相关产品推荐

