Python多进程中如何选择性终止进程及利用空闲核心?
并行计算任务填充与进程池管理问题
我使用Python的multiprocessing模块执行密集型并行计算,将计算逻辑抽象为runcode()中的sleep函数,这类任务的耗时难以预估。启动第一批主任务后,经常出现部分主任务已完成、仅剩余少量慢任务的情况。为避免空等,我希望提交次级填充任务来利用这段时间,但主任务全部完成后,无论次级任务是否完成都要终止它们(即调用pool.terminate())。由于存在处理读写的监听函数需始终保持活跃,因此不能盲目终止整个进程池。
我想到两个方案:
- 方案1:将次级任务加入同一进程池,通过标记等方式选择性终止,但不确定具体实现方法。
- 方案2:创建第二个进程池来运行次级任务,主任务完成后直接终止该进程池(代码中带
###的部分为相关实现思路)。但我曾看到过不建议使用多进程池的说法,却找不到相关依据。
此外,两个方案都存在一个共同问题:若盲目提交大量次级任务,可能主任务完成前次级任务就已结束。理想状态是仅在主任务未完成且有空闲核心时(即下文的ThereAreFreeCores())才提交次级任务,请问是否有实现方法?
import multiprocessing import multiprocessing.pool from contextlib import ExitStack import time import random class BoundedQueuePool: def __init__(self, limit, semaphore_type): self._semaphore = semaphore_type(limit) def release(self, result, callback=None): self._semaphore.release() if callback: callback(result) def apply_async(self, func, args=(), kwds={}, callback=None, error_callback=None): self._semaphore.acquire() callback_fn = self.release if callback is None else lambda result: self.release(result, callback=callback) error_callback_fn = self.release if error_callback is None else lambda result: self.release(result, callback=callback) return super().apply_async(func, args, kwds, callback=callback_fn, error_callback=error_callback_fn) class BoundedQueueProcessPool(BoundedQueuePool, multiprocessing.pool.Pool): def __init__(self, *args, max_waiting_tasks=None, **kwargs): multiprocessing.pool.Pool.__init__(self, *args, **kwargs) if max_waiting_tasks is None: max_waiting_tasks = self._processes elif max_waiting_tasks < 0: raise ValueError(f'Invalid negative max_waiting_tasks value: {max_waiting_tasks}') limit = self._processes + max_waiting_tasks BoundedQueuePool.__init__(self, limit, multiprocessing.BoundedSemaphore) def listener (q, csv_names): ''' Listens for results from each run and prints it to file. ''' with ExitStack() as stack: files = {key: stack.enter_context(open(csv_names[key], "w")) for key in csv_names} for key in csv_names: files[key].write('Value\n') files[key].flush() while True: m = q.get() if m == 'STOP': break ff = m[0] files[ff].write(m[1]) files[ff].flush() return def runcode(q): x = random.random() time.sleep(x*10) q.put([1, str(x)+'\n']) return iterlimit = 40 csv_names = {1:"./file1.csv"} count = 0 manager = multiprocessing.Manager() q = manager.Queue() pool = BoundedQueueProcessPool(multiprocessing.cpu_count() + 2) pool2 = BoundedQueueProcessPool(multiprocessing.cpu_count() + 2) ### pool.apply_async(listener, (q,csv_names)) # Activating Listener - Handles read and write while count <= iterlimit: # Lets call these runs as master runs count += 1 pool.apply_async(runcode, (q,)) while len(pool._cache) > 1: # Waiting for runs to finish # New chunk I want to add # if ThereAreFreeCores(): ### # pool2.apply_async(runcode, (q,)) ### Lets call these runs secondary runs continue # Now ALL master runs are finished but secondary runs MIGHT not be done pool.close() pool2.terminate() ### q.put('STOP') pool.join()
内容的提问来源于stack exchange,提问作者Abilash
相关产品推荐
相关产品推荐

