← 返回题目列表

进程池开多少个进程合适?chunksize 该怎么设?

中等 第 17 / 27 题 更新于 2026/08/01
并行调优chunksizecpu_count阿姆达尔定律

简化版

并行不是「进程越多越快」——它有三笔固定成本:进程启动、参数/返回值的 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.maxchunksize 是最容易被忽略的调优点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=2resources.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)),而 imapimap_unorderedExecutor.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))本进程被允许用哪些 CPUos.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)),而 imapimap_unorderedProcessPoolExecutor.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.mapimap(但它在调用时就把所有任务提交进队列,内存并不省),as_completedimap_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_memorynp.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_countimap/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 是十倍量级的调优点★:单次调度固定开销 1050μ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.imapPool.imap_unorderedProcessPoolExecutor.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」,就是为了给负载均衡留出余量。
  • 误区:imapimap_unordered 只是返回顺序不同,性能差不多。 imap 为了保证输出顺序和输入一致,必须缓存那些「已经算完但前面还没算完」的结果——如果第 1 个任务特别慢,后面 999 个的结果会全部堆在内存里等着,既占内存又拖慢整体(消费者拿不到任何结果)。imap_unordered 谁先完成谁先返回,不需要缓存、不需要等待,因此更快也更省内存。所以准则是:只要业务不依赖顺序,一律用 imap_unordered;确实需要知道每个结果对应哪个输入时,让任务返回 (index, result) 自带标识,主进程再按索引归位——这比让框架替你保序高效得多。concurrent.futures 里的对应关系是 ex.mapimapas_completedimap_unordered
  • 误区:开的进程数等于核心数,CPU 就能跑满。 常见的三个「跑不满」原因。① 序列化瓶颈:主进程忙于 pickle/unpickle 和读写管道,成了单点瓶颈——现象是主进程 CPU 100% 而 worker 们都在等。② 阿姆达尔定律:读数据、分发、汇总结果、写输出这些串行部分限制了上限——串行占 10% 时,再多核也只能快 10 倍。③ 隐藏的线程超订:numpy/scipy 底层的 OpenBLAS、MKL 自己会起 N 个线程,你再开 N 个进程就是 N×N 个线程抢 N 个核,上下文切换的开销吃掉全部收益——必须在 import numpy 之前设置 OMP_NUM_THREADS=1MKL_NUM_THREADS=1OPENBLAS_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=1chunksize 是十倍量级的调优点:单次调度固定开销 10~50μs,任务本身只有 1μs 时开销占 95%,设成 1000 后降到 2%;Pool.map/starmap 会自动估算(len / (进程数 × 4),乘 4 是为了负载均衡),但 imapimap_unorderedProcessPoolExecutor.map 默认都是 1,必须手动设(线程池的 chunksize 无效)。调法是:让每块耗时落在 10ms~1s、块数 ≥ 进程数 × 4、任务不均匀时调小(否则长尾——总耗时取决于最慢的 worker)。接口选择上,只要顺序无所谓就用 imap_unorderedimap 为保序会把已完成的结果全缓存在内存里等)。序列化优化三招:别传大对象改传 id、用 initializer 让每个 worker 只加载一次、大数组走 shared_memory/memmap,返回值也要在子进程里先聚合。最后记住阿姆达尔定律:串行部分占 5% 时 8 核最多快 5.9 倍,而「汇总阶段」常常就是那个被忽略的串行瓶颈。