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

Python多进程异步任务阻塞及进度保存问题咨询

多进程任务处理相关问题

一、进程池Pool相关问题

假设进程池有p0-p3四个进程,共10个任务task(0)至task(9),其中p0执行task(0)耗时极久,针对以下代码:

from multiprocessing import Pool

def task(s):
    # 执行一些操作
    return res

pool = Pool(4)
status = []
res = []
for i in data :
    status.append(pool.apply_async(task, (i,)))

for i in status :
    res.append(i.get())
    # 每30秒用pickle保存res

问题1:主进程会在第一个res.append(i.get())处阻塞吗?

是的。i.get()会强制阻塞主进程,直到对应的异步任务完成并返回结果。第一个i对应task(0),由于p0执行该任务耗时极久,主进程会一直卡在这行代码,直到task(0)执行完毕。

问题2:若p1完成task(1)而p0仍在处理task(0),p1会继续处理task(4)等后续任务吗?

会。进程池的工作逻辑是:当某个进程完成当前任务后,会自动从任务队列中取下一个未执行的任务。p1完成task(1)后,会立即接手task(4)(初始分配的是task(0)-task(3)),不会因p0未完成而闲置。

问题3:若第一个问题答案为是,如何提前获取其他任务结果,最后再获取task(0)的结果?

可以将task(0)的异步任务对象单独拆分,先处理其他任务的结果,最后再等待task(0)。或者改用concurrent.futures的as_completed方法,按任务完成顺序获取结果,无需等待慢任务。

示例代码(改用ProcessPoolExecutor):

from concurrent.futures import ProcessPoolExecutor, as_completed

with ProcessPoolExecutor(4) as ex:
    futures = [ex.submit(task, i) for i in data]
    res = []
    for future in as_completed(futures):
        res.append(future.result())
        # 每30秒保存res的逻辑

如果坚持用multiprocessing.Pool,可以手动分离慢任务:

from multiprocessing import Pool

pool = Pool(4)
status = []
for i in data :
    status.append(pool.apply_async(task, (i,)))

# 分离task(0)的任务对象
task0_future = status[0]
other_futures = status[1:]

# 先处理其他已完成的任务
res = []
for future in other_futures:
    if future.ready():
        res.append(future.get())

# 最后等待task(0)的结果
res.append(task0_future.get())

二、ProcessPoolExecutor阻塞问题排查

改用以下代码后,出现主进程阻塞但子进程仍在运行的问题:

import concurrent.futures
import datetime
import os
import shutil
import pickle

# 类中方法示例
def some_method(self):
    futuresList = []
    with concurrent.futures.ProcessPoolExecutor(4) as ex :
        for i in self.inBuffer :
            futuresList.append(ex.submit(wrapper, i))
        
        for i in concurrent.futures.as_completed(futuresList) :
            print("getting res of something")
            (word, r) = i.result()
            print("finishing i.result")
            self.resDict[word] = r
            print("finished getting res of {}".format(word))
            self.logger.info("{} --> {}".format(word, r))
            cur = datetime.now()
            if (cur - self.timeStmp).total_seconds() > 30 :
                self.outputPickle()
                self.timeStmp = datetime.now()

def outputPickle(self):
    if os.path.exists(os.path.join(self.wordDir, self.outFile)) :
        if os.path.exists(os.path.join(self.wordDir, "{}_backup".format(self.outFile))):
            os.remove(os.path.join(self.wordDir, "{}_backup".format(self.outFile)))
        shutil.copy(os.path.join(self.wordDir, self.outFile), os.path.join(self.wordDir, "{}_backup".format(self.outFile)))
    
    with open(os.path.join(self.wordDir, self.outFile), 'wb') as f:
        pickle.dump(self.resDict, f)

现象

初始运行正常,某时刻起日志文件长时间未更新(远超单个wrapper的120秒耗时),但wrapper仍在输出"message by wrapper"(每个wrapper最多输出一次),且.pkl文件时间戳未变化。添加print语句后得到日志:

getting res of something
finishing i.result
finished getting res of CNICnanotubesmolten
getting res of something
finishing i.result
finished getting res of CNN0
getting res of something
message by wrapper
message by wrapper
message by wrapper
message by wrapper
message by wrapper

阻塞原因分析

从日志来看,主进程卡在了i.result()调用前(已打印"getting res of something",但未打印"finishing i.result"),结合子进程仍在输出信息,核心原因是子进程的结果无法被主进程正常接收或反序列化,具体可能性如下:

  1. 结果序列化失败:wrapper返回的(word, r)中,r包含无法被pickle序列化的对象(如未实现__reduce__的自定义类实例、文件句柄、网络连接等)。子进程完成任务后,序列化结果时出错,导致主进程的i.result()一直阻塞等待。
  2. 子进程异常未传递:wrapper执行中抛出未捕获的异常,但由于进程间通信机制限制,异常无法正确传递到主进程,导致i.result()一直处于阻塞状态,而子进程的print输出是在异常抛出前执行的。
  3. 文件IO隐性阻塞:outputPickle中的文件复制、删除操作可能因文件被其他进程占用而阻塞,但从日志看主进程卡在i.result(),该可能性较低。

排查建议

  • 给i.result()添加异常捕获,确认是否有异常:
    try:
        (word, r) = i.result()
    except Exception as e:
        print(f"Result error: {e}")
        self.logger.error(f"Task failed: {e}")
        continue
    
  • 检查wrapper返回的r是否可被pickle序列化:在wrapper中尝试pickle.dumps(r),查看是否抛出异常。
  • 给wrapper添加全局异常捕获:
    def wrapper(i):
        try:
            # 原有逻辑
            return (word, r)
        except Exception as e:
            print(f"Wrapper error: {e}")
            return (i, None)
    
  • 查看系统进程状态,确认子进程是否为僵尸进程。

内容的提问来源于stack exchange,提问作者user19631495

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 19:10:20