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

Python多进程任务pruning后CPU使用率骤降的问题与优化咨询

多进程处理大规模数据时CPU使用率骤降的问题与优化方案

问题背景

我开发的应用需要处理数百万条文件数据点,流程如下:

  • 从文件读取数据 → 转换为数据结构A
  • 通过toString()去重 → 转换为更节省内存的数据结构B(采用延迟加载)
  • 为提升处理速度,使用两个multiprocessing.Pool分别并行生成数据结构A和B

最初用imap_unordered延迟生成A,再通过apply_async提交A转B的任务,但处理约1万条数据后出现内存错误(推测是AsyncResult堆积持有B的结果导致)。于是实现了定期清理逻辑:从AsyncResult中取出结果并移除引用,但清理后CPU使用率直接从满核降至20%左右且无法恢复。

改用回调函数简化代码后,CPU使用率仍然骤降;换成ThreadPool后问题有所缓解(线程共享内存避免了进程间的数据复制)。需要解决CPU使用率下降的问题,或寻找更优实现方式。

初始版本代码

import multiprocessing as mp

p = mp.Pool(4) # 处理数据结构A的生成
q = mp.Pool(32) # 处理数据结构B的生成
n = 0 # 触发清理逻辑的计数器
structureBs = [] # 最终输出结构列表
Aresults = [] # AsyncResult对象列表
Astrings = set() # 用于去重的字符串集合

for structureA in p.imap_unordered(generateStructA, readFile(inputFile)):
    Astring = structureA.toString()
    if Astring not in Astrings:
        Astrings.add(Astring)
        Aresults.append(q.apply_async(generateStructB, (structureA, )))
        n += 1
        if n % 1000 == 0:
            for res in Aresults:
                if res.ready():
                    structureB = res.get()
                    structureBs.append(structureB)
                    Aresults.remove(res)

# 处理未完成的B生成任务
for res in Aresults:
    structureB = res.get()
    structureBs.append(structureB)

回调函数版本代码

import multiprocessing as mp

p = mp.Pool(4) # 处理数据结构A的生成
q = mp.Pool(32) # 处理数据结构B的生成
structureBs = [] # 最终输出结构列表
Astrings = set() # 用于去重的字符串集合

def callback(result):
    structureBs.append(result)

for structureA in p.imap_unordered(generateStructA, readFile(inputFile)):
    Astring = structureA.toString()
    if Astring not in Astrings:
        Astrings.add(Astring)
        q.apply_async(generateStructB, (structureA, ), callback=callback) 

核心问题分析

  1. 进程间数据传递开销过大
    多进程模式下,structureA在进程间传递需要经过序列化(默认用pickle)和反序列化。如果structureA体积较大,这个过程会占用大量CPU和IO资源,导致父进程在数据传递上耗时过多,子进程陷入等待,CPU利用率下降。
  2. 双Pool调度不匹配
    生成A的Pool仅4个进程,生成B的Pool有32个,当A的生成速度跟不上B的处理速度时,B的Pool会出现大量空闲进程,整体CPU利用率自然下降。
  3. 清理/回调逻辑干扰调度
    初始版本的清理循环中,父进程频繁遍历Aresults检查任务状态,打乱了Pool的正常调度节奏;回调函数版本中,回调在父进程执行,当B的生成速度远快于父进程处理回调的速度时,回调队列堆积,父进程被占满无法及时提交新任务,导致子进程空闲。

优化方案

1. 匹配Pool进程数,平衡任务节奏

先单独测试generateStructA和generateStructB的单任务耗时,根据耗时比例调整两个Pool的进程数。比如如果生成A的耗时是生成B的8倍,4:32的比例是合理的;若比例不符,就调整进程数,确保A的生成速度能喂饱B的Pool,避免B的进程空闲。

2. 用队列实现流式处理,减少AsyncResult堆积

用multiprocessing.Queue替代apply_async的回调,让B的Pool进程直接从队列取A的数据,避免父进程的中间调度开销:

import multiprocessing as mp

def produce_A(queue, input_file):
    astrings = set()
    for data in readFile(input_file):
        structA = generateStructA(data)
        astr = structA.toString()
        if astr not in astrings:
            astrings.add(astr)
            queue.put(structA)
    # 放入结束标记
    queue.put(None)

def consume_B(queue, output_list):
    while True:
        structA = queue.get()
        if structA is None:
            # 传递结束标记给下一个进程
            queue.put(None)
            break
        structB = generateStructB(structA)
        output_list.append(structB)

if __name__ == "__main__":
    manager = mp.Manager()
    queue = manager.Queue(maxsize=100)  # 设置队列大小,避免内存爆炸
    structureBs = manager.list()

    # 启动A的生成进程池
    a_pool = mp.Pool(4)
    a_pool.apply_async(produce_A, args=(queue, inputFile))
    a_pool.close()

    # 启动B的消费进程池
    b_pool = mp.Pool(32)
    for _ in range(32):
        b_pool.apply_async(consume_B, args=(queue, structureBs))
    b_pool.close()

    a_pool.join()
    b_pool.join()

    # 转换为普通列表
    structureBs = list(structureBs)

队列设置maxsize可以避免A生成过快导致内存占用过高,同时让A的进程在队列满时阻塞,自然平衡生产和消费速度。

3. 合并处理步骤,减少进程间数据传递

如果generateStructA和去重的开销不大,可以把去重逻辑放到B的生成流程中,或者直接在A的生成进程中完成A转B的操作,减少一次进程间的数据传递:

def process_data(data):
    structA = generateStructA(data)
    astr = structA.toString()
    return astr, generateStructB(structA)

if __name__ == "__main__":
    astrings = set()
    structureBs = []
    with mp.Pool(32) as pool:
        for astr, structB in pool.imap_unordered(process_data, readFile(inputFile)):
            if astr not in astrings:
                astrings.add(astr)
                structureBs.append(structB)

这样只需要传递一次处理后的数据,减少IPC开销,同时用一个Pool统一调度,避免两个Pool的调度冲突。

4. 优化序列化方式

默认的pickle序列化效率不高,可以换成cloudpickle或dill,或者自定义__reduce__方法优化structureA的序列化过程,减少进程间数据传递的时间。

5. 合理选择多进程/多线程

如果generateStructA和generateStructB是CPU密集型任务,线程池受GIL限制无法真正并行,还是应该用多进程,但要优化数据传递问题;如果是IO密集型任务,线程池会更高效(无需进程间数据复制)。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 15:01:01