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

Python multiprocessing apply_async未并行运行问题求助

解决方案

核心原因

主进程提交完任务后立刻进入CPU密集型的训练循环,抢占了大部分CPU资源,导致进程池的子进程无法被操作系统调度执行,只能等主进程当前训练步骤完成后才能运行,最终表现为串行。另外,每次epoch重复创建Pool和Manager队列的额外开销也会影响并行效率。

具体修改步骤

  1. 将进程池和队列的创建移到epoch循环外,避免每次epoch重复初始化,减少资源开销。
  2. 保存异步任务结果,让子进程与主进程训练并行:提交任务后不立即阻塞,让子进程后台预处理数据,主进程同时从队列取数据训练,最后统一等待所有任务完成。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 23:17:20