Skip to content

进程池与 CPU 密集任务:什么时候绕开线程

对应视频:EP09《进程池与 CPU 密集任务》
视频入口:B站 EP09 定时稿(计划公开:2026-08-13 18:49)。 对应合集:Python 并发实战:从基础模型到语音工程。 本文为视频的工程展开版,补充代码模板、排查链路、状态字段和检查清单。 最后核验:2026/08/09 08:26:33

视频解释进程池为什么适合一部分 CPU 密集任务。文章补充成本判断、任务粒度和音频处理里的边界。

进程池成本模型图

这张图提醒你同时计算收益和成本:进程池能利用多核,但启动、序列化和内存放大都可能吃掉收益。

适用判断

问题如果答案是“是”
任务主要是纯 Python CPU 计算吗进程池可能有效
任务能拆成相互独立的小批次吗适合并行分发
参数和结果容易序列化吗进程间传输成本可控
每个任务耗时足够长吗能覆盖调度和序列化开销
是否需要共享巨大模型或缓存小心内存放大

最小模板

python
from concurrent.futures import ProcessPoolExecutor


def cpu_heavy(n: int) -> int:
    total = 0
    for value in range(n):
        total += value * value
    return total


def main() -> None:
    inputs = [500_000, 600_000, 700_000, 800_000]
    with ProcessPoolExecutor(max_workers=4) as pool:
        for result in pool.map(cpu_heavy, inputs):
            print(result)


if __name__ == "__main__":
    main()

和音频处理的关系

音频任务推荐方向原因
ffmpeg 转码有界外部命令池主要工作在外部进程
纯 Python 特征计算进程池CPU 密集且可拆分
NumPy 特征计算先测量底层可能已有原生并行
PyTorch 推理GPU 队列或 batch瓶颈可能在 GPU/显存
大模型重复加载谨慎多进程内存可能放大

工程补充:怎么把这篇用到项目里

读图方式

  • 进程池适合大粒度 CPU 任务,任务太小时,序列化和调度成本可能比计算本身还高。
  • 如果每个进程都加载大模型或大缓存,RSS 会被放大,机器可能先被内存打满。
  • 音频特征提取、波形计算、CPU 转码前处理可以试进程池;GPU 推理通常要谨慎。

排查路径

  • 先用 cProfilepy-spy 确认瓶颈真在 CPU 计算,而不是磁盘或 API 等待。
  • 比较不同 max_workerschunksize,记录吞吐、CPU util、RSS 和单任务耗时。
  • 如果 worker 启动慢,考虑 initializer 预加载;如果参数太大,考虑传路径而不是传大对象。

问题、证据和解决动作

问题现场先查什么解决动作
小任务提交特别多单任务耗时、pickle 时间合并任务或设置 chunksize
内存突然翻倍每进程 RSS、模型加载次数减少 worker、共享只读文件或预加载策略
CPU 没跑满磁盘 util、输入读取耗时先优化 I/O 或批量读取
结果回传很慢返回对象大小结果落盘传路径,避免大对象跨进程

落地边界

  • 进程池不是线程池的高级版,它是为 CPU 密集和隔离准备的。
  • 生产环境要考虑进程退出、异常重启和中间结果恢复。

常见失败模式

错误写法现象修正
小任务大量丢进进程池反而变慢增大任务粒度
传输巨大对象序列化耗时高传路径、ID 或小参数
每个进程加载大模型内存爆涨重新设计 worker 或 batch
不测原生库行为重复并行导致过度抢 CPU控制底层线程数并压测

实战检查清单

  • 是否确认任务是 CPU 密集,而不是 I/O 等待。
  • 是否知道参数和结果如何序列化。
  • 是否测过不同任务粒度。
  • 是否观察了总 CPU、内存和进程数。
  • 是否比较了线程池、进程池和顺序版本。

读完之后能完成什么

  • 能判断 CPU 密集任务是否适合进程池。
  • 能说出进程池的启动、序列化和内存成本。
  • 能为音频转码、特征提取、本地推理选择不同执行模型。

配套示例代码

源码文件:examples/ep09/processpool_cpu_template.py

python
from concurrent.futures import ProcessPoolExecutor


def cpu_heavy(n: int) -> int:
    total = 0
    for value in range(n):
        total += value * value
    return total


def main() -> None:
    inputs = [500_000, 600_000, 700_000, 800_000]
    with ProcessPoolExecutor(max_workers=4) as pool:
        for result in pool.map(cpu_heavy, inputs):
            print(result)


if __name__ == "__main__":
    main()
别急,先让缓存热一下。