求助:如何进一步优化multiprocess/multiprocessing库对CPU密集型类方法的提速效果
求助:如何进一步优化multiprocess/multiprocessing库对CPU密集型类方法的提速效果
我现在用multiprocess库加速类里的CPU密集任务,模拟处理500页文档,串行跑大概18-19秒,用多进程最多只能跑到6秒左右,但我电脑有16核,目标是跑到2秒以内。目前换了标准库的multiprocessing后有小幅提升,但还是远不如预期,想问问有没有其他优化方法?
代码实现
dummy.py
from multiprocess import Pool class DummyClass: def __init__(self, workers=1): self.workers = workers # 模拟CPU密集任务 def _process_one(self, page): count = 0 while count < 2_000_000: count += 1 return page # 用"multiprocess"多进程处理 def multiprocess(self, pages): with Pool(processes=self.workers) as pool: async_results = pool.map_async(self._process_one, pages) extraction = async_results.get() return extraction # 串行处理 def sequential(self, pages): extraction = [] for page in pages: extraction.append(self._process_one(page)) return extraction
test.py
import time from dummy import DummyClass # 串行测试 def dummy_sequential(): dummy_extractor = DummyClass() extraction = dummy_extractor.sequential(range(500)) return extraction # 多进程测试 def dummy_multiprocess(workers): dummy_extractor = DummyClass(workers=workers) extraction = dummy_extractor.multiprocess(range(500)) return extraction if __name__ == "__main__": # 串行测试结果 ini = time.time() dummy_sequential() fin = time.time() print(f"sequential: {fin - ini:.2f}s") # 多进程测试结果(multiprocess库) print("\nUsing multiprocess:") for i in range(2, 9): ini = time.time() dummy_multiprocess(workers=i) fin = time.time() print(f"{i} workers: {fin - ini:.2f}s") # 多进程测试结果(标准库multiprocessing) print("\nUsing multiprocessing:") # 需先把dummy.py的导入改为from multiprocessing import Pool for i in range(2, 9): ini = time.time() dummy_multiprocess(workers=i) fin = time.time() print(f"{i} workers: {fin - ini:.2f}s")
现有测试结果
# 串行 sequential: 18.77s # multiprocess库 2 workers: 10.71s 3 workers: 7.68s 4 workers: 8.11s 5 workers: 7.36s 6 workers: 6.99s 7 workers: 7.18s 8 workers: 7.20s # 标准库multiprocessing 2 workers: 10.34s 3 workers: 7.27s 4 workers: 6.55s 5 workers: 6.70s 6 workers: 6.80s 7 workers: 6.26s 8 workers: 6.30s
优化建议(Stack Overflow老玩家视角)
兄弟,16核机器只跑到6秒确实没发挥出硬件实力,核心问题大概率是任务粒度太小+类方法序列化开销,给你几个落地的优化方向,挨个试试:
1. 给进程池加chunksize参数(最快见效)
你现在每个进程处理单个page(0.04秒/任务),进程间调度、序列化的开销占比太高了!map_async支持chunksize参数,让每个进程处理一批page,直接减少调度次数:
def multiprocess(self, pages): # 按进程数均分任务批次,避免出现空批次 chunksize = max(1, len(pages) // self.workers) with Pool(processes=self.workers) as pool: # 加上chunksize参数 async_results = pool.map_async(self._process_one, pages, chunksize=chunksize) extraction = async_results.get() return extraction
比如16核的话,500页均分后每个进程处理31-32页,单进程任务时间从0.04秒变成1.2秒左右,开销占比直接下降,CPU能跑满。
2. 把任务函数移出类,减少序列化开销
Windows下multiprocessing/multiprocess用spawn启动子进程,会序列化整个类实例传给子进程,这部分开销不小。把_process_one改成独立函数,避免序列化整个类:
# 移到DummyClass外面的独立函数 def process_one(page): count = 0 while count < 2_000_000: count += 1 return page class DummyClass: # ...其他代码不变... def multiprocess(self, pages): chunksize = max(1, len(pages) // self.workers) with Pool(processes=self.workers) as pool: # 调用独立的process_one函数 async_results = pool.map_async(process_one, pages, chunksize=chunksize) extraction = async_results.get() return extraction
这样子进程只需要序列化函数本身,而不是整个类实例,能省不少序列化时间。
3. 换用concurrent.futures.ProcessPoolExecutor试试
标准库的ProcessPoolExecutor在Windows下的spawn模式优化可能更好,API也更简洁,试试替换Pool:
from concurrent.futures import ProcessPoolExecutor class DummyClass: # ...其他代码不变... def multiprocess(self, pages): chunksize = max(1, len(pages) // self.workers) with ProcessPoolExecutor(max_workers=self.workers) as executor: # map方法自动支持chunksize extraction = list(executor.map(process_one, pages, chunksize=chunksize)) return extraction
4. 验证CPU使用率,排除系统干扰
打开任务管理器看看运行时CPU是不是跑满100%:
- 如果没跑满:可能是后台程序(比如Windows Defender)占了资源,或者任务粒度还是太小,继续调大chunksize;
- 如果跑满但速度还是上不去:试试把
process_one的循环次数改到2000万,看看是不是任务计算量不够大(毕竟16核CPU单循环很快),大任务下并行收益会更明显。
5. 试试用loky库(可选)
loky是专门解决Windows下spawn开销问题的进程池实现,安装后直接替换,用法和concurrent.futures完全一致:
pip install loky
from loky import ProcessPoolExecutor
备注:内容来源于stack exchange,提问作者zest16
相关产品推荐
相关产品推荐

