如何将含内部逻辑的嵌套循环转为多进程并行迭代对象
问题
我已经掌握了这种多进程并行的写法:
with PoolExec(max_workers=int(config['MULTIPROCESSING']['proceses_count']),initializer=self.initPool,initargs=(arg0,arg1,arg2,arg3,arg4,arg5,arg6,arg7,arg8,arg9,)) as GridCells10mX10mIteratorPool.__poolExec: self.__chunkSize = PoolUtils.getChunkSizeForLenOfIterables(lenOfIterablesList=self.__maxNumOfCellsVertically*self.__maxNumOfCellsHorizontally,cpuCount=int(config['MULTIPROCESSING']['cpu_count'])) for res in GridCells10mX10mIteratorPool.__poolExec.map(self.run,[(i,j) for i in range(0,1800,10) for j in range(0,2000,10)] ,chunksize=self.__chunkSize): # 处理结果
但现在有一段嵌套循环,外层和内层循环后各有两行状态判断逻辑,不知道怎么转换成上面的多进程写法:
for x in range(row,row + gVerticalStep): if rowsCnt == gVerticalStep: rowsCnt = 0 for y in range(col,col + gHorizontalStep): if colsCnt == gHorizontalStep: colsCnt = 0 # 针对(x,y)的核心处理逻辑
解决方案
首先明确:多进程环境下,每个子进程有独立内存空间,rowsCnt、colsCnt这类状态变量没法在进程间直接共享。所以得把状态逻辑调整为每个任务独立处理,或者重构逻辑,让状态不需要跨任务传递。
1. 生成所有待处理的(x,y)迭代器
先把嵌套循环里的所有(x,y)对用列表推导式生成,和你之前的写法一致:
# 生成所有需要处理的(x,y)元组 task_iter = [(x, y) for x in range(row, row + gVerticalStep) for y in range(col, col + gHorizontalStep)]
2. 重构run函数,封装状态逻辑
把原来循环内的rowsCnt、colsCnt判断逻辑移到run函数里,分两种情况处理:
情况1:状态基于x/y的循环位置
如果rowsCnt是记录当前x在循环中的索引(比如每循环gVerticalStep次重置),完全可以通过计算得到,不用维护变量:
def run(self, args): x, y = args # 计算当前x对应的循环索引,取模gVerticalStep得到rowsCnt rowsCnt = (x - row) % gVerticalStep if rowsCnt == gVerticalStep - 1: # 最后一次循环时触发重置逻辑 # 写rowsCnt重置后的操作,比如初始化局部变量 pass # 同理计算colsCnt colsCnt = (y - col) % gHorizontalStep if colsCnt == gHorizontalStep - 1: # colsCnt重置后的操作 pass # 执行(x,y)对应的核心业务逻辑 # ... return result
情况2:状态是跨任务的累积值(不推荐)
如果rowsCnt/colsCnt是需要跨多个(x,y)任务的累积状态,map就不适用了——因为map是无状态的任务分发。这种情况最好用Queue手动管理任务,或者改用线程(但线程受GIL限制,CPU密集型任务别用)。
3. 替换为多进程map写法
把生成的task_iter传入map,和你之前的写法对齐:
with PoolExec(max_workers=int(config['MULTIPROCESSING']['proceses_count']), initializer=self.initPool, initargs=(arg0, arg1, arg2, arg3, arg4, arg5, arg6, arg7, arg8, arg9,)) as GridCells10mX10mIteratorPool.__poolExec: total_tasks = gVerticalStep * gHorizontalStep self.__chunkSize = PoolUtils.getChunkSizeForLenOfIterables( lenOfIterablesList=total_tasks, cpuCount=int(config['MULTIPROCESSING']['cpu_count']) ) for res in GridCells10mX10mIteratorPool.__poolExec.map(self.run, task_iter, chunksize=self.__chunkSize): # 处理每个任务的结果 # ...
关键提醒
- 多进程里别依赖全局共享状态变量,每个进程内存隔离,子进程修改的变量不会同步到主进程或其他子进程。
- 如果原来的
rowsCnt/colsCnt只是用来做循环内的初始化,优先通过计算x/y位置推导状态,别搞共享状态,成本太高。
内容的提问来源于stack exchange,提问作者Amrmsmb
相关产品推荐
相关产品推荐

