如何获取Python multiprocessing.Pool执行时剩余待处理输入项数量
解决方案
问题根源
你遇到的问题由两个原因导致:
- 多进程并发向标准输出写入内容时没有同步机制,直接打印会出现内容错位、乱序
multiprocessing.Value的自增操作不是原子操作,不加锁的场景下会出现计数丢失、统计不准的问题
方案1:子进程内带锁统计(适合任务函数内部需要获取剩余任务量的场景)
通过初始化方法将共享计数器和进程锁传入子进程,操作计数和打印时加锁保证原子性,锁的开销极小,不会引入额外的计算负担:
import multiprocessing def init_globals(cnt, lck, total_cnt): global counter, lock, total counter = cnt lock = lck total = total_cnt def f(x): # 业务逻辑放在锁外,不影响并行效率 res = x * 2 # 这里替换为你的实际业务逻辑 with lock: counter.value += 1 remaining = total - counter.value print(f"Done with task: {counter.value}, remaining tasks: {remaining}") return res if __name__ == '__main__': rangeval = range(1000) total = len(rangeval) counter = multiprocessing.Value('i', 0) lock = multiprocessing.Lock() pool = multiprocessing.Pool( processes=multiprocessing.cpu_count(), initializer=init_globals, initargs=(counter, lock, total) ) pool.map(f, iterable=rangeval) pool.close() pool.join()
方案2:主进程侧统计(更推荐,无共享变量开销)
如果不需要在子进程的任务函数内部获取剩余任务量,只需要全局统计进度,可以将统计逻辑完全放到主进程执行,用imap替代map迭代任务结果,主进程单线程统计和打印,完全不会出现乱序问题,也不需要任何共享变量或锁:
import multiprocessing def f(x): # 仅保留业务逻辑,无任何统计/打印相关代码 return x * 2 # 替换为实际业务逻辑 if __name__ == '__main__': rangeval = range(1000) total = len(rangeval) pool = multiprocessing.Pool(processes=multiprocessing.cpu_count()) completed = 0 # 用imap迭代返回结果,每完成一个任务就会返回一次 for result in pool.imap(f, iterable=rangeval): completed +=1 remaining = total - completed print(f"Done with task: {completed}, remaining tasks: {remaining}") pool.close() pool.join()
这个方案不会对任务执行的并行效率产生任何影响,是性能最优的实现方式。
内容的提问来源于stack exchange,提问作者Vance Pyton
相关产品推荐
相关产品推荐

