多进程
CPU 密集任务的解决方案。每个进程有独立的 Python 解释器,完全绕开 GIL,真正利用多核 CPU。核心难点是进程间通信(IPC)。
1. Process 基本用法
1.1 和 Thread 的对比
| 对比维度 | threading.Thread | multiprocessing.Process |
|---|---|---|
| 内存 | 共享进程内存 | 各自独立内存 |
| GIL | 受 GIL 限制 | 每个进程有独立 GIL,不受影响 |
| 创建开销 | 低(KB 级) | 高(MB 级) |
| 适合场景 | I/O 密集 | CPU 密集 |
| 通信方式 | 共享变量 + 锁 | IPC(Queue / Pipe / 共享内存) |
1.2 基本代码
from multiprocessing import Process
import os
def worker(n):
print(f"进程 {n},PID: {os.getpid()}")
if __name__ == "__main__": # Windows 必须有此保护!
processes = [Process(target=worker, args=(i,)) for i in range(4)]
for p in processes:
p.start()
for p in processes:
p.join()
1.3 if __name__ == "__main__" 的 Windows 保护
Linux 默认使用 fork(复制当前进程),没有这个问题。但为了跨平台兼容,永远加上这个保护。
2. 进程间通信(IPC)
进程不共享内存,每个进程有独立的 Python 解释器。要交换数据,必须用 IPC 机制。
2.1 Queue(消息队列,最常用)
from multiprocessing import Process, Queue
def producer(q):
for i in range(5):
q.put(f"消息 {i}")
q.put(None) # 哨兵值,通知消费者结束
def consumer(q):
while True:
item = q.get()
if item is None:
break
print(f"处理: {item}")
if __name__ == "__main__":
q = Queue()
p1 = Process(target=producer, args=(q,))
p2 = Process(target=consumer, args=(q,))
p1.start(); p2.start()
p1.join(); p2.join()
Queue 是进程安全的,内部已处理好锁和序列化。适合生产者-消费者模式。
2.2 Pipe(双向管道,两点通信)
from multiprocessing import Process, Pipe
def sender(conn):
conn.send({"type": "greeting", "data": "hello"})
conn.close()
def receiver(conn):
msg = conn.recv()
print(f"收到: {msg}")
conn.close()
if __name__ == "__main__":
parent_conn, child_conn = Pipe() # 创建一对连接
p1 = Process(target=sender, args=(child_conn,))
p2 = Process(target=receiver, args=(parent_conn,))
p1.start(); p2.start()
p1.join(); p2.join()
Pipe 比 Queue 快,但只支持两点通信。Queue 支持多个生产者和消费者。
2.3 Value / Array(共享内存)
from multiprocessing import Process, Value, Array
import ctypes
def increment(shared_val):
for _ in range(100_000):
shared_val.value += 1
if __name__ == "__main__":
# 共享整数
counter = Value(ctypes.c_int, 0)
processes = [Process(target=increment, args=(counter,)) for _ in range(4)]
for p in processes: p.start()
for p in processes: p.join()
print(counter.value) # 注意:多个进程同时修改可能有竞态条件
共享内存性能最高,但需要自己处理同步(Value 有一个内置的锁参数 lock=True)。适合简单数据。
2.4 Manager(共享复杂结构)
from multiprocessing import Process, Manager
def worker(shared_dict):
shared_dict["key"] = "value"
if __name__ == "__main__":
with Manager() as manager:
shared_dict = manager.dict()
shared_list = manager.list()
p = Process(target=worker, args=(shared_dict,))
p.start()
p.join()
print(dict(shared_dict)) # {'key': 'value'}
Manager 可以共享 dict、list 等复杂结构,但有序列化开销(底层通过服务器进程代理)。性能比 Queue 和共享内存低。
IPC 选型总结
| 方式 | 适用场景 | 速度 | 复杂度 |
|---|---|---|---|
| Queue | 生产者-消费者,多对多 | 中 | 低 |
| Pipe | 两点通信 | 快 | 低 |
| Value/Array | 共享简单数据(计数器、数组) | 最快 | 中(需自己加锁) |
| Manager | 共享复杂结构(dict、list) | 慢 | 低 |
3. ProcessPoolExecutor
3.1 和 ThreadPoolExecutor 统一接口
from concurrent.futures import ProcessPoolExecutor
import math
def compute(n):
"""CPU 密集:计算大数的阶乘。"""
return math.factorial(n)
if __name__ == "__main__":
numbers = [100_000, 200_000, 300_000, 400_000]
# 和 ThreadPoolExecutor 写法完全一样
with ProcessPoolExecutor(max_workers=4) as executor:
results = list(executor.map(compute, numbers))
for n, r in zip(numbers, results):
print(f"{n}! 有 {len(str(r))} 位")
这就是 concurrent.futures 的精髓:写一遍代码,选线程还是进程只改一个类名。
3.2 max_workers 怎么选
# ProcessPoolExecutor 默认 max_workers = CPU 核心数
# 不建议超过核心数太多——每个进程都是独立的 Python 解释器,内存开销大
with ProcessPoolExecutor() as executor: # 默认核心数
pass
with ProcessPoolExecutor(max_workers=os.cpu_count()) as executor: # 显式指定
pass
4. multiprocessing.Pool 高级功能
concurrent.futures 覆盖 95% 场景,但 multiprocessing.Pool 提供了两个额外功能:
4.1 imap(惰性返回,节省内存)
from multiprocessing import Pool
def compute(x):
return x ** 2
if __name__ == "__main__":
with Pool(processes=4) as pool:
# map:全部计算完再返回(可能占用大量内存)
results = pool.map(compute, range(100_000))
# imap:惰性返回,每次返回一个(节省内存)
for result in pool.imap(compute, range(100_000)):
process(result) # 边算边处理
4.2 imap_unordered(不保序,更快)
with Pool(processes=4) as pool:
# imap_unordered:谁先完成先返回,不保证顺序
for result in pool.imap_unordered(compute, range(100_000)):
process(result) # 可能是 0, 4, 1, 9, 16... 不是 0, 1, 4, 9, 16
当顺序不重要时,imap_unordered 比 imap 更快——不需要等慢的先完成。
5. Pickle 序列化的限制
| 不可 pickle 的对象 | 原因 |
|---|---|
lambda | 匿名函数,无法序列化 |
| 闭包(捕获了外部变量的函数) | 外部变量可能不可序列化 |
文件句柄(open() 返回的对象) | 底层是 OS 资源 |
| 数据库连接 | 同上 |
threading.Lock | 锁是线程级资源 |
# 错误示范
def my_func(x):
return x ** 2
# lambda 不能传给 ProcessPoolExecutor
with ProcessPoolExecutor() as executor:
executor.map(lambda x: x ** 2, range(10)) # PicklingError!
# 正确做法:传普通函数
def square(x):
return x ** 2
with ProcessPoolExecutor() as executor:
executor.map(square, range(10)) # OK
6. 小结
| 知识点 | 要记住的 |
|---|---|
| Process vs Thread | 进程独立内存、绕开 GIL、开销高 |
| Windows 保护 | 永远加 if __name__ == "__main__" |
| IPC 选型 | Queue 最通用,Pipe 最快,Manager 最灵活 |
| ProcessPoolExecutor | 接口和 ThreadPoolExecutor 统一,改类名即可切换 |
| Pickle 限制 | lambda、闭包、文件句柄不能传 |
上一篇:python-multithreading,I/O 密集的解决方案。 下一篇:python-asyncio,当并发量大到线程扛不住时,协程登场。