Python多进程任务pruning后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)
核心问题分析
- 进程间数据传递开销过大
多进程模式下,structureA在进程间传递需要经过序列化(默认用pickle)和反序列化。如果structureA体积较大,这个过程会占用大量CPU和IO资源,导致父进程在数据传递上耗时过多,子进程陷入等待,CPU利用率下降。 - 双Pool调度不匹配
生成A的Pool仅4个进程,生成B的Pool有32个,当A的生成速度跟不上B的处理速度时,B的Pool会出现大量空闲进程,整体CPU利用率自然下降。 - 清理/回调逻辑干扰调度
初始版本的清理循环中,父进程频繁遍历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

