Python multiprocessing apply_async未并行运行问题求助
解决方案
核心原因
主进程提交完任务后立刻进入CPU密集型的训练循环,抢占了大部分CPU资源,导致进程池的子进程无法被操作系统调度执行,只能等主进程当前训练步骤完成后才能运行,最终表现为串行。另外,每次epoch重复创建Pool和Manager队列的额外开销也会影响并行效率。
具体修改步骤
- 将进程池和队列的创建移到epoch循环外,避免每次epoch重复初始化,减少资源开销。
- 保存异步任务结果,让子进程与主进程训练并行:提交任务后不立即阻塞,让子进程后台预处理数据,主进程同时从队列取数据训练,最后统一等待所有任务完成。
- epoch开始前清空队列,避免上一轮残留数据干扰当前训练。
修改后的代码示例
from multiprocessing import Manager, Pool import os import time # 提前初始化进程池和队列,复用资源 ma1 = Manager() q = ma1.Queue(maxsize=40) p = Pool(32) for epoch in range(num_epochs): timer = d2l.Timer() """write to queue""" # 清空队列,确保当前epoch用干净的队列 while not q.empty(): try: q.get(block=False) except: pass # 提交所有预处理任务,保存异步结果 async_results = [] for i in range(giveashuffleindx.nb_updates_per_epoch): res = p.apply_async( write_to_queue, args=(trainset, trainset_volume_manager, batchsize, max_length, i, giveashuffleindx.shuffledata_indices, q, False), error_callback=print_error ) async_results.append(res) print('Subprocesses started, beginning training...') metric = d2l.Accumulator(2) """get data from queue and train""" t = 0 while t < giveashuffleindx.nb_updates_per_epoch: batch = q.get(True) print('count {} from {}'.format(t, q)) optimizer.zero_grad() X, Y, Mask, batch_len = [x.to(devices) for x in batch] X = X.to(torch.float32) state = net.init_state() Y_hat, _ = net(X, state) B, T, C = Y_hat.shape Y_hat = Y_hat.view(B*T, C) Y = Y.view(B*T, C) l = F.cross_entropy(Y_hat, Y) l.backward() d2l.grad_clipping(net, 1) optimizer.step() with torch.no_grad(): metric.add(l, X.size(0)) t += 1 # 等待当前epoch的所有预处理任务完成,避免影响下一轮 for res in async_results: res.get() print(f"Epoch {epoch+1} done, loss: {metric[0]/metric[1]}") # 所有训练完成后,关闭进程池 p.close() p.join()
额外注意事项
- 如果
prepare_batch中存在大量纯Python循环(绑定GIL的操作),可以尝试用concurrent.futures.ProcessPoolExecutor替代multiprocessing.Pool,但多数场景下Pool足够满足需求。 - 检查机器CPU使用率,若被其他任务占用过多,并行效果会打折扣,确保有足够空闲核心留给子进程。
内容的提问来源于stack exchange,提问作者yi yang
相关产品推荐
相关产品推荐

