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

基于新进程的多进程处理:固定并行数下单次进程使用方案咨询

实现固定并发数的单任务独立进程方案

核心思路

使用multiprocessing.Semaphore控制同时运行的进程数量,每次启动新进程前获取信号量,进程执行完成后释放信号量,确保同一时刻最多有指定数量的进程在运行;每个任务对应一个全新的Process实例,任务结束后进程自动销毁。

代码实现示例

import multiprocessing
import time

def compute_task(task_id):
    # 替换为你的实际计算逻辑
    print(f"任务{task_id}启动,进程ID: {multiprocessing.current_process().pid}")
    time.sleep(2)  # 模拟计算耗时
    print(f"任务{task_id}完成")

if __name__ == "__main__":
    MAX_CONCURRENT = 3  # 固定并行进程数
    TOTAL_TASKS = 30     # 总任务量
    semaphore = multiprocessing.Semaphore(MAX_CONCURRENT)
    process_list = []

    for task_id in range(TOTAL_TASKS):
        # 阻塞直到有可用的并发名额
        semaphore.acquire()
        
        # 包装任务函数,确保进程结束后释放信号量
        def wrapped_task(tid):
            try:
                compute_task(tid)
            finally:
                semaphore.release()
        
        # 为每个任务启动新进程
        proc = multiprocessing.Process(target=wrapped_task, args=(task_id,))
        process_list.append(proc)
        proc.start()

    # 等待所有进程执行完毕
    for proc in process_list:
        proc.join()

关键细节说明

  • 信号量控制并发:Semaphore的初始值设为最大并行数,acquire()会占用一个名额,release()会释放名额,确保不会超过设定的并发上限。
  • 独立进程保证:每个任务都创建新的Process对象,进程执行完任务后自动退出销毁,完全满足“每个输入任务使用全新进程”的要求。
  • 避免死锁:用try...finally包裹任务逻辑,无论任务正常完成还是抛出异常,都会释放信号量,防止因异常导致信号量无法释放的死锁问题。

额外建议

  • 若任务需要传递大量数据,优先使用multiprocessing.Queue或管道进行进程间通信,减少直接传参带来的内存拷贝开销。
  • 在compute_task中添加异常捕获逻辑,避免单个任务崩溃影响整个程序的运行。

内容的提问来源于stack exchange,提问作者Simon P.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 14:25:39