Python多进程:高效维持最大并发进程数及嵌套函数问题
首先,咱们来解决你最核心的需求:保持最多N个进程并行运行,一个进程结束后立刻启动新任务。你之前的轮询方案确实效率不高,还会占用额外的CPU资源,Python标准库其实提供了更优雅、高效的工具来搞定这个事儿。
首选方案:用concurrent.futures.ProcessPoolExecutor
这是Python 3.2+引入的高级API,专门用来管理进程/线程池,它会自动帮你调度任务——当一个进程完成当前任务,马上就会从任务队列里取下一个任务执行,完全不需要你手动轮询进程状态。
举个例子,直接替换你的逻辑:
import concurrent.futures import time def processf(num): print(f"in process: {num}") now = time.time() while time.time() - now < 2: pass # 模拟密集计算任务 if __name__ == "__main__": max_parallel_processes = 4 # 最多同时跑4个进程 total_tasks = 20 # 总共有20个任务要执行 # 创建进程池,上下文管理器会自动帮你关闭/清理进程 with concurrent.futures.ProcessPoolExecutor(max_workers=max_parallel_processes) as executor: # 把所有任务交给executor调度,它会自动维持进程数 executor.map(processf, range(total_tasks))
这个代码逻辑清晰,而且完全实现了你的需求:始终保持最多4个进程在运行,一个结束就补新的。如果需要处理任务返回值或者异常,还可以用executor.submit配合回调函数:
import concurrent.futures import time def processf(num): print(f"in process: {num}") now = time.time() while time.time() - now < 2: pass return f"Process {num} done!" def handle_result(future): # 任务完成时的回调函数 print(future.result()) if __name__ == "__main__": max_parallel_processes = 4 total_tasks = 20 with concurrent.futures.ProcessPoolExecutor(max_workers=max_parallel_processes) as executor: for num in range(total_tasks): future = executor.submit(processf, num) future.add_done_callback(handle_result)
备选方案:用multiprocessing.Pool
如果你习惯用multiprocessing模块,或者需要兼容Python 2,Pool的apply_async方法也能实现类似效果,通过回调函数来触发新任务:
import multiprocessing import time def processf(num): print(f"in process: {num}") now = time.time() while time.time() - now < 2: pass return num def on_task_done(result): print(f"Process {result} finished") global task_counter task_counter += 1 if task_counter < total_tasks: # 任务没做完,继续提交新任务 pool.apply_async(processf, args=(task_counter,), callback=on_task_done) if __name__ == "__main__": max_parallel_processes = 4 total_tasks = 20 task_counter = max_parallel_processes # 先启动max个任务 pool = multiprocessing.Pool(max_workers=max_parallel_processes) # 初始化启动第一批任务 for i in range(max_parallel_processes): pool.apply_async(processf, args=(i,), callback=on_task_done) pool.close() pool.join()
不过这种方式需要手动维护任务计数器,不如ProcessPoolExecutor简洁,所以更推荐前者。
你提到把processf嵌入main函数里会报错AttributeError: can't pickle local object 'main.<locals>.processf',这是因为Python的multiprocessing默认用pickle来序列化函数,但pickle没办法序列化嵌套在函数里的本地函数——子进程找不到这个函数的模块级别定义,自然无法加载它。
最简单的解决方法:把函数移到模块级别
直接把processf从main里拿出来,作为模块的顶层函数,这样pickle就能正常序列化它了,同时你也不需要用列表z这种hack方式来维护计数器,进程池会帮你搞定任务调度。
复杂场景的 workaround:用cloudpickle
如果你确实需要把函数嵌套在main里,可以用第三方库cloudpickle(需要先安装pip install cloudpickle),它支持序列化本地函数:
import multiprocessing import time import cloudpickle def main(): def processf(num): print(f"in process: {num}") now = time.time() while time.time() - now < 2: pass # 修改multiprocessing的序列化器为cloudpickle multiprocessing.set_start_method('spawn') # Windows系统必须用spawn启动方式 multiprocessing.reduction.ForkingPickler = cloudpickle.CloudPickler max_parallel_processes = 4 total_tasks = 20 with multiprocessing.Pool(max_workers=max_parallel_processes) as pool: pool.map(processf, range(total_tasks)) if __name__ == "__main__": main()
不过这种方法依赖第三方库,而且跨平台兼容性需要测试,所以还是推荐第一种方案。
关于你用线程套进程的方案
你的第二个方案确实避免了轮询的问题,但其实完全没必要用线程来启动进程——进程池本身已经帮你做好了进程调度的工作,线程在这里反而增加了不必要的复杂度。而且用列表z维护计数器的方式也不够优雅,用进程池的话根本不需要手动管理这个状态。
内容的提问来源于stack exchange,提问作者user7638202

