首页
/ CPython multiprocessing 完全指南:进程级并行、进程池与进程间通信精讲

CPython multiprocessing 完全指南:进程级并行、进程池与进程间通信精讲

2026-09-07 12:55:03作者:戚魁泉Nursing

multiprocessing 是 CPython 标准库中实现基于进程的并行(process-based parallelism)的核心包:它以一套与 threading 高度相似的 API 创建真实操作系统进程,从而绕开 CPython 全局解释器锁(GIL)的限制,让程序能够真正利用多核 CPU,并同时支持 POSIX 与 Windows。本文以仓库中的官方参考文档 Doc/library/multiprocessing.rst 为骨架,结合标准库 Lib/multiprocessing/ 的真实源码实现,系统讲解启动方式(start methods)、进程对象、队列/管道通信、同步原语、共享状态、Manager 代理、进程池与编程准则;读完你将能够根据平台特性选择合适的启动方式,并用进程池与消息传递设计出可稳定运行、可跨平台的数据并行程序。

multiprocessing 的设计定位与包结构

multiprocessing 的核心卖点并不只是"新开进程",而是threading 的编程模型平移到进程模型上Process 对标 threading.ThreadQueue 对标 queue.QueueLock/RLock/Semaphore 等同步原语全部有对应实现。差异在于 threading 共享同一进程内存,受 GIL 约束;而 multiprocessing 通过子进程隔离内存,从而"side-stepping the Global Interpreter Lock",可以把任务真正铺开到多颗处理器上。

从源码结构看,CPython 把该包拆分为若干职责单一的子模块(均在 Lib/multiprocessing/ 下):

一个最典型的数据并行示例——把 f 作用到多个输入上并收集结果:

from multiprocessing import Pool

def f(x):
    return x*x

if __name__ == '__main__':
    with Pool(5) as p:
        print(p.map(f, [1, 2, 3]))

输出:

[1, 4, 9]

注意两件事:其一,f 定义在模块顶层而非 REPL 里,以便子进程 import 该模块;其二,程序用 with Pool(5) 作为上下文管理器,退出时自动回收工作进程(详见后文"进程池"一节)。

相比 threadingmultiprocessing 还额外提供了没有对应物的能力:你可以对一个运行中的进程调用 Process.terminate 强行终止、用 Process.interrupt()(Python 3.14 新增)发送 SIGINT,或直接 Process.kill() 发送 SIGKILL。如果你的场景只是"往后台提交一批任务并拿结果",官方文档建议优先考虑更上层的 concurrent.futures.ProcessPoolExecutor:它把"提交任务"与"等待结果"解耦得更干净,且其 Future 与其他生态兼容性更好。

说明:本仓库为 CPython 主干,Include/patchlevel.h 表明当前版本为 3.16.0a0,因此本文中标注为 3.14/3.15 新增或变更的 API 均已在当前源码中生效。

用 Process 创建子进程

进程的创建遵循 "create Processstart()" 模式,Process 全面复刻 threading.Thread 的 API。最小示例:

from multiprocessing import Process

def f(name):
    print('hello', name)

if __name__ == '__main__':
    p = Process(target=f, args=('bob',))
    p.start()
    p.join()

扩展一下,用 os.getpid() / os.getppid() 观察父子进程的真实 PID:

from multiprocessing import Process
import os

def info(title):
    print(title)
    print('module name:', __name__)
    print('parent process:', os.getppid())
    print('process id:', os.getpid())

def f(name):
    info('function f')
    print('hello', name)

if __name__ == '__main__':
    info('main line')
    p = Process(target=f, args=('bob',))
    p.start()
    p.join()

if __name__ == '__main__' 包裹的必要性来自 multiprocessing主模块安全导入要求:spawn/forkserver 启动子进程时需要重新导入主模块,若导入过程无条件触发建进程就会递归失控。若你直接在交互式 REPL 里定义 f 再创建进程,spawn 方式的子进程在反序列化时找不到 __main__ 里的 f,会以 AttributeError 崩溃——这正是"目标函数必须位于可导入模块中"这一约束的直观体现。此约束从 Python 3.14 起不再仅限 spawn:3.14 之后 fork 在任何平台上都不再是默认启动方式,所以即便用 fork 也建议遵守同样的代码组织规范。

Process 构造参数与常用成员

构造签名(详见 Doc/library/multiprocessing.rst):

Process(group=None, target=None, name=None, args=(), kwargs={}, *, daemon=None)
  • group:永远应为 None,仅为兼容 threading.Thread 保留;
  • target:子进程 run() 时调用的可调用对象,默认 None 表示不调用任何东西;
  • name:进程名;缺省时自动生成形如 Process-1:1 的名字(N_1:N_2:...:N_k 表示第 k 层第 N 个子进程);
  • args/kwargs:传给 target 的位置/关键字参数;所有参数都必须可 pickle
  • daemon:关键字专用参数,True/False 直接设定守护标志,None(默认)则继承创建者进程的标志。

子类化时,构造函数里第一件事必须是调用 super().__init__()Process 的常用方法/属性包括:

成员 语义
run() 代表进程活动;默认调用构造时传入的 target(*args, **kwargs),子类可覆写
start() 安排 run() 在新进程中执行;每个进程对象最多调用一次
join([timeout]) 阻塞直至进程结束;timeout 为正数时最多等 timeout 秒;超时返回 None,需用 exitcode 判断是否真的终止;不可 join 自身(会死锁),也不可 join 尚未 start 的进程
name 字符串名,仅用于标识,无语义,可重复
is_alive() 进程是否存活
daemon 布尔守护标志,须在 start() 前设置;守护进程退出时其守护子进程会被一并终止;守护进程不允许再创建子进程;它不是 Unix daemon,只是"父进程退出即被终止、不会被 join"的普通进程
pid 子进程 PID;spawn 前为 None
exitcode 退出码:正常返回为 0sys.exit(N)Nrun() 内未捕获异常为 1;被信号 N 终止为 -N
authkey 进程认证密钥(字节串)。主进程初始化时用 os.urandom 生成随机串,Process 对象默认继承父进程的 authkey,可手工改为其它字节串
sentinel 进程结束时变为 ready 的系统对象句柄(POSIX 上为文件描述符,Windows 上为 OS handle),可配合 multiprocessing.connection.wait 一次等待多个事件
interrupt() Python 3.14 新增:向子进程发送 SIGINT(默认效果等价于在子进程中抛 KeyboardInterrupt,因此子进程若吞掉 KeyboardInterrupt 将不会被终止);默认会把 exitcode 置为 1
terminate() POSIX 发送 SIGTERM,Windows 调用 TerminateProcess不会执行退出处理器与 finally,其子孙进程会成为孤儿;若该进程正在用管道/队列/锁,资源可能损坏
kill() terminate,但 POSIX 用 SIGKILL(Python 3.7 起)
close() 释放资源;进程仍在运行则抛 ValueError;成功后多数方法/属性再访问会抛 ValueError

一个完整的生命周期演示(来自文档 doctest,展示了 spawn 上下文下进程状态与退出码):

>>> import multiprocessing, time, signal
>>> mp_context = multiprocessing.get_context('spawn')
>>> p = mp_context.Process(target=time.sleep, args=(1000,))
>>> print(p, p.is_alive())
<...Process ... initial> False
>>> p.start()
>>> print(p, p.is_alive())
<...Process ... started> True
>>> p.terminate()
>>> time.sleep(0.1)
>>> print(p, p.is_alive())
<...Process ... stopped exitcode=-SIGTERM> False
>>> p.exitcode == -signal.SIGTERM
True

需要留意:startjoinis_aliveterminateexitcode只能由创建该进程对象的进程调用

相关异常

  • ProcessError:本模块所有异常基类;
  • BufferTooShortConnection.recv_bytes_into() 提供的缓冲区太小时抛出,e.args[0] 携带完整消息字节串;
  • AuthenticationError:认证失败时抛出;
  • TimeoutError:带超时的方法超时时抛出(注意这与 queue.Empty/queue.Full 不同,后者是队列超时的信号)。

Contexts 与三种启动方式

multiprocessing 依据平台支持三种启动方式(start method)——spawnforkforkserver,这是决定跨平台兼容性的第一关键概念。

spawn

父进程启动一个全新的 Python 解释器作为子进程。子进程只继承运行 run() 所需的必要资源,尤其不会继承父进程多余的 fd 与句柄;代价是启动较慢。POSIX 与 Windows 均可用,且是 Windows 与 macOS 的默认方式。从 3.14 起它不再是任何平台的默认(POSIX 上被 forkserver 取代),但在 Windows/macOS 仍是默认。

fork

父进程调用 os.fork() 复制出子进程。子进程开始时与父进程几乎完全一致、继承父进程全部资源。安全地 fork 一个多线程进程是有问题的:3.12 起,若 Python 检测到当前进程存在多个线程,os.fork() 会抛 DeprecationWarning(详见 os.fork 文档);3.8 起 macOS 上因系统库可能启动线程导致子进程崩溃,已不再推荐。3.14 起 fork 在任意平台上都不再是默认启动方式,需要 fork 的代码必须显式通过 get_context/set_start_method 指定。仅 POSIX 可用。

forkserver

程序启动并选定 forkserver 后,会先派生一个单线程的 fork server 进程;此后每当需要新进程,父进程连接该 server 并请求其 fork。由于 server 单线程,用它 fork 通常是安全的,且不会继承多余资源。适用于支持通过 Unix 管道传递 fd 的 POSIX 平台(如 Linux),并且 3.14 起成为 POSIX 平台的默认启动方式(同时保留性能并规避常见多线程 fork 不兼容问题,参见 CPython gh-84559)。

各版本变迁要点:3.4 起 spawn 在全部 POSIX 平台可用、forkserver 在部分平台可用,且 Windows 上子进程不再继承父进程全部可继承句柄;3.8 起 macOS 默认 spawn;3.14 起 POSIX 默认从 fork 改为 forkserver。

另外,在 POSIX 上用 spawn/forkserver 还会额外启动一个 resource tracker 进程:它记录程序各进程创建的"已 unlink 的命名系统资源"(命名信号量、SharedMemory 等)。当所有进程退出后,tracker 会 unlink 掉残留对象。正常退出时通常没有残留;但若进程被信号杀死,可能泄漏部分命名信号量或共享内存段——它们要等到下次重启才会被系统回收,而命名信号量数量有限、共享内存段占用主存,故值得重视。

如何在代码中选择启动方式

set_start_method 设置(每个程序最多调用一次,应置于 if __name__ == '__main__' 内):

import multiprocessing as mp

def foo(q):
    q.put('hello')

if __name__ == '__main__':
    mp.set_start_method('spawn')
    q = mp.Queue()
    p = mp.Process(target=foo, args=(q,))
    p.start()
    print(q.get())
    p.join()

另一种更推荐的做法是用 get_context 拿到上下文对象。上下文与模块 API 完全一致,且允许同一程序内混用多种启动方式:

import multiprocessing as mp

def foo(q):
    q.put('hello')

if __name__ == '__main__':
    ctx = mp.get_context('spawn')
    q = ctx.Queue()
    p = ctx.Process(target=foo, args=(q,))
    p.start()
    print(q.get())
    p.join()

注意:不同上下文的对象互不兼容——例如 fork 上下文创建的锁不能传给 spawn/forkserver 启动的进程。在 Lib/multiprocessing/context.py 的源码中可以看到默认上下文的组装逻辑:POSIX 上默认 _default_context 指向 forkserver,其余平台(含 Windows)默认指向 spawn;模块顶层导出的 ProcessQueueLock 等实质上都是"默认上下文"的方法。因此:

  • 库的开发者若使用 multiprocessingProcessPoolExecutor应当允许使用方注入自己的上下文,切勿在库内部私自固定某一种启动方式,以免与使用方应用冲突;若确有要求,必须在文档中写明。
  • 警告:spawn/forkserver 通常不能用于"冻结"后的可执行文件(PyInstaller、cx_Freeze 等在 POSIX 上生成的二进制);fork 在无线程场景下也许可行。

全局启动方式的管理

模块级提供了如下函数来查询与操纵全局启动方式:

  • get_all_start_methods():返回支持列表,首项为默认;可选值为 'fork''spawn''forkserver'
  • get_context(method=None):返回带模块全量 API 的上下文对象;method=None 返回默认上下文,并且若全局启动方式尚未设定,会顺手把它设为系统默认。
  • get_start_method(allow_none=False):返回当前启动方式名称,可能为 'fork'/'spawn'/'forkserver'/Noneallow_none=True 时不触发默认值设置。
  • set_start_method(method, force=False):设置启动方式;若已设置过且 force=False 则抛 RuntimeErrormethod=None, force=True 会把启动方式重置为 None。必须在 if __name__ == '__main__' 中保护,且至多调用一次。
  • set_executable(executable):指定用于启动子进程的解释器路径(默认 sys.executable),嵌入式场景常需 set_executable(os.path.join(sys.exec_prefix, 'pythonw.exe'));3.11 起接受 path-like 对象。
  • set_forkserver_preload(module_names, *, on_error='ignore'):让 forkserver 主进程预导入模块清单,fork 出的进程直接继承其已导入状态,避免每个子进程重复导入开销。须在 forkserver 启动(创建 Pool 或 start Process)之前调用;仅对 forkserver 有意义。on_error 控制预导入 ImportError 的处理:"ignore"(默认,静默)、"warn"(在 forkserver 子进程 stderr 发 ImportWarning)、"fail"(forkserver 带 traceback 退出,后续建进程将得到 EOFError/ConnectionError);该参数为 3.15 新增。

全局启动方式还有一个隐含机制:多个 multiprocessing 函数/方法在实例化某些对象时会隐式把全局启动方式设为系统默认(若尚未设置)。全局启动方式只能被设置一次,因此若你需要偏离系统默认,必须在使用这些函数/方法或创建对象之前主动设置。

另外 Lib/multiprocessing/context.py 在模块导入阶段依据平台构造了 ForkContextSpawnContextForkServerContextDefaultContext 再按操作系统选择默认,这解释了文档里"默认值随平台变化"在源码层面的落地方式。

进程间通信:Queues 与 Pipes

多进程编程的原则是尽量用消息传递而非共享状态multiprocessing 提供两类通道。

Queue

Queuequeue.Queue 的近亲,线程与进程均安全,且所有放入的对象都会被 pickle 序列化

from multiprocessing import Process, Queue

def f(q):
    q.put([42, None, 'hello'])

if __name__ == '__main__':
    q = Queue()
    p = Process(target=f, args=(q,))
    p.start()
    print(q.get())    # prints "[42, None, 'hello']"
    p.join()

Pipe

Pipe() 返回一对由管道连接、默认**双向(duplex)**的连接对象:

from multiprocessing import Process, Pipe

def f(conn):
    conn.send([42, None, 'hello'])
    conn.close()

if __name__ == '__main__':
    parent_conn, child_conn = Pipe()
    p = Process(target=f, args=(child_conn,))
    p.start()
    print(parent_conn.recv())   # prints "[42, None, 'hello']"
    p.join()

危险点:若两个进程/线程同时从管道同一端读写,数据可能损坏;不同端同时使用则无此风险。send 序列化对象、recv 重新创建对象。

三种队列类型

QueueSimpleQueueJoinableQueue 都是模仿 queue.Queue 的多生产者多消费者 FIFO:

  • Queue([maxsize]):基于管道加若干锁/信号量实现。进程首次 put 时启动一个 feeder 线程,把缓冲中的对象冲刷到管道。实现 queue.Queuetask_donejoinshutdown 外的全部方法。实例化它可能隐式设置全局启动方式。
  • SimpleQueue():简化的队列,接近"加了锁的 Pipe",只有 close/empty/get/put
  • JoinableQueue([maxsize])Queue 子类,额外提供 task_done()join()每个被取出的任务都必须调用一次 task_done,否则统计未完成任务数的信号量最终可能溢出并抛异常。

队列的方法语义(阻塞与超时)沿用 queue.Queue

  • put(obj[, block[, timeout]])block=True(默认)且 timeout=None 时阻塞到有空位;timeout 为正数时最多等待并抛 queue.Fullblock=False 时有空位立即放入,否则立即抛 queue.Full(忽略 timeout)。队列已关闭时抛 ValueError(3.8 起,替代原 AssertionError)。
  • get([block[, timeout]]):对称语义,超时抛 queue.Empty;队列已关闭时抛 ValueError(3.8 起替代 OSError)。
  • get_nowait() = get(False)put_nowait(obj) = put(obj, False)
  • qsize():近似长度,不可靠;在 macOS 等 sem_getvalue() 未实现的平台可能抛 NotImplementedError
  • empty()/full():受并发语义影响,均不可靠;关闭的队列上调用 empty() 可能抛 OSError

此外还有通常用不到的清理方法:

  • close():释放内部资源,之后不可再 get/put/empty;feeder 线程会在冲刷完缓冲后退出(GC 时自动调用)。
  • join_thread():仅 close() 之后可用,阻塞直到 feeder 线程退出,确保缓冲全部写入管道。默认情况下,非创建者的进程退出时会尝试 join 队列的 feeder 线程。
  • cancel_join_thread():阻止 join_thread 阻塞(允许进程"不冲刷缓冲就退出"),很可能造成已入队数据丢失,通常不需要使用。

几点实用认知(文档特别强调):

  1. 对象放入队列后是"先 pickle、再由后台线程异步刷入管道"。因此刚 put 完立刻调 empty() 可能仍返回 Trueget_nowait() 可能暂抛 Empty——存在"无限小"的延迟。
  2. 多个进程并发入队时,另一端收到的顺序可能乱序;但同一进程入队的对象彼此必然保序
  3. 若进程在操作 Queue 时被 terminate()/os.kill 杀死,队列数据很可能损坏,波及后续所有使用者。
  4. 子进程入队后若不调用 cancel_join_thread(),它不会在缓冲数据全部刷入管道前终止——于是 join() 该进程可能死锁(除非确认所有入队项都已被消费)。用 Manager 创建的队列无此问题。

multiprocessing 用标准的 queue.Empty / queue.Full 表达超时,二者不在 multiprocessing 命名空间内,需要 from queue import Empty, Full

Connection 对象

Connection 可看作"面向消息的已连接 socket",用于收发可 pickle 对象或字节串,通常由 Pipe 产生。主要方法:

方法 说明
send(obj) 发送对象(须可 pickle);超大 pickle(约 32 MiB 以上,视 OS 而定)可能抛 ValueError
recv() 阻塞接收并重建对象;对端已关闭且无数据时抛 EOFError
fileno() 返回所用 fd/handle
close() 关闭连接(GC 时自动调用)
poll([timeout]) 是否有可读数据:不给 timeout 立即返回;数字秒数;None 无限阻塞
send_bytes(buf[, offset[, size]]) 以完整消息发送 bytes-like 对象的字节数据
recv_bytes([maxlength]) 接收完整字节消息;maxlength 超限时抛 OSError 且连接不可再读
recv_bytes_into(buf[, offset]) 读入 buf 并返回字节数;缓冲过短抛 BufferTooShorte.args[0] 为完整消息

连接对象 3.3 起支持上下文管理协议(__exit__ 自动 close),且连接对象自身也可通过 send/recv 在进程间传递。示例:

>>> from multiprocessing import Pipe
>>> a, b = Pipe()
>>> a.send([1, 'hello', None])
>>> b.recv()
[1, 'hello', None]
>>> b.send_bytes(b'thank you')
>>> a.recv_bytes()
b'thank you'
>>> import array
>>> arr1 = array.array('i', range(5))
>>> arr2 = array.array('i', [0] * 10)
>>> a.send_bytes(arr1)
>>> count = b.recv_bytes_into(arr2)
>>> assert count == len(arr1) * arr1.itemsize
>>> arr2
array('i', [0, 1, 2, 3, 4, 0, 0, 0, 0, 0])

安全警告recv() 会自动反序列化收到的数据——若非可信进程发来,这是一个安全隐患(unpickle 不可信数据有风险)。对于非 Pipe 产生的连接,应先做认证再使用 recv/send(见下节"Authentication keys")。另外,进程在读写管道中途被杀死会使数据损坏、消息边界不可判定。

进程间同步原语

多进程程序通常不像多线程那样依赖同步原语,但 multiprocessing 还是提供了 threading 的完整对应物。最经典的场景是用锁保护共享输出:

from multiprocessing import Process, Lock

def f(l, i):
    l.acquire()
    try:
        print('hello world', i)
    finally:
        l.release()

if __name__ == '__main__':
    lock = Lock()
    for num in range(10):
        Process(target=f, args=(lock, num)).start()

没有锁时,各进程的输出会互相穿插。本包内的同步原语(详见 Lib/multiprocessing/synchronize.pyDoc/library/multiprocessing.rst):

  • Lock()非递归锁,实为工厂函数,返回默认上下文初始化的 multiprocessing.synchronize.Lock。任何进程/线程都可释放它(与 threading.Lock.release 不同,对未锁定的锁 releaseValueError)。支持上下文管理器。acquire(block=True, timeout=None):注意首参名为 block(与 threading 不同);timeout 为负数等价于 0,None 表示无限;block=False 时 timeout 被忽略。
  • RLock()递归锁,必须由获得它的进程/线程释放,重复 acquire 递增递归层级,逐次 release 递减;非持有者或未锁定状态下 release 抛 AssertionError
  • Semaphore([value]):信号量,首个参数同样名为 blockget_value() 获取当前值(macOS 无 sem_getvalue(),可能抛 NotImplementedError)。
  • BoundedSemaphore([value]):有界信号量,locked() 返回是否已锁(二者 3.14 起新增 locked())。
  • Condition([lock]):条件变量,threading.Condition 的别名;lock 应是本模块的 Lock/RLock;3.3 起支持 wait_for()
  • Event()threading.Event 的克隆。
  • Barrier(parties[, action[, timeout]]):栅栏(3.3 起)。
  • 平台注意:macOS 不支持 sem_timedwait,带 timeout 的 acquire() 会用睡眠循环模拟;某些功能依赖宿主系统的共享信号量实现,缺失时 multiprocessing.synchronize 模块被禁用、导入即抛 ImportError(CPython 老 issue 3770)。

3.14 起 LockRLockSemaphoreBoundedSemaphore 均提供 locked() 方法用于查询当前状态。

进程间共享状态

官方建议尽量避免共享状态——这在多进程场景下比多线程更关键。但确有需要时有两种途径。

共享内存:Value / Array

Value/Array 把数据放进共享内存映射:

from multiprocessing import Process, Value, Array

def f(n, a):
    n.value = 3.1415927
    for i in range(len(a)):
        a[i] = -a[i]

if __name__ == '__main__':
    num = Value('d', 0.0)
    arr = Array('i', range(10))
    p = Process(target=f, args=(num, arr))
    p.start()
    p.join()
    print(num.value)
    print(arr[:])

输出:

3.1415927
[0, -1, -2, -3, -4, -5, -6, -7, -8, -9]
  • Value(typecode_or_type, *args, lock=True):返回从共享内存分配的 ctypes 对象(默认包上同步包装,通过 value 属性访问)。typecode_or_type 可为 ctypes 类型或 array 模块风格的单字符 typecode(如 'd' 双精度浮点、'i' 有符号整数)。lock=True 创建新递归锁,传入 Lock/RLock 对象则复用,lock=False 则不做自动保护、非进程安全。注意 lock 是 keyword-only 参数。像 counter.value += 1 这种"读-改-写"并非原子操作,要原子自增必须持有锁:
with counter.get_lock():
    counter.value += 1
  • Array(typecode_or_type, size_or_initializer, *, lock=True):返回共享内存 ctypes 数组。typecode 不支持 'w''c'ctypes.c_char 的别名。size_or_initializer 为整数表示长度(初始清零),为序列则按其初始化并定长。c_char 数组同时有 value(遇到空字节即截断,类似 C 字符串语义)与 raw(对应整个数组的 bytes)两个属性。

需要更灵活的结构时可用 multiprocessing.sharedctypes 模块在共享内存中分配任意 ctypes 对象,其函数有:

  • RawArray(typecode_or_type, size_or_initializer) / RawValue(typecode_or_type, *args):无锁的裸对象,读写可能非原子
  • Array(...) / Value(...):据 lock 参数返回裸对象或同步包装(比同名的模块级 Value/Array 更底层的版本,额外接受 ctx 参数选择上下文,lock/ctx 均为 keyword-only);
  • copy(obj):把 ctypes 对象复制到共享内存;
  • synchronized(obj, lock=None, ctx=None):给 ctypes 对象加进程安全包装,多出 get_obj()(返回被包对象)与 get_lock()(返回同步锁)两个方法;通过包装访问比裸访问慢得多。3.5 起包装对象支持上下文管理器协议。

关键陷阱:指针存在于共享内存中通常指向"某一特定进程地址空间"的地址,在另一进程里解引用很可能会崩溃——不要把指针放进共享内存当作跨进程通信手段。

文档给出 ctypes 语法 ↔ sharedctypes 语法的对照表(设 MyStructctypes.Structure 子类):

ctypes sharedctypes 用类型 sharedctypes 用 typecode
c_double(2.4) RawValue(c_double, 2.4) RawValue('d', 2.4)
MyStruct(4, 6) RawValue(MyStruct, 4, 6)
(c_short * 7)() RawArray(c_short, 7) RawArray('h', 7)
(c_int * 3)(9, 2, 8) RawArray(c_int, (9, 2, 8)) RawArray('i', (9, 2, 8))

综合示例——子进程用锁同时修改数值、数组与自定义结构:

from multiprocessing import Process, Lock
from multiprocessing.sharedctypes import Value, Array
from ctypes import Structure, c_double

class Point(Structure):
    _fields_ = [('x', c_double), ('y', c_double)]

def modify(n, x, s, A):
    n.value **= 2
    x.value **= 2
    s.value = s.value.upper()
    for a in A:
        a.x **= 2
        a.y **= 2

if __name__ == '__main__':
    lock = Lock()
    n = Value('i', 7)
    x = Value(c_double, 1.0/3.0, lock=False)
    s = Array('c', b'hello world', lock=lock)
    A = Array(Point, [(1.875,-6.25), (-5.75,2.0), (2.375,9.5)], lock=lock)
    p = Process(target=modify, args=(n, x, s, A))
    p.start()
    p.join()
    print(n.value)
    print(x.value)
    print(s.value)
    print([(a.x, a.y) for a in A])

输出:

49
0.1111111111111111
HELLO WORLD
[(3.515625, 39.0625), (33.0625, 4.0), (5.640625, 90.25)]

服务器进程:Manager 与代理对象

Manager() 返回的 manager 控制一个服务器进程:它持有共享对象(referent),其它进程通过**代理(proxy)**操纵它们。Manager() 支持共享类型:listdictset(3.14 新增)、NamespaceLockRLockSemaphoreBoundedSemaphoreConditionEventBarrierQueueValueArray。示例:

from multiprocessing import Process, Manager

def f(d, l, s):
    d[1] = '1'
    d['2'] = 2
    d[0.25] = None
    l.reverse()
    s.add('a')
    s.add('b')

if __name__ == '__main__':
    with Manager() as manager:
        d = manager.dict()
        l = manager.list(range(10))
        s = manager.set()
        p = Process(target=f, args=(d, l, s))
        p.start()
        p.join()
        print(d)
        print(l)
        print(s)

输出:

{0.25: None, 1: '1', '2': 2}
[9, 8, 7, 6, 5, 4, 3, 2, 1, 0]
{'a', 'b'}

优点:可支持任意对象类型,且单个 manager 可跨机器、跨网络共享(见下);缺点:比共享内存。manager 进程在自身被 GC 或父进程退出时自动关闭。

代理对象的行为要点

  • 代理引用其它进程中的共享对象(referent);同一 referent 可被多个代理引用;str(proxy) 返回 referent 的表示,而 repr(proxy) 返回代理本身的表示(如 <ListProxy object, typeid 'list' at 0x...>)。
  • 代理可被 pickle,因此可在进程间传递,也可嵌套——manager.dict() 能塞进 manager.list(),嵌套修改会实时同步。
  • 陷阱:若 referent 内部塞的是普通(非代理)的 list/dict,对它们的就地修改不会经 manager 同步(代理无从感知)。要生效需重新赋值以触发代理的 __setitem__
lproxy = manager.list()
lproxy.append({})      # 放入一个普通 dict
d = lproxy[0]
d['a'] = 1             # 就地修改不会同步
lproxy[0] = d          # 重新赋值,触发代理通知
  • 代理不支持按值比较manager.list([1,2,3]) == [1,2,3] 恒为 False;比较前请先取 referent 的拷贝。
  • Namespace:无公开方法、但属性可写的类型;通过代理访问时,以下划线 _ 开头的属性属于代理本身而非 referent。
  • 底层方法:proxy._callmethod(methodname, args, kwds) 在 manager 进程执行 getattr(obj, methodname)(*args, **kwds),远程异常原样重抛、manager 进程内其它异常包装为 RemoteErrorproxy._getvalue() 返回 referent 拷贝。
  • 清理机制:代理用 weakref 回调,被 GC 时向所属 manager 注销自己;referent 没有代理引用后从 manager 进程删除。
  • 代理对象不要跨线程裸用,除非加锁保护(多进程各用同一代理则没有问题)。

自定义与远程 Manager

BaseManager(address=None, authkey=None, serializer='pickle', ctx=None, *, shutdown_timeout=1.0) 是自定义 manager 的基类:address 为监听地址(None 时自动选取);authkey 为连接认证密钥(None 时用 current_process().authkey,否则必须是字节串);serializer'pickle''xmlrpclib'(后者实际用 xmlrpc.client);shutdown_timeout(3.11 起)是 shutdown() 等待 manager 进程退出的超时,超时先 terminate、再超时则 kill。

关键方法:start([initializer[, initargs]]) 启动子进程;get_server() 返回 Server(支持 serve_forever(),且带 address 属性);connect() 让本地对象连接远程 manager 进程;shutdown() 停止进程(仅当用 start() 启动过,可多次调用);类方法 register(typeid[, callable[, proxytype[, exposed[, method_to_typeid[, create_method]]]]]) 注册类型:

  • typeid:共享对象的类型标识字符串;
  • callable:为该类型创建对象的可调用;仅用 connect() 连接、或 create_method=False 时可省略;
  • proxytypeBaseProxy 子类,用于生成代理;None 时自动生成;
  • exposed:允许代理访问的方法名序列;None 时用 proxytype._exposed_,再没有则暴露所有"公开方法"(可调用且不以 _ 开头);
  • method_to_typeid:方法名→typeid 映射,决定哪些方法返回代理(不在映射中的按值拷贝);
  • create_method:是否生成与 typeid 同名的创建方法,默认 True

manager 对象自 3.3 起支持上下文管理协议(__enter__ 启动服务器进程,__exit__shutdown)。示例(自定义 manager 注册类与生成器代理):

from multiprocessing.managers import BaseManager

class MathsClass:
    def add(self, x, y):
        return x + y
    def mul(self, x, y):
        return x * y

class MyManager(BaseManager):
    pass

MyManager.register('Maths', MathsClass)

if __name__ == '__main__':
    with MyManager() as manager:
        maths = manager.Maths()
        print(maths.add(4, 3))         # prints 7
        print(maths.mul(7, 8))         # prints 56

仓库还提供了一份展示 exposed 白名单、自定义 GeneratorProxy_exposed_ = ['__next__'])、注册模块级函数等高级用法的完整可运行示例 Doc/includes/mp_newtype.py

远程 manager:可以把 manager server 跑在一台机器上,让其它机器(防火墙放行时)作为客户端访问。服务器端公开一个共享 Queue 并常驻服务:

>>> from multiprocessing.managers import BaseManager
>>> from queue import Queue
>>> queue = Queue()
>>> class QueueManager(BaseManager): pass
>>> QueueManager.register('get_queue', callable=lambda:queue)
>>> m = QueueManager(address=('', 50000), authkey=b'abracadabra')
>>> s = m.get_server()
>>> s.serve_forever()

客户端(此处假设服务器主机名为 foo.bar.org)接入并读写:

>>> from multiprocessing.managers import BaseManager
>>> class QueueManager(BaseManager): pass
>>> QueueManager.register('get_queue')
>>> m = QueueManager(address=('foo.bar.org', 50000), authkey=b'abracadabra')
>>> m.connect()
>>> queue = m.get_queue()
>>> queue.put('hello')

另一客户端:

>>> m = QueueManager(address=('foo.bar.org', 50000), authkey=b'abracadabra')
>>> m.connect()
>>> queue = m.get_queue()
>>> queue.get()
'hello'

本地进程同样可以经由该模式访问(把 Process 子类里 run()q.put(...) 结果放入由 manager 暴露的队列)。注意 authkey 起着口令作用,且它不会在线路上明文传输(见下文认证机制)。

进程池:Pool

Pool 代表一组工作进程,把任务分发给它们执行。完整方法面(官方"Using a pool of workers"一节示例):

from multiprocessing import Pool, TimeoutError
import time
import os

def f(x):
    return x*x

if __name__ == '__main__':
    # start 4 worker processes
    with Pool(processes=4) as pool:
        # print "[0, 1, 4,..., 81]"
        print(pool.map(f, range(10)))
        # print same numbers in arbitrary order
        for i in pool.imap_unordered(f, range(10)):
            print(i)
        # evaluate "f(20)" asynchronously
        res = pool.apply_async(f, (20,))      # runs in *only* one process
        print(res.get(timeout=1))             # prints "400"
        # evaluate "os.getpid()" asynchronously
        res = pool.apply_async(os.getpid, ()) # runs in *only* one process
        print(res.get(timeout=1))             # prints the PID of that process
        # launching multiple evaluations asynchronously *may* use more processes
        multiple_results = [pool.apply_async(os.getpid, ()) for i in range(4)]
        print([res.get(timeout=1) for res in multiple_results])
        # make a single worker sleep for 10 seconds
        res = pool.apply_async(time.sleep, (10,))
        try:
            print(res.get(timeout=1))
        except TimeoutError:
            print("We lacked patience and got a multiprocessing.TimeoutError")
        print("For the moment, the pool remains available for more work")
    # exiting the 'with'-block has stopped the pool
    print("Now the pool is closed and no longer available")

要点总结:

  • Pool 的方法只能由创建它的进程调用(源码 Lib/multiprocessing/pool.py 中的任务派发与结果回收由内部 task handler/result handler 线程完成,外部并发调用不被支持)。
  • Pool([processes[, initializer[, initargs[, maxtasksperchild[, context]]]]])
    • processes:工作进程数;None 时默认用 os.process_cpu_count()(3.13 起替代 os.cpu_count());
    • initializer/initargs:每个 worker 启动时执行 initializer(*initargs)
    • maxtasksperchild:worker 完成多少个任务后被新进程替换,以便释放资源(Apache/mod_wsgi 中常见的"按工作量轮换进程"模式);默认 None 即 worker 存活到池关闭;
    • context:启动 worker 所用的上下文,通常不必手动指定;
    • maxtasksperchild 传非正整数会抛 ValueError(见 Lib/multiprocessing/pool.py)。
  • 阻塞式:apply(func[, args[, kwds]])(只在一个 worker 中执行并阻塞到结果)、map(func, iterable[, chunksize])(并行版内置 map,只支持单个 iterable;把 iterable 切成块派发;超长 iterable 慎用,可能内存占用高)。
  • 异步式:apply_async / map_async,返回 AsyncResult,支持 callback(结果就绪时回调,单参数)与 error_callback(失败时回调异常实例);回调应尽快返回,否则处理结果的线程会被阻塞。
  • 惰性式:imap(func, iterable, chunksize=1, *, buffersize=None)imap_unordered(...)(结果顺序任意,仅单 worker 时保序)。chunksize 大时显著加速;chunksize=1 时迭代器的 next(timeout) 支持超时抛 multiprocessing.TimeoutErrorbuffersize 用于限制在途任务数以支持"边消费边提交"(建议至少设为进程数)。
  • 多参数式:starmap(func, iterable[, chunksize])starmap_async(...):iterable 元素解包为 func 的实参,[(1,2),(3,4)] 得到 [func(1,2), func(3,4)]
  • 生命周期:close() 禁止新任务、等已提交任务完成后退 worker;terminate() 立即停止 worker(对象被 GC 时也会自动 terminate);join() 等待 worker 退出(先 close/terminate)。3.3 起 Pool 支持上下文管理协议(__exit__terminate)。
  • 资源管理警告:池有内部资源必须妥善管理——推荐用 with 或手动 close/terminate依赖 GC 终结池是不正确的,CPython 不保证池析构会被调用(__del__ 语义)。
  • 关于 REPL:池的功能要求子进程能 import __main__ 模块,因此在交互式解释器里直接定义 f 后建池会失败,报 AttributeError: Can't get attribute 'f' on <module '__main__' ...>(父进程可能需手动中断)。

AsyncResultapply_async/map_async/starmap_async 的返回对象)方法:get([timeout])(超时抛 TimeoutError,远程异常原样重抛)、wait([timeout])ready()successful()(未就绪时调用抛 ValueError,3.7 起由 AssertionError 改为 ValueError)。

池的另一个完整示例见 Doc/includes/mp_pool.py,它同时演示了 apply_asyncimapimap_unorderedmap 的结果顺序差异,以及如何用 try/except 预期捕获 worker 中的 ZeroDivisionError

multiprocessing.dummy 与 ThreadPool

multiprocessing.dummy 复刻了 multiprocessing 的全部 API,但只是 threading 的包装。其 Pool 返回 ThreadPool——Pool 的子类,接口完全一致但使用工作线程

ThreadPool([processes[, initializer[, initargs]]])

Pool 不同,ThreadPool 不接受 maxtasksperchildcontext;资源同样需用 with 或手动 close/terminate 管理。其接口按进程池设计、且自带 AsyncResult 类型,与其它库不互通;官方建议一般场景优先使用 concurrent.futures.ThreadPoolExecutor

connection 模块:Listener、Client 与多路等待

消息传递通常用队列或 Pipe,而 multiprocessing.connection 提供的是针对 socket / Windows 命名管道的高层"面向消息" API,并内置基于 hmac 的摘要认证与多连接轮询。

主要对象:

  • Listener([address[, family[, backlog[, authkey]]]]):封装"正在监听"的绑定 socket 或命名管道。family{'AF_INET','AF_UNIX','AF_PIPE'}(TCP / Unix 域 socket / Windows 命名管道),仅 AF_INET 保证可用;family=None 时按 address 格式推断(address=None 时选被认为最快者);AF_UNIX 且 address 为 None 时,socket 建在 tempfile.mkstemp 创建的私有临时目录。Windows 上监听地址不要用 '0.0.0.0'(不可连接),应改用 '127.0.0.1'backlog(默认 1)传给 socket 的 listen()authkey 为 HMAC 认证密钥;None 则不做认证。方法:accept()(返回 Connection,认证失败抛 AuthenticationError)、close();属性:addresslast_accepted。3.3 起支持上下文管理协议。
  • Client(address[, family[, authkey]]):连接上述监听者并返回 Connectionfamily 通常可从 address 格式推断;带 authkey 时做 HMAC 质询认证。
  • deliver_challenge(connection, authkey) / answer_challenge(connection, authkey):认证握手两端——发起方发随机消息等应答,应答等于用 authkey 计算的消息摘要则回 welcome 消息,否则抛 AuthenticationError
  • wait(object_list, timeout=None):等待列表中任一对象 ready,返回就绪子集。可等待对象:可读的 Connection、已连接且可读的 socket.socketProcess.sentinel。POSIX 上它几乎等价 select.select(list, [], [], timeout),但不会被 EINTR 中断;Windows 上对象须为可等待句柄或带 fileno() 返回 socket/pipe 句柄的对象(注意 pipe/socket handle 本身不是 waitable handle)。

地址格式

  • 'AF_INET'(hostname, port) 元组;
  • 'AF_UNIX':文件系统上的文件名(字符串);
  • 'AF_PIPE'r'\\.\pipe\{PipeName}';连接远程机器 ServerName 上的命名管道用 r'\\{ServerName}\pipe\{PipeName}'任何以两个反斜杠开头的字符串默认视为 AF_PIPE 而非 AF_UNIX

服务端示例(以 b'secret password' 为认证密钥监听并发送数据):

from multiprocessing.connection import Listener
from array import array

address = ('localhost', 6000)     # family is deduced to be 'AF_INET'

with Listener(address, authkey=b'secret password') as listener:
    with listener.accept() as conn:
        print('connection accepted from', listener.last_accepted)
        conn.send([2.25, None, 'junk', float])
        conn.send_bytes(b'hello')
        conn.send_bytes(array('i', [42, 1729]))

客户端:

from multiprocessing.connection import Client
from array import array

address = ('localhost', 6000)

with Client(address, authkey=b'secret password') as conn:
    print(conn.recv())                  # => [2.25, None, 'junk', float]
    print(conn.recv_bytes())            # => 'hello'
    arr = array('i', [0, 0, 0, 0, 0])
    print(conn.recv_bytes_into(arr))    # => 8
    print(arr)                          # => array('i', [42, 1729, 0, 0, 0])

wait 同时监听多个进程的管道:

from multiprocessing import Process, Pipe, current_process
from multiprocessing.connection import wait

def foo(w):
    for i in range(10):
        w.send((i, current_process().name))
    w.close()

if __name__ == '__main__':
    readers = []
    for i in range(4):
        r, w = Pipe(duplex=False)
        readers.append(r)
        p = Process(target=foo, args=(w,))
        p.start()
        # 父进程立刻关闭写端,确保 p 是唯一持有写端句柄的进程,
        # 这样 p 关闭写端后 wait() 能及时报告读端 ready。
        w.close()
    while readers:
        for r in wait(readers):
            try:
                msg = r.recv()
            except EOFError:
                readers.remove(r)
            else:
                print(msg)

Authentication keys 认证密钥机制

由于 Connection.recv 自动反序列化收到的数据,对不可信来源是安全风险,Listener/Client 因而使用 hmac 做摘要认证。认证密钥是字节串,可视为口令:连接建立后两端互相证明"知道同一个密钥",证明过程不发送密钥本身。若要求认证但未指定密钥,默认使用 current_process().authkey(由 os.urandom 生成、被子进程继承,故同程序的进程天然共享一把密钥);也可以用 os.urandom 自行生成合适密钥。

该认证只保护按地址可达的 Listener/Client 连接,不适用于 Pipe() 的匿名管道与 Queue 内部管道。multiprocessing 把"同一用户的所有本地进程"视为可信(多数 OS 上它们本就可互相访问管道 fd);需要同用户进程间隔离的应用必须在操作系统层面隔离(不同账号运行 worker 或沙箱化)。

杂项实用函数与日志

  • active_children():返回当前进程所有存活的子进程列表(附带 join 已完成进程的副作用)。
  • cpu_count():返回系统 CPU 数。不等于当前进程可用的 CPU 数——后者用 os.process_cpu_count()(或 len(os.sched_getaffinity(0)))。无法确定时抛 NotImplementedError。3.13 起返回值可被 -X cpu_count 选项或 PYTHON_CPU_COUNT 环境变量覆盖,因为它只是 os 相关 API 的薄封装。
  • current_process():返回当前进程对应的 Process 对象(类比 threading.current_thread)。
  • parent_process():返回当前进程的父 Process;主进程为 None(3.8 起)。
  • freeze_support():为"冻结"成可执行文件的程序提供支持(经 py2exe、PyInstaller、cx_Freeze 测试)。应紧跟 if __name__ == '__main__' 之后调用;省略时运行冻结产物会抛 RuntimeError。仅 spawn 启动方式、且确为冻结场景时有意义,普通运行下调用无副作用:
from multiprocessing import Process, freeze_support

def f():
    print('hello world!')

if __name__ == '__main__':
    freeze_support()
    Process(target=f).start()
  • 缺失项说明:multiprocessing 没有 threading.active_countthreading.enumeratethreading.settracethreading.setprofilethreading.Timerthreading.local 的对应物。

Loggingget_logger() 返回本模块内部 logger(首次创建时级别 NOTSET、无默认 handler、默认不向 root logger 传播);Windows 上子进程只继承父进程 logger 的级别,其它定制不会继承。log_to_stderr(level=None) 额外为它挂一个输出到 sys.stderr、格式为 '[%(levelname)s/%(processName)s] %(message)s' 的 handler。示例:

>>> import multiprocessing, logging
>>> logger = multiprocessing.log_to_stderr()
>>> logger.setLevel(logging.INFO)
>>> logger.warning('doomed')
[WARNING/MainProcess] doomed
>>> m = multiprocessing.Manager()
[INFO/SyncManager-...] child process calling self.run()
[INFO/SyncManager-...] created temp directory /.../pymp-...
[INFO/SyncManager-...] manager serving at '/.../listener-...'
>>> del m
[INFO/MainProcess] sending shutdown message to manager
[INFO/SyncManager-...] manager exiting with exitcode 0

注意:logging 包不使用进程共享锁,视 handler 类型不同,跨进程日志消息可能互相穿插。

编程指南:避免踩坑的权威准则

文档 Doc/library/multiprocessing.rst 的 "Programming guidelines" 一节总结了多年沉淀的工程经验,分两种场景。

适用于所有启动方式

  1. 尽量避免共享状态:不要在进程间搬移大量数据;通信优先用队列/管道,而非底层同步原语。
  2. 可 pickle 性:传给代理方法的参数必须可 pickle。
  3. 代理的线程安全:不要不加锁地跨线程使用同一代理(不同进程用同一代理没问题)。
  4. 清理僵尸进程:POSIX 上结束但未 join 的进程成为 zombie;每次新起进程或调用 active_children() 都会 join 已完成的进程,is_alive() 也会。但仍建议对启动的进程显式 join()
  5. 能继承就别来回 pickle:spawn/forkserver 下许多 multiprocessing 对象要求可 pickle;但一般应避免把共享对象经管道/队列传给其它进程,而应让需要它的进程从祖先进程继承
  6. 慎用 terminateProcess.terminate() 很可能破坏其正在使用的锁、信号量、管道、队列。只对"不使用共享资源"的进程考虑 terminate。
  7. join 使用队列的进程会死锁:入队进程会等 feeder 线程把缓冲冲刷完才退出(可用 cancel_join_thread() 绕过)。必须确保 join 前队列中所有已入队项最终都会被取出,因为非守护进程退出时会自动 join。典型死锁反例:
from multiprocessing import Process, Queue

def f(q):
    q.put('X' * 1000000)

if __name__ == '__main__':
    queue = Queue()
    p = Process(target=f, args=(queue,))
    p.start()
    p.join()                    # this deadlocks
    obj = queue.get()

修复方法:把最后两行交换(先 getjoin),或干脆删掉 join()

  1. 把资源显式传给子进程:fork 下子进程虽能通过全局变量访问父进程资源,仍应把对象作为构造参数传入——既兼容 Windows 与其它启动方式,又保证子进程存活期间对象不会被父进程 GC(若对象被 GC 时释放资源,这点可能至关重要)。也就是把
lock = Lock()
for i in range(10):
    Process(target=f).start()          # 依赖全局 lock

改写为

lock = Lock()
for i in range(10):
    Process(target=f, args=(lock,)).start()
  1. 别用"类文件对象"替换 sys.stdin:进程引导时会执行 sys.stdin.close() 并把 sys.stdin 重定向为 /dev/nullos.open(os.devnull, os.O_RDONLY)closefd=False),以规避多进程间的坏 fd 冲突。但若你的应用把 sys.stdin 换成带输出缓冲的类文件对象,多个进程并发 close() 可能重复冲刷数据造成损坏。自实现的缓存可用"缓存时记录 pid、pid 变化即丢弃缓存"来保证 fork 安全:
@property
def cache(self):
    pid = os.getpid()
    if pid != self._pid:
        self._pid = pid
        self._cache = []
    return self._cache

仅适用于 spawn / forkserver 的额外约束

  1. 更强的可 pickle 要求Process 的所有参数必须可 pickle;子类覆写 Process.__init__ 时须保证实例在调用 start() 时可被 pickle。
  2. 全局变量语义:子进程读到的全局变量值可能与父进程调用 start() 时不同;但模块级常量没问题。
  3. 主模块的安全导入:必须保证主模块能被全新解释器安全导入且不产生副作用(不直接开新进程),否则 spawn/forkserver 下运行会得到 RuntimeError。反例:
from multiprocessing import Process

def foo():
    print('hello')

p = Process(target=foo)   # 模块级直接 start,spawn 时导入主模块即递归
p.start()

正例是用 if __name__ == '__main__': 保护入口(并可选地加 freeze_support()set_start_method('spawn')):

from multiprocessing import Process, freeze_support, set_start_method

def foo():
    print('hello')

if __name__ == '__main__':
    freeze_support()
    set_start_method('spawn')
    p = Process(target=foo)
    p.start()

这样新解释器能安全导入该模块后再执行 foo()在主模块中创建 pool 或 manager 时也有同样限制

综合示例:三种典型模式

仓库在 Doc/includes/ 提供了三个可直接运行的示例,分别对应三类经典场景:

  1. mp_newtype.py:演示自定义 Manager 与代理。注册 Foo 类并分别以全公开方法、exposed=('g', '_h') 白名单两种方式暴露;为生成器 baz 定制 GeneratorProxy_exposed_ = ['__next__'],用 _callmethod 转发);还演示了把 operator 模块注册后经代理调用其公开函数。

  2. mp_pool.py:演示 Pool 的四种任务形态——apply_async(结果保持提交顺序)、imap(顺序迭代)、imap_unordered(乱序到达)、map(阻塞至全部完成),并用 TimeError/ZeroDivisionError 分别演练了超时与 worker 异常在 applymapimapIMapIterator.next() 中的传播。

  3. mp_workers.py:演示用队列给一组 worker 进程喂任务并回收结果的经典生产者-消费者架构——主进程把 (函数, 参数) 任务放入 task_queue,4 个 worker 循环 input.get(),以 'STOP' 哨兵对象收尾退出,结果经 done_queue 汇总;注意 worker 的入队进程需在退出前由主进程取尽结果,避免 join 死锁。

小结

multiprocessing 的价值在于用 threading 式 API 换取真正的多进程并行,从而绕开 GIL、吃满多核,同时把"数据并行(Pool)、消息传递(队列/管道/连接)、共享状态(共享内存/Manager 代理)、网络化共享(远程 Manager 与 Listener/Client)"整合为一套自洽的编程模型。选用要点可归纳为:

  • 启动方式:跨平台可移植优先选 spawn;POSIX 纯计算场景默认已是性能与安全兼顾的 forkserver;只有确定无多线程且需要"继承式"语义时才显式选 fork(3.14 起不再是任何平台的默认)。
  • 通信优先于共享:能用 Queue/Pipe 消息传递,就不要手工加锁维护共享可变状态。
  • 必须共享时:标量/同构数组用 Value/Array/sharedctypes(性能好);异构、需灵活嵌套或跨机器访问时用 Manager 代理(慢但通用)。
  • 进程池:永远用 with Pool(...) 或显式 close/terminate 管理资源;大任务列表用 imap/imap_unordered 配较大 chunksizebuffersize 控流。
  • 守住两条红线:目标函数放可导入模块、入口用 if __name__ == '__main__' 保护;join 任何入过队的进程前先确保队列被消费干净。
登录后查看全文
热门项目推荐
相关项目推荐