You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Ubuntu下Python动态向进程池添加TSP任务至收敛的技术问询

旅行商问题并行采样优化方案

背景

在Ubuntu环境下使用Python解决旅行商(Travelling Salesman, TSP)路由问题:因每日客户集合会变化,需对网络进行随机采样并求解每个样本,通过足够样本估算期望平均成本,直至样本平均收敛(由CONVERGED()函数判断,返回True表示已收敛)。

现有实现问题

为加速计算,需并行执行采样任务,且仅在未收敛时生成并求解新样本,先后尝试两种方案:

  1. 版本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()会等待当前任务完成才进入下一轮循环,无法充分利用多核资源。

  1. 版本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()

但任务耗时差异极大,需等待批量中最慢的任务完成才提交下一批,效率受限。

需求

能否修改实现逻辑,无需设定固定批量大小,让程序在核心空闲且未收敛时自动添加新的采样任务?

补充细节及疑问

  1. 任务特性:每个任务输入为随机种子,无返回值,仅将结果写入文件,CONVERGED()通过遍历文件中的结果判断收敛状态。
  2. 实现疑问:尝试用全局变量_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
  1. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 20:10:27