文章
合集Python 并发编程第 3 / 7 篇

多进程

CPU 密集任务的解决方案。每个进程有独立的 Python 解释器,完全绕开 GIL,真正利用多核 CPU。核心难点是进程间通信(IPC)。

1. Process 基本用法

1.1 和 Thread 的对比

对比维度threading.Threadmultiprocessing.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,当并发量大到线程扛不住时,协程登场。