进程池开多少个进程合适?chunksize 该怎么设?
简化版
并行不是「进程越多越快」——它有三笔固定成本:进程启动、参数/返回值的 pickle 序列化、进程间数据传输。任务本身越小,这三笔成本占比越高,甚至可能比串行还慢。判断能不能并行受益的第一步是算账:单个任务的计算时间要远大于「序列化 + 传输」的时间(经验值:单任务至少几十毫秒才值得进程池)。进程数怎么定:CPU 密集型用 os.cpu_count()(或略少 1~2 个留给主进程和系统),IO 密集型可以远超核心数。但 os.cpu_count() 在容器里是个大坑——它返回的是宿主机的核心数,而不是 cgroup 限制的配额,在一个限了 2 核的 Pod 里可能返回 64,导致你开 64 个进程互相争抢、大量上下文切换反而更慢;正确做法是用 os.process_cpu_count()(Python 3.13+) 或 len(os.sched_getaffinity(0)),容器里还要读 cgroup 的 cpu.max。chunksize 是最容易被忽略的调优点:Pool.map/imap 默认 chunksize=1 意味着每个元素都要单独经历一次「序列化 → 传输 → 反序列化」,处理 100 万个小任务时调度开销可能占掉 90% 以上的时间;把 chunksize 设成几百到几千,同样的任务能快十倍以上。核心记忆:先算账再并行;cpu_count 在容器里不可信;小任务多必须调大 chunksize。
详细版
三种并行执行接口对比:
| 接口 | 返回 | 顺序 | 内存 | chunksize 默认 |
|---|---|---|---|---|
Pool.map(f, it) | 列表(全部算完才返回) | 保序 | 全部结果进内存 | 自动估算 |
Pool.imap(f, it, cs) | 迭代器(边算边给) | 保序 | 流式 | 1 |
Pool.imap_unordered | 迭代器 | 谁先完谁先给 | 流式 | 1 |
Pool.starmap(f, its) | 列表 | 保序 | 全部进内存 | 自动估算 |
Executor.map(f, it, chunksize=) | 迭代器 | 保序 | 流式(但提交时全部入队) | 1 |
import os, time, multiprocessing as mp
from concurrent.futures import ProcessPoolExecutor
# ① ★进程数:cpu_count 在容器里不可信★
print(os.cpu_count()) # ★宿主机核心数★(容器里可能是 64)
print(len(os.sched_getaffinity(0))) # ★进程实际可用的 CPU(Linux)★
# Python 3.13+:
# print(os.process_cpu_count()) # ★推荐:已考虑亲和性★
def usable_cpus() -> int:
# 1) cgroup v2 配额(容器里最准)
try:
quota, period = open("/sys/fs/cgroup/cpu.max").read().split()
if quota != "max":
return max(1, int(int(quota) / int(period)))
except OSError:
pass
# 2) CPU 亲和性
try:
return len(os.sched_getaffinity(0))
except AttributeError:
pass
return os.cpu_count() or 1
# ② ★chunksize 的威力(同样的任务,差十倍)★
def square(x):
return x * x
data = range(1_000_000)
with mp.Pool(4) as pool:
t0 = time.perf_counter()
r1 = pool.map(square, data, chunksize=1) # ★每个元素单独调度★
t1 = time.perf_counter()
r2 = pool.map(square, data, chunksize=10000) # ★一次发一万个★
t2 = time.perf_counter()
# 实测量级:chunksize=1 约 12s;chunksize=10000 约 0.9s ← ★十倍以上差距★
# 原因:每次调度都要 pickle 一个 int + 走队列 + unpickle,
# 而 x*x 本身只要几十纳秒 → ★调度开销占 99%★
# ③ ★Pool.map 会自动估算 chunksize(但 imap/Executor.map 不会)★
# Pool.map 内部:chunksize, extra = divmod(len(iterable), len(pool) * 4)
# if extra: chunksize += 1
# → 100 万元素 / (4 进程 × 4) = 62500
# ★imap / imap_unordered / Executor.map 的默认 chunksize 都是 1★→ 必须手动设
# ④ 三种接口的选择
with mp.Pool(4) as pool:
res = pool.map(work, items) # 要全部结果、顺序重要
for r in pool.imap(work, items, chunksize=100): # ★流式,省内存★
handle(r)
for r in pool.imap_unordered(work, items, 100): # ★谁先完谁先处理,最快★
handle(r)
# ⑤ ★序列化开销:大对象传参是灾难★
import numpy as np
big = np.zeros((10000, 10000)) # 800MB
# pool.map(process, [big] * 4) # ✗ ★每个任务都 pickle 800MB★
# ✓ 用共享内存 / 只传索引 / 让子进程自己读文件
from multiprocessing import shared_memory
# ⑥ 测量:并行到底值不值
def measure(n_workers, chunksize):
t0 = time.perf_counter()
with ProcessPoolExecutor(n_workers) as ex:
list(ex.map(work, data, chunksize=chunksize))
return time.perf_counter() - t0
# ★永远和串行版本对比★:serial = timeit(list(map(work, data)))
⚠️ 三个必须记住的点:①
os.cpu_count()返回的是「机器有多少个 CPU」,不是「我能用多少个 CPU」。在 Docker/K8s 里,--cpus=2或resources.limits.cpu: "2"是通过 cgroup 的 CPU 配额实现的(限制的是时间片总量,不是可见核心数),所以os.cpu_count()照样返回宿主机的 64——按它开 64 个进程,实际只有 2 核的算力,结果是大量上下文切换、内存翻 32 倍、性能反而暴跌。Python 3.13 起提供了os.process_cpu_count()(考虑了 CPU 亲和性),但它仍然读不到 cgroup 配额,容器里最稳的做法是读/sys/fs/cgroup/cpu.max或从环境变量注入。②chunksize的默认值不统一:Pool.map/starmap会自动估算(约len(iterable) / (进程数 × 4)),而imap、imap_unordered、Executor.map的默认值都是 1——这意味着你用ProcessPoolExecutor.map处理一百万个小任务时,每个元素都要单独走一遍 pickle + 队列 + unpickle,调度开销可能占掉 99% 的时间。③ 返回值也要序列化:只关注入参大小是不够的——一个返回 100MB DataFrame 的任务,结果同样要 pickle → 通过管道传回 → unpickle,这笔开销经常比计算本身还大,所以并行任务应该在子进程里就把结果聚合/压缩成小对象再返回。
完整版教学
一、先算账:并行的三笔固定成本
并行执行一个任务的完整成本:
① ★进程启动★(只在池创建时付一次,但 spawn 很贵)
fork: ~1ms spawn: ~50~200ms(要重启解释器 + 重新 import)
② ★参数序列化 + 传输 + 反序列化★(★每个任务都要付★)
pickle 一个小 int: ~1μs;一个 1MB 的 DataFrame: ~10ms
走管道传输: 取决于大小
③ ★返回值序列化 + 传输 + 反序列化★(同样每个任务都要付)
④ 调度开销:队列锁、任务分发
★ 判断公式(粗略但管用):
并行有收益 ⟺ 单任务计算时间 ≫ (参数序列化 + 返回值序列化 + 调度)
经验阈值:
单任务 < 1ms → ★几乎必然更慢★(除非 chunksize 开得很大)
单任务 1~10ms → 要调大 chunksize 才有收益
单任务 > 100ms → ★并行收益明显★
单任务 > 1s → 随便怎么写都有收益
一个真实的算例(★说明"并行更慢"是常态★):
任务:对 100 万个整数求平方(每个约 50 纳秒)
串行: 100万 × 50ns = ★0.05 秒★
进程池(4),chunksize=1:
每个任务的调度+序列化开销 ≈ 12μs
100万 × 12μs / 4 进程 = ★3 秒★ ← ★慢了 60 倍★
进程池(4),chunksize=10000:
100 个块 × (序列化 1 万个 int ≈ 2ms + 计算 0.5ms) / 4 = ★0.06 秒★
→ 勉强追平串行,但★仍然没有收益★(任务本身太小)
★ 结论:对"每个元素计算量极小"的任务,
正确答案不是调 chunksize,而是★根本不要用多进程★
→ 改用 numpy 向量化(0.05s → 0.001s)才是数量级的提升
什么任务真正适合进程池:
✓ 图像处理(每张图几十~几百毫秒)
✓ 文件解析(每个文件几十毫秒以上)
✓ 模型推理 / 特征计算
✓ 加密、压缩、哈希大文件
✓ 独立的模拟/回测任务
✗ 简单算术、字符串处理、小字典操作(★用向量化或直接串行★)
并行不是免费的午餐,它有三笔固定成本:进程启动(fork 约 1ms、spawn 约 50~200ms)、每个任务都要付的参数序列化 + 传输 + 反序列化、以及返回值的同样开销。判断公式很朴素:并行有收益的前提是「单任务计算时间远大于序列化和调度开销」,经验阈值是单任务小于 1ms 几乎必然更慢,大于 100ms 才收益明显。那个算例很能说明问题:对 100 万个整数求平方,串行只要 0.05 秒,而用进程池加 chunksize=1 要 3 秒——慢了 60 倍;即使把 chunksize 调到 10000 也只是勉强追平。结论是:对「每个元素计算量极小」的任务,正确答案不是调 chunksize,而是根本不要用多进程——改用 numpy 向量化才是数量级的提升。真正适合进程池的是图像处理、文件解析、模型推理、加密压缩这类单任务几十毫秒以上的工作。
二、进程数怎么定:cpu_count 的陷阱
★ 三个"CPU 数量"完全不是一回事:
os.cpu_count() ★机器上有多少个逻辑 CPU★(含超线程)
← ★不考虑 cgroup 限制、不考虑亲和性★
len(os.sched_getaffinity(0)) ★本进程被允许使用哪些 CPU★(Linux,考虑 taskset)
os.process_cpu_count() ★Python 3.13+,等价于上面那个,跨平台★
cgroup 的 cpu.max ★容器实际分到的 CPU 时间配额★
★ 容器里的经典事故:
K8s 里 resources.limits.cpu: "2"
→ 通过 cgroup ★限制 CPU 时间配额★(每 100ms 只能用 200ms 的 CPU 时间)
→ ★不改变可见的核心数★ → os.cpu_count() 仍然返回宿主机的 64
→ 代码里 ProcessPoolExecutor()(默认 = cpu_count)开了 64 个进程
→ 后果:① 内存 × 64(每个进程都要加载解释器和模块)★可能直接 OOM★
② 64 个进程抢 2 核的时间片 → ★大量上下文切换★
③ 被 cgroup 限流(throttling)→ 延迟尖刺
→ 实测:限 2 核时,开 2 个进程比开 64 个★快 3~5 倍★
✓ 正确的获取方式(按优先级):
def usable_cpus():
# cgroup v2
try:
q, p = open("/sys/fs/cgroup/cpu.max").read().split()
if q != "max": return max(1, round(int(q) / int(p)))
except OSError: pass
# cgroup v1
try:
q = int(open("/sys/fs/cgroup/cpu/cpu.cfs_quota_us").read())
p = int(open("/sys/fs/cgroup/cpu/cpu.cfs_period_us").read())
if q > 0: return max(1, round(q / p))
except OSError: pass
try: return len(os.sched_getaffinity(0))
except AttributeError: return os.cpu_count() or 1
✓ 或者★由部署层通过环境变量注入★(最省事、最可控):
workers = int(os.getenv("WORKER_PROCESSES", usable_cpus()))
进程数的经验值:
┌────────────────────┬──────────────────────────────────────┐
│ 纯 CPU 密集 │ ★N 或 N-1★(留一个给主进程/系统) │
│ CPU + 少量 IO │ N ~ 1.5N │
│ IO 为主(但用进程) │ 2N ~ 4N(★更该用线程或 asyncio★) │
│ 内存受限 │ ★由内存决定★:可用内存 / 单进程峰值内存 │
│ 超线程机器 │ 物理核心数往往比逻辑核心数更优(实测决定)│
└────────────────────┴──────────────────────────────────────┘
★ 别忘了"隐藏的并行":
numpy/scipy 底层的 OpenBLAS/MKL ★自己会起 N 个线程★
→ 你再开 N 个进程 → ★N × N 个线程抢 N 个核★ → 严重超订、性能暴跌
✓ 必须显式限制:
os.environ["OMP_NUM_THREADS"] = "1" # ★在 import numpy 之前设★
os.environ["MKL_NUM_THREADS"] = "1"
os.environ["OPENBLAS_NUM_THREADS"] = "1"
→ 这是数据处理场景最常见的性能坑之一
进程数的第一个陷阱是三个「CPU 数量」完全不是一回事:os.cpu_count() 是机器有多少逻辑 CPU(不考虑 cgroup 和亲和性)、len(os.sched_getaffinity(0)) 是本进程被允许用哪些 CPU、os.process_cpu_count()(3.13+)是它的跨平台版本、而 cgroup 的 cpu.max 才是容器实际分到的配额。容器里的经典事故就是:K8s 限了 2 核(通过限制 CPU 时间配额实现,不改变可见核心数),而 ProcessPoolExecutor() 默认按 cpu_count() 开了 64 个进程——内存翻 64 倍可能直接 OOM、64 个进程抢 2 核导致大量上下文切换、还会被 cgroup 限流产生延迟尖刺,实测比开 2 个进程慢 3~5 倍。第二个常被忽略的是**「隐藏的并行」**:numpy 底层的 OpenBLAS/MKL 自己会起 N 个线程,你再开 N 个进程就是 N×N 个线程抢 N 个核,严重超订——必须在 import numpy 之前设 OMP_NUM_THREADS=1 等环境变量。
三、chunksize:最容易被忽略的十倍差距
chunksize 是什么:★一次分发给一个 worker 多少个元素★
chunksize=1 每个元素单独:pickle → 队列 → unpickle → 计算 → 结果回传
chunksize=1000 一次打包 1000 个:★1 次调度开销摊到 1000 个元素上★
为什么默认 1 会慢十倍(★算清这笔账★):
单次调度的固定开销 ≈ 10~50μs(pickle + 队列锁 + IPC + unpickle)
任务本身 = 0.1ms
chunksize=1: 开销/总时间 = 20μs / 120μs = ★17%★
任务本身 = 0.001ms(1μs)
chunksize=1: 开销/总时间 = 20μs / 21μs = ★95%!★
chunksize=1000: 开销/总时间 = 20μs / 1020μs = ★2%★
→ ★任务越小,chunksize 越要大★
各接口的默认值(★必须记住★):
Pool.map / Pool.starmap ★自动估算★:
chunksize = ceil(len(iterable) / (进程数 × 4))
Pool.imap / imap_unordered ★默认 1★ ← 手动设!
Executor.map(ProcessPool) ★默认 1★ ← 手动设!
Executor.map(ThreadPool) chunksize ★无效★(线程池不需要分块,直接忽略)
怎么选 chunksize(★工程做法★):
① 先估:让"每个块的处理时间"落在 ★10ms ~ 1s★ 之间
chunksize ≈ 目标块耗时 / 单任务耗时
单任务 0.1ms、目标 100ms → chunksize ≈ 1000
② 再校验总块数:★块数应该 ≥ 进程数 × 4★(保证负载均衡)
chunksize = min(估算值, len(data) // (workers * 4))
③ 任务耗时不均匀时 → ★调小 chunksize★(大块会导致长尾)
例:处理不同大小的图片,有的 10ms 有的 10s
→ chunksize 太大时,一个 worker 可能拿到一堆大图 → ★其他 worker 空闲等它★
④ 实测微调(★不同数据分布最优值差很多★)
★ chunksize 的两个副作用:
① ★内存★:一次要 pickle chunksize 个元素 + 缓存 chunksize 个结果
chunksize=100000 且每个元素 1MB → ★100GB★
② ★长尾★:块越大,最后一个块的落单时间越长
→ 总耗时 ≈ max(各 worker 完成时间),而不是平均值
一个平衡的通用公式:
def auto_chunksize(n_items, n_workers, target_chunk_ms=100, item_ms=1):
by_time = max(1, int(target_chunk_ms / item_ms))
by_balance = max(1, n_items // (n_workers * 4)) # ★保证块数够多★
return min(by_time, by_balance)
chunksize 决定「一次分发给一个 worker 多少个元素」,它的影响可以是十倍量级。算清这笔账就明白了:单次调度的固定开销约 10~50μs——任务本身 0.1ms 时开销占 17%,任务本身只有 1μs 时开销占 95%,而把 chunksize 设成 1000 后开销降到 2%。各接口的默认值不统一是最大的坑:Pool.map/starmap 会自动估算(约 len / (进程数 × 4)),而 imap、imap_unordered、ProcessPoolExecutor.map 默认都是 1(线程池的 chunksize 参数则完全无效,会被忽略)。工程做法是三步:先按「让每块处理时间落在 10ms~1s」估算、再校验块数至少是进程数的 4 倍(保证负载均衡)、任务耗时不均匀时要调小(否则一个 worker 拿到一堆大任务,其他 worker 空闲等它,形成长尾)。还要注意两个副作用:内存(一次 pickle chunksize 个元素)和长尾(总耗时取决于最慢的 worker 而非平均值)。
四、map / imap / imap_unordered 的取舍
三者的行为差异:
pool.map(f, items)
→ ★阻塞直到全部完成★,返回完整列表
→ ★所有输入先物化成列表★(内存 O(n)),所有结果也在内存里
→ 顺序和输入一致
pool.imap(f, items, chunksize)
→ 返回★迭代器★,边算边给(可以立刻开始处理前面的结果)
→ ★保序★:即使第 2 个先算完,也要等第 1 个才能 yield
→ 内存友好(★但仍会缓存乱序完成的结果直到能按序给出★)
pool.imap_unordered(f, items, chunksize)
→ ★谁先完成谁先返回★
→ ★最快、内存最省★(不需要缓存等待排序)
→ 代价:结果顺序和输入无关 → ★要自己带上标识★
★ 保序的隐藏成本(面试爱问):
imap 保序意味着:如果第 1 个任务特别慢,后面 999 个都算完了,
它们的结果★全部缓存在内存里等着★ → 内存可能爆
→ 只要业务不依赖顺序,★优先 imap_unordered★
→ 需要对应关系时,让任务返回 (index, result):
for idx, res in pool.imap_unordered(work_with_index, enumerate(items), 100):
results[idx] = res
选择决策树:
结果要全部拿到 + 数据量不大 → map(最简单)
数据量大 / 想边出边处理 → imap 或 imap_unordered
★顺序无所谓★ → ★imap_unordered(最优)★
需要顺序但数据量大 → imap_unordered + 自带索引
要能提前终止(找到就停) → imap_unordered + break(★注意要 terminate 池★)
concurrent.futures 的对应物:
ex.map(f, items, chunksize=N) ≈ imap(保序、流式)
★但注意:ex.map 在调用时就把★所有任务提交进队列★(内存不省)
as_completed(futures) ≈ imap_unordered
★ ProcessPoolExecutor.map 的 chunksize 默认是 1 —— 大数据量必须手动设
提前终止的正确姿势:
with mp.Pool(4) as pool:
for r in pool.imap_unordered(check, candidates, chunksize=50):
if r.found:
pool.terminate() # ★立刻停掉所有 worker★
break
★ 不 terminate 的话,剩余任务会继续跑完(浪费算力)
★ with 块退出时 Pool 会调 terminate;但 ProcessPoolExecutor 的
shutdown(cancel_futures=True) 只能取消★未开始★的
三个接口的核心差异是返回时机和顺序:map 阻塞到全部完成并把所有输入和结果都放内存;imap 返回保序的迭代器(边算边给);imap_unordered 谁先完成谁先返回。这里有个保序的隐藏成本值得记住:imap 为了保序,如果第 1 个任务特别慢,后面已经算完的 999 个结果全部缓存在内存里等着——数据量大时可能爆内存。所以只要业务不依赖顺序,优先用 imap_unordered(最快、内存最省);需要对应关系时让任务返回 (index, result) 自带标识即可。concurrent.futures 的对应关系是:ex.map ≈ imap(但它在调用时就把所有任务提交进队列,内存并不省),as_completed ≈ imap_unordered。最后一个实用点:提前终止时要调 pool.terminate(),否则剩余任务会继续跑完、白白浪费算力。
五、序列化开销:并行性能的隐形杀手
每个任务的完整数据流:
主进程:pickle(参数) → 写管道 → 子进程:unpickle → 计算
→ pickle(返回值) → 写管道 → 主进程:unpickle
★ 两次 pickle + 两次 unpickle + 两次 IPC,全都算在"并行开销"里
pickle 的实际速度(量级感受):
小 int / 短字符串 ~1μs
1000 元素的 list ~50μs
1MB 的 numpy 数组 ~1ms(★纯内存拷贝,还算快★)
1MB 的 DataFrame ~10ms(★结构复杂,慢得多★)
100MB 的 DataFrame ~1s ← ★比很多计算本身还久★
含大量 Python 对象的嵌套结构 ★极慢★(每个对象都要单独处理)
★ 三条优化原则:
① ★别传大对象,传"怎么拿到它"★
✗ pool.map(process, [huge_df] * 100) # ★100 次 pickle 大对象★
✓ pool.map(process_by_id, range(100)) # 子进程自己从文件/DB 读
✓ 或用 initializer 在每个 worker 里加载一次:
def init(path):
global DATA
DATA = load(path) # ★每个进程只加载一次★
Pool(4, initializer=init, initargs=(path,))
② ★返回值也要瘦身★
✗ 返回完整的处理后 DataFrame(100MB × 1000 个任务)
✓ 在子进程里就聚合成统计量 / 写文件后只返回路径
③ ★大数组用共享内存,绕开 pickle★
from multiprocessing import shared_memory
shm = shared_memory.SharedMemory(create=True, size=arr.nbytes)
buf = np.ndarray(arr.shape, dtype=arr.dtype, buffer=shm.buf)
buf[:] = arr[:]
# 子进程用 shm.name 附加,★零拷贝★
✓ 或者:np.memmap(文件映射,多进程共享页缓存)
✓ 或者:Arrow / Parquet 等零拷贝格式
★ initializer 模式(★进程池最重要的优化手法之一★):
每个 worker 进程启动时执行一次初始化:
- 加载模型 / 大词表 / 配置
- 建立数据库连接(★fork 后必须在子进程里建★)
- 设置日志
→ 避免"每个任务都传一遍大对象"或"每个任务都重新建连接"
def init_worker():
global model, conn
model = load_model() # ★只在进程启动时执行一次★
conn = create_connection()
with ProcessPoolExecutor(4, initializer=init_worker) as ex:
ex.map(predict, batches, chunksize=10)
★ 什么时候 pickle 直接失败:
lambda、局部函数、闭包、生成器、打开的文件/socket、线程锁、
嵌套定义的类 → ★spawn 模式下全部传不了★
✓ 改成模块级函数 + functools.partial 绑定参数
序列化是并行性能的隐形杀手——每个任务要经历两次 pickle + 两次 unpickle + 两次 IPC。要有量级感:小 int 约 1μs,1MB 的 DataFrame 约 10ms、100MB 的约 1 秒——比很多计算本身还久。三条优化原则:① 别传大对象,传「怎么拿到它」(传 id 让子进程自己读,或用 initializer 在每个 worker 里加载一次);② 返回值也要瘦身(在子进程里就聚合成统计量,或写文件后只返回路径);③ 大数组用 shared_memory 或 np.memmap 绕开 pickle(零拷贝)。其中 initializer 是进程池最重要的优化手法之一:每个 worker 启动时执行一次初始化(加载模型、建数据库连接、设置日志),避免「每个任务都传一遍大对象」或「每个任务都重新建连接」。最后记住 spawn 模式下 lambda、闭包、局部函数、打开的文件和锁全都无法 pickle,要改成模块级函数配 functools.partial。
六、实测方法与决策清单
★ 测量的正确姿势(不测就是瞎调):
① 永远和★串行版本★对比
serial_time = time_it(lambda: [work(x) for x in data])
→ 如果并行没有明显快过它,就★不要并行★
② 扫描参数网格
for workers in (1, 2, 4, 8, 16):
for cs in (1, 10, 100, 1000, 10000):
print(workers, cs, measure(workers, cs))
→ ★最优值往往出乎意料★
③ 用真实数据规模(小数据测不出调度开销的影响)
④ 关注三个指标:总耗时、峰值内存、CPU 利用率
★CPU 利用率上不去 = 有瓶颈★(序列化?IO?锁?)
★ 加速比的天花板:阿姆达尔定律
加速比 = 1 / (串行部分 + 并行部分/N)
串行占 10% → 无限个核最多也只能快 ★10 倍★
串行占 5%、8 核 → 1/(0.05 + 0.95/8) = ★5.9 倍★(不是 8 倍)
★ 这里的"串行部分"包括:读数据、分发任务、汇总结果、写结果
→ ★所以"汇总阶段"经常是真正的瓶颈★
常见的"并行了但没变快"原因排查表:
┌────────────────────────────┬────────────────────────────────┐
│ 现象 │ 可能原因 │
├────────────────────────────┼────────────────────────────────┤
│ 比串行还慢 │ ★任务太小★,调度开销占主导 │
│ CPU 利用率低 │ 序列化瓶颈 / IO 等待 / 块数太少 │
│ 只用了 2 核(明明开了 16) │ ★容器 CPU 限制★ / 亲和性 │
│ 内存爆 │ chunksize 太大 / map 全量物化 │
│ 加速比远低于核心数 │ ★阿姆达尔★:汇总/读写是串行的 │
│ 最后几秒只有一个进程在跑 │ ★长尾★:chunksize 太大或任务不均│
│ 开了 N 个进程反而变慢 │ ★BLAS 线程超订★(N×N 个线程) │
└────────────────────────────┴────────────────────────────────┘
★ 最终决策清单:
□ 先算账:单任务耗时 > 10ms 吗?不到就别并行(考虑向量化)
□ 进程数:用 ★process_cpu_count / cgroup 配额★,不要裸用 cpu_count
□ import numpy 前设 ★OMP_NUM_THREADS=1★(防线程超订)
□ chunksize:imap/Executor.map ★必须手动设★(默认 1)
□ 目标:每块 10ms~1s,块数 ≥ 进程数 × 4
□ 任务不均匀 → 调小 chunksize
□ 大对象:用 initializer 预加载 / shared_memory / 只传 id
□ 返回值瘦身(在子进程里聚合)
□ 顺序不重要 → imap_unordered
□ ★实测对比串行★,扫参数网格
调优的第一条纪律是永远和串行版本对比——如果并行没有明显快过串行,就不要并行。同时要理解加速比的天花板:阿姆达尔定律——串行部分占 5%、8 核时最多只能快 5.9 倍而不是 8 倍,而这里的「串行部分」包括读数据、分发任务、汇总结果、写结果,所以汇总阶段经常才是真正的瓶颈。那张排查表值得记住:比串行还慢通常是任务太小、只用了 2 核往往是容器 CPU 限制、最后几秒只剩一个进程在跑是长尾(chunksize 太大或任务不均)、开了 N 个进程反而变慢多半是 BLAS 线程超订。最终决策清单里最关键的四条:先算账(单任务 >10ms 才值得并行)、进程数用 cgroup 配额而不是裸 cpu_count、imap/Executor.map 必须手动设 chunksize、大对象用 initializer 预加载而不是反复传参。
记忆钩子:「并行有三笔固定成本:进程启动、★每个任务都要付的参数 pickle★、★返回值 pickle★——所以第一步永远是★算账★:单任务计算时间必须远大于序列化+调度开销(经验值:★<1ms 几乎必然更慢、>100ms 才收益明显★)。对 100 万个整数求平方,串行 0.05 秒而进程池 chunksize=1 要 3 秒(★慢 60 倍★),这种情况正确答案是 ★numpy 向量化而不是多进程★。★进程数的头号坑:os.cpu_count() 返回的是机器有多少核,不是我能用多少核★——容器的 CPU 限制是通过 ★cgroup 时间配额★实现的、不改变可见核心数,所以限 2 核的 Pod 里它照样返回 64,按它开进程会内存翻 64 倍 + 上下文切换 + 被限流;要用 ★os.process_cpu_count()(3.13+) / sched_getaffinity / 读 cgroup 的 cpu.max★。还要防★隐藏的并行★:numpy 底层 BLAS 自己起 N 个线程,再开 N 个进程就是 N×N 抢 N 核 → ★import numpy 前设 OMP_NUM_THREADS=1★。★chunksize 是十倍量级的调优点★:单次调度固定开销 10
50μs,任务只有 1μs 时开销占 95%,设成 1000 后降到 2%;★Pool.map/starmap 会自动估算(len/(进程数×4)),但 imap、imap_unordered、ProcessPoolExecutor.map 默认都是 1,必须手动设★(线程池的 chunksize 无效)。调法:让每块耗时落在 10ms1s,且块数 ≥ 进程数×4;★任务不均匀要调小★(否则长尾——总耗时取决于最慢的 worker)。接口选择:★顺序无所谓一律用 imap_unordered★(imap 为保序会把已完成的结果全缓存在内存里等)。序列化优化三招:★别传大对象传 id、用 initializer 每个 worker 只加载一次、大数组走 shared_memory/memmap★;返回值也要在子进程里先聚合。最后记住★阿姆达尔定律★:串行占 5% 时 8 核最多快 5.9 倍,而『汇总阶段』常常就是那个串行瓶颈。」
七、常见误区与追问
- 误区:
os.cpu_count()就是我能用的 CPU 数量。 它返回的是机器上有多少个逻辑 CPU,完全不考虑 cgroup 配额和 CPU 亲和性。容器的 CPU 限制(--cpus=2、K8s 的limits.cpu: "2")是通过限制 CPU 时间配额实现的——每个调度周期内只能用固定的 CPU 时间,可见的核心数并不改变,所以os.cpu_count()在一个限了 2 核的 Pod 里照样返回宿主机的 64。按它创建进程池的后果是三重的:内存翻 64 倍(每个进程都要加载解释器和全部模块,很可能直接 OOM)、64 个进程争抢 2 核的时间片造成大量上下文切换、以及被 cgroup 限流产生延迟尖刺——实测比老老实实开 2 个进程慢 3~5 倍。正确做法是用os.process_cpu_count()(3.13+)、len(os.sched_getaffinity(0)),或直接读/sys/fs/cgroup/cpu.max;最省事的是由部署层通过环境变量注入。 - 误区:
chunksize是个小优化,用默认值就行。 它经常是十倍量级的差异,而且默认值在不同接口之间不统一:Pool.map/starmap会自动估算(约len(iterable) / (进程数 × 4)),但Pool.imap、Pool.imap_unordered和ProcessPoolExecutor.map的默认值都是 1。默认 1 意味着每一个元素都要单独走一遍「pickle → 队列 → IPC → unpickle」,而单次调度的固定开销是 10~50μs——当任务本身只要几微秒时,99% 的时间花在调度上。实测处理 100 万个小任务,chunksize=1要 12 秒、chunksize=10000只要 0.9 秒。所以用imap/Executor.map处理大量小任务时必须手动设置chunksize(顺带一提,线程池的chunksize参数是无效的,会被直接忽略,因为线程池不需要跨进程分发)。 - 误区:chunksize 越大越好,反正减少了调度开销。 有两个反向代价。① 内存:一次要 pickle
chunksize个元素的参数、还要缓存这一块的所有结果——chunksize=100000且每个元素 1MB 就是 100GB。② 长尾:块越大,任务分配的粒度越粗,越容易出现「其他 worker 都空闲了,最后一个 worker 还在啃一个大块」的情况;并行任务的总耗时取决于最慢的那个 worker,而不是平均值。任务耗时不均匀时这个问题尤其严重(比如处理大小差异很大的图片,一个 worker 可能连续拿到几张巨图)。工程上的平衡是:让每块的处理时间落在 10ms~1s,同时保证总块数至少是进程数的 4 倍(这也正是Pool.map自动估算公式里× 4的用意)——Pool.map之所以除以「进程数 × 4」,就是为了给负载均衡留出余量。 - 误区:
imap和imap_unordered只是返回顺序不同,性能差不多。imap为了保证输出顺序和输入一致,必须缓存那些「已经算完但前面还没算完」的结果——如果第 1 个任务特别慢,后面 999 个的结果会全部堆在内存里等着,既占内存又拖慢整体(消费者拿不到任何结果)。imap_unordered谁先完成谁先返回,不需要缓存、不需要等待,因此更快也更省内存。所以准则是:只要业务不依赖顺序,一律用imap_unordered;确实需要知道每个结果对应哪个输入时,让任务返回(index, result)自带标识,主进程再按索引归位——这比让框架替你保序高效得多。concurrent.futures里的对应关系是ex.map≈imap、as_completed≈imap_unordered。 - 误区:开的进程数等于核心数,CPU 就能跑满。 常见的三个「跑不满」原因。① 序列化瓶颈:主进程忙于 pickle/unpickle 和读写管道,成了单点瓶颈——现象是主进程 CPU 100% 而 worker 们都在等。② 阿姆达尔定律:读数据、分发、汇总结果、写输出这些串行部分限制了上限——串行占 10% 时,再多核也只能快 10 倍。③ 隐藏的线程超订:numpy/scipy 底层的 OpenBLAS、MKL 自己会起 N 个线程,你再开 N 个进程就是 N×N 个线程抢 N 个核,上下文切换的开销吃掉全部收益——必须在 import numpy 之前设置
OMP_NUM_THREADS=1、MKL_NUM_THREADS=1、OPENBLAS_NUM_THREADS=1。排查时先看「主进程 CPU 是否打满」和「总线程数是否远超核心数」这两个信号。 - 追问:
Pool.map的默认 chunksize 是怎么算的,为什么要除以「进程数 × 4」? CPython 的实现是chunksize, extra = divmod(len(iterable), len(self._pool) * 4),有余数就加一——也就是把数据切成大约「进程数 × 4」个块。乘以 4 是为了负载均衡:如果只切成「进程数」个块,每个 worker 恰好拿一块,一旦某块特别慢(任务耗时不均),其他 worker 做完就只能干等,总耗时被最慢的那块决定;切成 4 倍数量的块后,快的 worker 可以继续领取后续的块,形成天然的动态负载均衡。倍数选 4 是经验权衡:太小(=1)负载不均,太大则调度开销上升。这也解释了为什么手动设chunksize时要遵守「块数 ≥ 进程数 × 4」这条约束。另外注意Pool.map会先把整个可迭代对象物化成列表(因为要算长度),所以它不适合处理无法一次装进内存的数据流——那种情况要用imap并手动指定chunksize。 - 追问:进程池里怎么高效地共享一个大对象(比如几 GB 的模型)? 按代价从低到高有四种方案。①
initializer预加载(最常用):ProcessPoolExecutor(n, initializer=init, initargs=(path,)),每个 worker 进程启动时加载一次并存进模块级全局变量,之后所有任务直接用——避免了「每个任务都 pickle 传一遍」。②fork+ COW:Linux 上父进程先加载好,fork 出的子进程理论上共享物理内存;但引用计数会破坏 COW(读对象也会改ob_refcnt),需要配合gc.freeze(),而且 fork 在多线程程序里不安全。③multiprocessing.shared_memory(3.8+):把数组数据放进共享内存段,子进程用SharedMemory(name=...)附加、用np.ndarray(..., buffer=shm.buf)零拷贝访问——适合 numpy 数组这类连续数值数据,注意要负责unlink()清理。④np.memmap/ Arrow / Parquet:把数据放在文件里做内存映射,多进程共享同一份页缓存,操作系统自动管理,最省心也最适合超出内存的数据。反模式是pool.map(f, [big_obj] * n)——每个任务都会完整 pickle 一遍大对象。 - 追问:什么时候「并行」这条路本身就是错的? 四种情况。① 单任务太小:每个元素只要几微秒的算术或字符串操作,调度开销必然占主导——正确解法是 numpy/pandas 向量化(把循环下沉到 C 层),能带来数量级提升,而不是把慢循环分给多个进程。② 瓶颈不在 CPU:任务其实卡在磁盘 IO 或网络上,多进程只会让 IO 争抢更严重——应该用线程池或 asyncio(IO 等待期间 GIL 会释放)。③ 数据传输量远大于计算量:比如「传 100MB 进去、算 10ms、传 100MB 出来」,序列化开销是计算的百倍——要么把计算下推到数据所在的地方,要么改用共享内存。④ 串行部分占比太高:读取和汇总占了大半时间,按阿姆达尔定律加速比上不去——先优化那部分。判断方法很简单:先写串行版本并 profile,看清楚时间花在哪里,再决定要不要并行——没有 profile 就并行,多半是在给自己制造复杂度。
八、加强记忆
并行有三笔固定成本:进程启动、每个任务都要付的参数 pickle、以及返回值 pickle——所以第一步永远是算账:单任务的计算时间必须远大于「序列化 + 调度」开销(经验值:单任务 <1ms 几乎必然更慢、>100ms 才收益明显)。那个算例要记住:对 100 万个整数求平方,串行 0.05 秒而进程池配 chunksize=1 要 3 秒(慢 60 倍)——这种情况的正确答案是 numpy 向量化而不是多进程。进程数的头号坑是 os.cpu_count() 返回「机器有多少核」而不是「我能用多少核」:容器的 CPU 限制通过 cgroup 时间配额实现、不改变可见核心数,所以限 2 核的 Pod 里它照样返回 64,按它开进程会导致内存翻倍、上下文切换和 cgroup 限流;应该用 os.process_cpu_count()(3.13+)、sched_getaffinity、或读 /sys/fs/cgroup/cpu.max。还要防隐藏的并行:numpy 底层的 BLAS 自己会起 N 个线程,再开 N 个进程就是 N×N 个线程抢 N 个核——必须在 import numpy 之前设 OMP_NUM_THREADS=1。chunksize 是十倍量级的调优点:单次调度固定开销 10~50μs,任务本身只有 1μs 时开销占 95%,设成 1000 后降到 2%;Pool.map/starmap 会自动估算(len / (进程数 × 4),乘 4 是为了负载均衡),但 imap、imap_unordered、ProcessPoolExecutor.map 默认都是 1,必须手动设(线程池的 chunksize 无效)。调法是:让每块耗时落在 10ms~1s、块数 ≥ 进程数 × 4、任务不均匀时调小(否则长尾——总耗时取决于最慢的 worker)。接口选择上,只要顺序无所谓就用 imap_unordered(imap 为保序会把已完成的结果全缓存在内存里等)。序列化优化三招:别传大对象改传 id、用 initializer 让每个 worker 只加载一次、大数组走 shared_memory/memmap,返回值也要在子进程里先聚合。最后记住阿姆达尔定律:串行部分占 5% 时 8 核最多快 5.9 倍,而「汇总阶段」常常就是那个被忽略的串行瓶颈。