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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 18:23:17