Ubuntu下Python动态向进程池添加TSP任务至收敛的技术问询
旅行商问题并行采样优化方案
背景
在Ubuntu环境下使用Python解决旅行商(Travelling Salesman, TSP)路由问题:因每日客户集合会变化,需对网络进行随机采样并求解每个样本,通过足够样本估算期望平均成本,直至样本平均收敛(由CONVERGED()函数判断,返回True表示已收敛)。
现有实现问题
为加速计算,需并行执行采样任务,且仅在未收敛时生成并求解新样本,先后尝试两种方案:
- 版本v1:使用
multiprocessing.Pool,代码如下:
manager = multiprocessing.Manager() q = manager.Queue() pool = multiprocessing.Pool(multiprocessing.cpu_count() + 2) while not <CONVERGED()>: job = pool.apply_async(<FUNCTION TO CALCULATE OUTPUT>, <ARGUMENTS>)) job.get()
调用job.get()会等待当前任务完成才进入下一轮循环,无法充分利用多核资源。
- 版本v2:改为批量提交100个任务,代码如下:
manager = multiprocessing.Manager() q = manager.Queue() pool = multiprocessing.Pool(multiprocessing.cpu_count() + 2) while not <CONVERGED()>: jobs = [] for i in range(100): jobs.append(pool.apply_async(<FUNCTION TO CALCULATE OUTPUT>, <ARGUMENTS>)) for job in jobs: job.get()
但任务耗时差异极大,需等待批量中最慢的任务完成才提交下一批,效率受限。
需求
能否修改实现逻辑,无需设定固定批量大小,让程序在核心空闲且未收敛时自动添加新的采样任务?
补充细节及疑问
- 任务特性:每个任务输入为随机种子,无返回值,仅将结果写入文件,
CONVERGED()通过遍历文件中的结果判断收敛状态。 - 实现疑问:尝试用全局变量
_CONVERGED跟踪收敛状态,但子进程修改后主进程无法感知;改用回调函数则能正常运行,需了解:- 回调函数相比在任务函数内判断收敛的优势?
- 全局变量无法更新的原因?
相关代码如下:
def MC(args1, args2, args3): manager = multiprocessing.Manager() q = manager.Queue() pool = BoundedQueueProcessPool(multiprocessing.cpu_count() + 2) # 启动监听进程 pool.apply_async(listener, (q, csv_name, message, data, csv_col)) # 用于按需读写文件的进程 count = 0 global _CONVERGED while not _CONVERGED and count <= iterlimit: print(count, '-', _CONVERGED) count += 1 pool.apply_async(mc_call, (data, target, q, count, csv_name)) pool.close() q.put([-2]) # 终止监听进程 pool.join() return None def mc_call(data: MDVRPI, target: str, q: int, seed:float, csv_name:str): global _CONVERGED # TODO : 下面这段可以改成回调函数,有什么好处吗? if magic_function_to_calc_if_converged(): _CONVERGED = True return None gen_data = generate_data(...) run_output = run_generated_data(...) q.put([run_output]) # 将结果传给监听进程写入文件 return None if __name__ == '__main__': _CONVERGED = False # 收敛标志 MC(Args1, Args2, Args3) return 0
- Listener进程:通过队列接收任务结果并写入CSV文件,实时刷新,代码如下:
def listener (q, csv_name, message, data, csv_col): ''' 监听每个任务的结果并写入文件。-1表示收敛,-2表示运行结束(关闭文件) ''' file = open(csv_name, "w") file.write(platform.node() + ', ,'+ message +'\n') file.write(','.join(csv_col) + '\n') file.write(','.join(['0' for i in range(__main__.csv_col)]) + '\n') file.flush() while True: m = q.get() if m['status'] == -1: # print('已收敛') pass elif m['status'] == -2: # print('运行结束') break file.write(','.join([str(i) for i in [m['seed'],m['cost'], m['time'], m['solverstatus'], *m['demands']]]) + '\n') file.flush() file.close() return
内容的提问来源于stack exchange,提问作者Abilash
相关产品推荐
相关产品推荐

