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
相关产品推荐
相关产品推荐

