GIL 锁的是线程不是你 — fork 复制父进程内存 / spawn 干净重启, 任务与结果都要跨过 pickle 序列化边界, 进程池摊薄启动成本
if __name__ == "__main__":maxtasksperchild 让 worker 定期重生自愈泄漏initializer 给每个 worker 建独立连接/加载模型PicklingErrorchunksize 降低每任务固定税GIL 保证同一进程里只有一个线程执行字节码, 于是 CPU 密集任务在线程里只能排队。多进程的思路很直接: 再开几个解释器, 每个进程有自己的 GIL 和内存空间, 4 核就是 4 份真并行。代价有两条: 一是出生成本 — fork 把父进程内存整个快照(锁也一起搬, 是死锁温床), spawn 干净重启但要重新 import; 二是序列化边界 — 进程内存互相看不见, 任务、参数、结果全要 pickle 过桥, 数据大了这笔税比计算本身还贵, 所以才需要进程池摊薄启动、chunksize 合并小任务、SharedMemory 零拷贝传大块。
def sha(data: bytes) -> str: ... # CPU 密集: 压缩/哈希 with Pool() as p: hs = p.map(sha, chunks) # 4 核 ≈ 3.8x # 等网络的活: GIL 等 IO 时释放, asyncio/线程就够
import multiprocessing as mp mp.set_start_method("fork") # 写时复制快照, 不重新 import print(mp.get_start_method()) # → fork # 关键: 快, 但父进程的锁/线程状态被原样照搬
if __name__ == "__main__": 保护, 否则子进程 import 主模块又起 Pool, 无限繁殖。
# spawn: 全新解释器, 重新 import 主模块 if __name__ == "__main__": # 关键: 不保护 → 无限繁殖 with mp.Pool() as p: p.map(work, data)
mp.set_start_method("forkserver") # server 进程只 import 一次, 之后所有子进程由它 fork: # 比 spawn 快, 又不带主进程的脏锁/脏线程
def crash(): os._exit(1) # worker 模拟崩溃 p = mp.Process(target=crash); p.start(); p.join() print(p.exitcode) # → 1, 主进程照常跑 # 换成线程崩溃: 整个进程陪葬
__main__ 里定义的类直接 PicklingError。函数按"模块名+限定名"序列化, 所以必须是模块级函数。
def work(x): return x * 2 # 模块级具名: OK with Pool() as p: p.map(work, [1, 2]) # → [2, 4] # 错: p.map(lambda x: x*2, ...) → PicklingError
map 全量阻塞收齐; imap/imap_unordered 边算边收流式消费; apply_async + callback/error_callback 异步提交 — 三种节奏覆盖绝大多数场景。
with Pool() as p: p.map(fn, xs) # 全量阻塞收齐 for r in p.imap_unordered(fn, xs): ... # 流式 p.apply_async(fn, args, callback=on_ok) # 异步+回调
# 每块一次 IPC, 不是每任务一次 p.map(fn, range(1_000_000), chunksize=10_000) # 关键: 小任务默认块太小, 手动调大常有数倍差距
Pool(4, maxtasksperchild=1000) # 关键: worker 干满 1000 个任务就重生, 泄漏随旧进程 # 归还 OS — 内存缓涨自愈阀, 代价是重启成本
Manager().dict() 返回代理对象, 每次读写都是到 server 进程的一次 IPC — 低频共享(配置热更)好用, 高频操作比 Queue 慢一个量级。
cfg = Manager().dict() cfg["rate"] = 10 # 关键: 这次赋值 = 到 server 的一次 IPC # 低频共享(配置热更)好用; 高频循环比 Queue 慢一个量级
close() 解除本进程映射, unlink() 真正回收 — 忘了 unlink 就打满 /dev/shm。
shm = SharedMemory(create=True, size=1024) head = bytes(shm.buf[:4]) # 直接读, 不过 pickle shm.close() # 解除本进程映射 shm.unlink() # 关键: 真正回收, 忘了打满 /dev/shm
p = mp.Process(target=poll, daemon=True) p.start() # 主进程退出时自动终止 # 限制: 不能再生子进程, 也没有正常清理阶段
ProcessPoolExecutor, 用 loop.run_in_executor(pool, fn, arg) 调 — 事件循环不堵, 核也用满, 两全。
pool = ProcessPoolExecutor() r = await loop.run_in_executor(pool, sha, data) # 关键: 事件循环不堵, 核也用满 — async 管 IO, 进程管 CPU
图片解码/重编码是纯 CPU 活, 线程被 GIL 摁住。进程池一开, 每个核一个解释器:
from concurrent.futures import ProcessPoolExecutor from pathlib import Path from PIL import Image def compress(path: str) -> int: img = Image.open(path) # 解码+重编码全程吃 CPU, GIL 主场 out = Path(path).with_suffix(".webp") img.save(out, "WEBP", quality=80) return Path(path).stat().st_size - out.stat().st_size if __name__ == "__main__": # spawn 模式的入场券, 不写必炸 with ProcessPoolExecutor(max_workers=4) as ex: # 4 核 = 4 个独立 GIL saved = sum(ex.map(compress, paths, chunksize=16)) print(f"freed {saved/1e9:.1f}GB") # 实测: 单进程 38min → 4 进程 10min (3.8x); 线程版只有 1.2x
map 要等全量结果才返回, 大文件直接把内存顶爆。imap_unordered 谁先算完谁先回来:
from multiprocessing import Pool def read_blocks(path: str, size: int): # 生成器: 逐块读, 不物化整个文件 with open(path, "rb") as f: while blk := f.read(size): yield blk def score(item: tuple[int, bytes]) -> tuple[int, float]: i, blk = item return i, heavy_model(blk) # 纯 CPU, 无共享状态, 随便并行 with Pool(4) as pool: jobs = enumerate(read_blocks("big.bin", 1 << 20)) # 1MB 一块 for i, s in pool.imap_unordered(score, jobs, chunksize=8): sink.write(f"{i}\t{s}\n") # 流式落盘, 全程只驻留几个块
内存恒定在几十 MB, 而且乱序回收意味着慢块不阻塞快块的结果落盘。
任务本身只有 0.2ms, 但每次分发都要 pickle 参数 + 回传结果, IPC 是每任务固定税:
import math from multiprocessing import Pool n = 1_000_000 # chunksize=1: 一百万次"发送-回传"往返, IPC 占一半以上时间 with Pool(4) as pool: r1 = pool.map(normalize, range(n), chunksize=1) # 21.4s with Pool(4) as pool: r2 = pool.map(normalize, range(n), chunksize=1000) # 5.4s, 同样的计算 # 起步值: math.ceil(n / (workers * 4)), 再压测微调 suggested = math.ceil(n / (4 * 4)) # 注意: chunksize 过大会让先领大块的 worker 忙死、别人饿死, 损失负载均衡
用 Queue 传大数组等于 pickle 序列化一份 + 对端再反序列化一份。共享内存只传一个名字:
from multiprocessing import Pool, shared_memory import numpy as np arr = np.random.rand(2000, 2000) # ~32MB shm = shared_memory.SharedMemory(create=True, size=arr.nbytes) shared = np.ndarray(arr.shape, dtype=arr.dtype, buffer=shm.buf) shared[:] = arr # 拷一次进共享区, 之后所有进程直接读 def worker(args: tuple[str, tuple]) -> float: name, shape = args shm = shared_memory.SharedMemory(name=name) # 按名字 attach, 不过 pickle view = np.ndarray(shape, dtype=np.float64, buffer=shm.buf) r = view.mean() shm.close() # 只解除本进程映射, 不回收区域 return float(r) # 传给 pool 的只有 shm.name 字符串 + shape, 不是 32MB 数组 shm.close(); shm.unlink() # 最后一个使用者负责回收, 防 /dev/shm 泄漏
同规模数据 Queue 传 10 次 ~1.8s, SharedMemory 方案 attach 后近乎免费。
同一份代码 Linux 好、Windows/macOS 挂: spawn 子进程重新 import 模块, 模块级初始化不会"继承"过来。把初始化收进 initializer:
# Windows/macOS 默认 spawn: 子进程重新 import 本模块 MODEL = None # fork 下拿到父进程的值; spawn 下就是 None! def init_worker(path: str) -> None: global MODEL MODEL = load_model(path) # 每个 worker 各自加载一次, 只在启动时付成本 def infer(text: str) -> str: return MODEL.predict(text) # 运行期不再碰全局初始化逻辑 if __name__ == "__main__": # 不写: 子进程 import 时又执行这里 → 无限起进程 with Pool(4, initializer=init_worker, initargs=("model.bin",)) as pool: results = pool.map(infer, texts)
批处理跑两小时后 worker RSS 从 2.1GB 爬到 3.5GB, 第三方库泄漏抓不动。定期让 worker "转世":
# worker 每处理 5000 个任务就重生: 泄漏的内存/句柄随旧进程归还 OS with Pool(8, maxtasksperchild=5000) as pool: for meta in pool.imap(process_page, pages, chunksize=32): progress.mark(meta) # 观测: 改完后 worker RSS 平稳在 2.1GB, 再没触发过 OOM kill # 代价: 重生要重新 import + 跑 initializer # 5000 是"自愈频率 vs 重启成本"的折中, 按单任务耗时压测取值
想在不重启 worker 的情况下调限流阈值。Manager 提供跨进程可见的代理 dict:
from multiprocessing import Manager, Pool import time CONFIG = None # 每个 worker 里都是同一个 Manager dict 的代理 def call_api(task: str) -> int: rate = CONFIG["rate_limit"] # 每次读是一次到 Manager 的 RPC time.sleep(1 / rate) return len(task) if __name__ == "__main__": with Manager() as mgr: CONFIG = mgr.dict({"rate_limit": 100}) with Pool(16) as pool: pool.map(call_api, tasks) # 运维通道执行 CONFIG["rate_limit"] = 200, worker 下次读取即生效 # 低频读没问题; 高频计数请换 Value/SharedMemory, 别走 RPC
fork 出的 worker 用父进程连接, 服务端把四个进程当同一会话: 事务串号、command out of sync 连环报错。每个 worker 自己建:
CONN = None def init_db(dsn: str) -> None: global CONN CONN = psycopg.connect(dsn) # 各自建连接: 会话独立, 锁也不共享 def query(sql: str) -> list: return CONN.execute(sql).fetchall() if __name__ == "__main__": # initializer 在每个 worker 启动时跑一次, 连接随 worker 生死 with Pool(4, initializer=init_db, initargs=(DSN,)) as pool: rows = pool.map(query, sqls) # 配 maxtasksperchild 时注意: worker 重生会重连一次, 池要给连接数留余量
异步提交最大的坑是异常静默: worker 里抛的错只存在于 result 对象里, 不 get 就蒸发。回调 + 结果对象双保险:
from multiprocessing import Pool def job(payload: dict) -> dict: try: return {"ok": True, "data": render(payload)} except Exception as e: # worker 内先兜住, 装进结果回传 return {"ok": False, "error": repr(e)} def on_done(res): cache.set(res["id"], res) def on_fail(exc): alert.fire(f"pool job crashed: {exc!r}") res = pool.apply_async(job, (payload,), callback=on_done, error_callback=on_fail) res.get(timeout=30) # 需要同步语义时再等; 不 get 也别忘 error_callback
Web 服务里直接调指纹/加密等 CPU 函数, 事件循环一卡所有并发请求陪葬。扔进进程池:
import asyncio from concurrent.futures import ProcessPoolExecutor async def handler(req): # 直接 fingerprint(req.body) 会堵死事件循环, 全部请求一起卡 loop = asyncio.get_running_loop() digest = await loop.run_in_executor(POOL, fingerprint, req.body) return {"digest": digest} if __name__ == "__main__": with ProcessPoolExecutor(max_workers=4) as POOL: asyncio.run(serve(handler)) # 复用同一个池: 每请求 fork 一个进程的话, 启动开销吃光并行收益
P99 从 3.4s 回到 180ms — 事件循环负责 IO 并发, 进程池负责 CPU 满载, 各干各的。
pool.map(lambda x: ..., data) 或带闭包的函数丢给子进程, 直接 PicklingError: Can't pickle local object。正解: 任务函数写成模块级具名函数, 状态当参数显式传。
# 错: pool.map(lambda x: x * 2, data) # → PicklingError: Can't pickle local object def work(x): return x * 2 # 对: 模块级具名函数 pool.map(work, data) # 状态当参数显式传
# 错: 线程持有 logging/DB 锁时 fork → 锁被复制无人释放 # 程序跑几小时随机卡死 mp.set_start_method("spawn") # 对: 干净子进程 # 或 fork 前确保没有并发线程在跑
if __name__ == "__main__": + 初始化全走 initializer, 在最小 spawn 环境自测。
# 错: 模块顶层直接建 Pool → spawn 平台无限起进程 if __name__ == "__main__": # 对: 入口保护 mp.set_start_method("spawn") main() # 初始化全走 initializer
pool.map(fn, huge_generator) 内部先把 iterable 转成 list, 10GB 流直接撑爆内存。正解: imap/imap_unordered 流式消费, 或自己分批喂 map。
# 错: pool.map(fn, huge_generator) → 内部先转 list # 10GB 流直接撑爆内存 for r in pool.imap_unordered(fn, huge_generator): # 对 save(r) # 流式: 边算边收不物化
# 错: worker 被 OOM kill, 主进程只剩 BrokenProcessPool 一句 def work(x): try: return transform(x), None except Exception as e: return None, repr(e) # 对: 装进结果 # OOM 真因去 dmesg | tail 查
# 错: worker 任务里又建 Pool → 进程数按乘积爆炸 def task(x): return heavy(x) # 对: 保持单层拓扑 # 要嵌套: 拆两级队列, 由主进程统一调度
try/finally 保证 unlink; 注册 atexit; 排查用 ls /dev/shm 清理孤儿区。
shm = SharedMemory(create=True, size=4096) try: use(shm) finally: shm.close(); shm.unlink() # 对: 崩了也回收 # 排查: ls /dev/shm 清理孤儿区
mgr.dict() 每次读写都是一次跨进程 RPC, 循环里累加百万次比 Queue 慢一个量级。正解: 高频计数用 Value/共享内存 + 锁, Manager 只留给低频配置类共享。
# 错: mgr.dict() 每次读写都是一次跨进程 RPC cnt = mp.Value("i", 0) # 对: 共享 ctypes, 无 RPC with cnt.get_lock(): cnt.value += 1 # Manager 只留给低频配置类共享
CONFIG["x"] = 1 改的是自己内存副本, 其他进程根本看不见。正解: 跨进程状态走 Manager/Queue/SharedMemory; 纯只读配置在 initializer 里加载一份。
CONFIG = {"n": 0}
def work(x):
CONFIG["n"] += 1 # 错: 改的是本进程副本
return x
# 对: 状态走 Manager/Queue/SharedMemory; 只读配置
# 在 initializer 里各加载一份RuntimeError: daemonic processes are not allowed to have children。正解: 去掉 daemon 标记并自己管理生命周期, 或把二级并行改成线程/池内任务。
# 错: 守护 worker 里再开进程 # → RuntimeError: daemonic processes are not allowed # to have children p = mp.Process(target=job, daemon=False) # 对: 自己管生命周期
pool.join() 挂住不返回, 因为还有 worker 没退出。正解: 顺序是先 pool.close()(不再收任务) 再 pool.join(); 要放弃就 terminate()。with 语句会自动做这套。
# 错: pool.join() 直接挂住 (worker 还没退) pool.close() # 对: 先关闸, 不再收任务 pool.join() # 对: 再等 worker 全部退出 # with Pool() as p: ... 自动做这套
QueueHandler 单点写出, 带进程名前缀。
# 错: worker 里直接 print → 输出拦腰交错, 事后没法看 qh = QueueHandler(log_q) # 对: 日志进队列 logging.getLogger().addHandler(qh) # 主进程单点写出, 带进程名前缀
pool.terminate(); 编排系统里配 stop_grace_period 覆盖清理窗口。
# 错: kill 主进程 → worker 变孤儿, 白吃 CPU 占着队列 signal.signal(signal.SIGTERM, lambda *a: pool.terminate()) # 对: 先收池 # shell 巡检: ps aux | grep python
# 错: ThreadPoolExecutor 正跑着又 fork → 线程"消失" # 但它们持有的锁/半成品状态全被带进子进程 with ProcessPoolExecutor() as pool: # 对: 先建进程池 pool.submit(heavy, data) # 再开线程, 或直接换 spawn
# 错: 每任务 0.1ms 计算 + 10MB 参数 # pickle 来回 50ms → 并行反而更慢 p.map(batch_work, batches) # 对: 粗粒度打包一次传一批 # 大只读数据放 SharedMemory, 只传名字
# 错: async def hash_loop(): 哈希循环写成协程 # → 单线程空转发烫, 吞吐没变 with ProcessPoolExecutor() as pool: # 对: 进程池 hashes = list(pool.map(sha, chunks))
pid = os.fork() # 错: 裸 fork 后调 malloc 家族可能死锁 if pid == 0: os.execv("/bin/ls", ["ls"]) # 对: 尽快 exec 换掉映像 # 一般场景: 直接用 multiprocessing 的 fork 抽象
set_start_method("spawn"), CI 覆盖率配 COVERAGE_PROCESS_START + 并行合并。
# conftest.py — 错: 默认 fork → 告警 + 子进程覆盖率丢失 mp.set_start_method("spawn") # 对: 测试里显式 spawn # CI: COVERAGE_PROCESS_START=1 + coverage combine
with Pool() as p; 手动管理就 try/finally 里 close/terminate + join; 定期 ps aux | grep 巡检。
# 错: 异常路径没清理 → worker 残留成僵尸, 句柄占满 with Pool() as p: # 对: 无论如何都 close/join p.map(work, data) # 手动管理: try/finally 里 close/terminate + join
q.close(); q.join_thread(), 或用 Pool 的结果回传通道代替自管 Queue。
q.put(big_obj) # 错: put 返回 ≠ 对端收到, feeder 还在刷 q.close() # 对: 退出前 q.join_thread() # 对: 等 feeder 刷完再退