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

Python多进程问题:如何实现8进程并发且任务完成后补启新任务

解决Python多进程固定并发数的问题

你现在的问题在于直接创建了100个进程,之后如果一次性启动它们,系统就会同时跑100个任务,完全没控制住并发数。而且还有个隐藏坑:你用的全局bits数组在多进程里是各自独立的拷贝,子进程里修改它,主进程根本看不到变化!

下面给你两种靠谱的修改方案,都能实现“同时跑8个,一个完成就补新任务”的需求:

方案一:用multiprocessing.Pool(最省心)

Pool本身就帮你管理了进程池的大小,指定processes=8就能固定同时运行8个任务。我们可以让子进程返回结果,由主进程来更新数组,完美避开共享变量的坑:

import multiprocessing
import numpy as np

def some_function(idd):
    # 这里替换成你的实际业务逻辑
    return idd * 2

def worker(idd):
    # 子进程只负责计算,返回任务ID和对应结果
    return idd, some_function(idd)

if __name__ == "__main__":
    bits = np.zeros(100)
    # 创建大小为8的进程池
    with multiprocessing.Pool(processes=8) as pool:
        # imap_unordered会在任务完成时立刻返回结果,实时更新数组
        for idd, result in pool.imap_unordered(worker, range(100)):
            bits[idd] = result
    # 此时bits已经存储了所有计算结果
    print(bits)

如果一定要用共享数组的方式,也可以用multiprocessing.Array实现跨进程共享:

import multiprocessing
import numpy as np

def some_function(idd):
    return idd * 2

def worker(shared_bits, idd):
    # 将共享内存数组转为numpy数组操作
    bits_np = np.frombuffer(shared_bits.get_obj(), dtype=np.float64)
    bits_np[idd] = some_function(idd)

if __name__ == "__main__":
    # 创建共享内存的float64类型数组,长度100
    shared_bits = multiprocessing.Array('d', 100)
    with multiprocessing.Pool(processes=8) as pool:
        # 用starmap传递多参数
        pool.starmap(worker, [(shared_bits, i) for i in range(100)])
    # 转换为numpy数组查看结果
    bits = np.frombuffer(shared_bits.get_obj(), dtype=np.float64)
    print(bits)

方案二:手动管理进程(更灵活)

如果你需要更精细的进程控制,不想用Pool,可以手动维护“运行中进程”列表,完成一个就启动新的:

import multiprocessing
import numpy as np

def some_function(idd):
    return idd * 2

def worker(result_queue, idd):
    # 将结果放入队列,供主进程读取
    result_queue.put((idd, some_function(idd)))

if __name__ == "__main__":
    bits = np.zeros(100)
    max_workers = 8
    result_queue = multiprocessing.Queue()
    jobs = []

    # 先启动第一批8个任务
    for i in range(max_workers):
        p = multiprocessing.Process(target=worker, args=(result_queue, i))
        p.start()
        jobs.append((p, i))
    
    # 处理剩余92个任务
    for i in range(max_workers, 100):
        # 等待一个进程完成
        done_idx = None
        for idx, (p, pid) in enumerate(jobs):
            if not p.is_alive():
                done_idx = idx
                p.join()
                break
        # 移除已完成的进程,处理结果
        del jobs[done_idx]
        idd, result = result_queue.get()
        bits[idd] = result
        # 启动新任务
        new_p = multiprocessing.Process(target=worker, args=(result_queue, i))
        new_p.start()
        jobs.append((new_p, i))
    
    # 处理最后一批剩余的进程
    for p, pid in jobs:
        p.join()
        idd, result = result_queue.get()
        bits[idd] = result
    
    print(bits)

最后提醒一句:一定要把主逻辑放到if __name__ == "__main__":代码块下,这是Windows系统多进程的强制要求,能避免无限创建子进程的问题。

内容的提问来源于stack exchange,提问作者A.Razavi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:52:52