JEPA4Japan · 教程

线程、进程与并行工作

2,290字 7分钟阅读 #Python

为 I/O 密集型和 CPU 密集型任务匹配安全的并发原语。

课程进度 课程大纲 已发布 24/24 课

重叠并不等于同时执行

并发是指多个任务在同一段时间内都处于进行状态。并行是指多个任务真正地在同一时刻运行,通常运行在不同的 CPU 核心上。

  1. 顺序执行 做完一道菜,再开始下一道
  2. 并发 一道菜等待时切换到另一道
  3. 并行 两位厨师同时工作
生命周期相互重叠与同时执行回答的是不同的问题。

会造成阻塞的文件或网络工作通常是 I/O 密集型的:它的大部分时间都花在等待上。耗时较长的纯 Python 计算则是 CPU 密集型的。线程通常有助于让阻塞等待相互重叠;进程可以让独立的 CPU 任务在不同核心上运行。调度和通信都需要时间,因此在宣称获得性能提升之前,请先进行测量。

线程共享同一个房间

同一个进程中的线程共享相同的地址空间、全局变量和堆对象。每个线程都有自己的调用栈,但两个引用可以指向同一个可变对象。

在标准的 CPython 3.12 构建中,GIL 只允许一个进程内的一个线程在同一时刻执行 Python 字节码。CPython 会在许多阻塞式 I/O 等待期间释放它,一些原生扩展也会释放它。因此,线程可以重叠执行 I/O,但多个线程通常无法让纯 Python 的 CPU 字节码跨多个核心运行。其他解释器或特殊构建可能有所不同。

  1. 一个进程 一个共享的内存房间
  2. 线程 共享房间里的对象
  3. GIL 一次只有一个线程运行 Python 字节码
  4. I/O 等待 另一个线程可以继续执行
GIL 是解释器层面的边界,并不是保护你的业务规则的锁。

GIL 不会保护一个包含多个步骤的不变量。请使用 Lock 包住完整的检查和修改操作:

from concurrent.futures import ThreadPoolExecutor
from threading import Lock

balance = 3
balance_lock = Lock()


def spend_one():
    global balance
    with balance_lock:
        if balance > 0:
            balance -= 1


with ThreadPoolExecutor(max_workers=4) as pool:
    list(pool.map(lambda _: spend_one(), range(10)))

print(balance)
0

这个锁保护了“余额永远不能支出到零以下”这条规则。请让临界区尽可能小。更好的做法是,在不必共享数据时,让工作线程返回值,再由一个所有者统一合并这些值。

Executor 管理工作线程;Future 携带结果

ThreadPoolExecutor 会重复使用一组数量受限的线程。即使工作线程完成任务的顺序不同,map() 仍会按照输入顺序返回结果。submit() 会返回一个 Future:

from concurrent.futures import ThreadPoolExecutor


def parse_minutes(text):
    minutes = int(text)
    if minutes < 0:
        raise ValueError("negative minutes")
    return minutes


with ThreadPoolExecutor(max_workers=1) as pool:
    future = pool.submit(parse_minutes, "-2")
    try:
        future.result()
    except ValueError as error:
        print(f"Worker failed: {error}")
Worker failed: negative minutes

result() 会返回结果值,或者在负责协调的线程中重新抛出工作线程产生的异常。务必检查每一个 Future。as_completed() 会按完成顺序提供结果,而这个顺序本来就是不稳定的;请附上标识符,并在公开输出之前排序。

进程使用各自独立的房间

工作进程拥有各自独立的解释器和内存。在标准 CPython 中,它们可以在不同核心上执行彼此独立的 Python CPU 任务。输入、可调用对象和结果必须通过序列化跨越进程边界,通常使用 pickle。

  1. 进程 A 拥有自己的内存
  2. 序列化的数据 在房间之间传递
  3. 进程 B 可以使用另一个核心
进程隔离使 CPU 并行成为可能,但也会增加启动、内存和传输成本。

进程任务应该是可以导入的顶层函数。参数和结果应当大小合理并且可序列化;打开的文件、锁、活动连接、lambda 和嵌套函数都不适合作为任务值。请将协调逻辑放在可移植的启动保护条件之后:

def main():
    run_process_pool()


if __name__ == "__main__":
    main()

如果没有这个保护条件,生成的子进程可能会导入该模块,然后递归地创建另一个进程池。

限制与停止都有明确的边界

增加工作线程或进程的数量,可能会加剧磁盘争用、增加内存使用、引发服务拒绝,并提高序列化开销。请根据实际容量选择 max_workers。对于大量微小的 CPU 任务,请批量处理;对于无穷的数据流,请使用有限批次或容量受限的生产者队列。仅仅固定工作线程或进程的数量,并不能限制已经提交的 Future 数量。

停止相关 API 的实际含义,比它们的名称看起来更加有限:

  • future.result(timeout=1) 限制的是调用者等待的时长;任务可能仍会继续运行。
  • future.cancel() 只有在调用开始之前才能成功。
  • shutdown(wait=True) 会停止接受新任务,并等待正在运行和待处理的任务。
  • shutdown(wait=False) 会更早返回,但调用仍会继续,而且 Python 依然会在进程退出前等待它们完成。
  • cancel_futures=True 会取消待处理的 Future,而不会取消已经开始运行的 Future。

请使用原生的 I/O 超时机制和协作式停止信号。Python 无法安全地强制任意一个正在运行的线程退出执行。

小项目:本地阅读摘要

创建 reading_digest.py。它完全离线,并且只使用标准库。一个数量受限的线程池读取临时文件;一个数量受限的进程池执行彼此独立的文本分析。排序会让输出保持确定性。

from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor
from pathlib import Path
from tempfile import TemporaryDirectory


def load_article(path):
    return path.name, path.read_text(encoding="utf-8")


def analyze_article(record):
    filename, text = record
    words = text.split()
    mentions = sum(
        word.strip(".,:;!?()").casefold() == "python" for word in words
    )
    return filename, len(words), mentions


def build_digest(paths):
    paths = list(paths)
    if not paths:
        return []
    with ThreadPoolExecutor(max_workers=min(2, len(paths))) as pool:
        loaded = list(pool.map(load_article, paths))
    with ProcessPoolExecutor(max_workers=min(2, len(loaded))) as pool:
        analyzed = list(pool.map(analyze_article, loaded))
    return sorted(analyzed)


def main():
    with TemporaryDirectory() as temporary:
        directory = Path(temporary)
        samples = {
            "a.txt": "Python waits. Python reads files.",
            "b.txt": "Processes analyze independent Python text.",
        }
        paths = []
        for name, text in samples.items():
            path = directory / name
            path.write_text(text, encoding="utf-8")
            paths.append(path)

        items = build_digest(reversed(paths))
        print("Reading digest")
        for filename, words, mentions in items:
            print(f"- {filename}: {words} words, python={mentions}")

        with ThreadPoolExecutor(max_workers=1) as pool:
            missing = pool.submit(load_article, directory / "missing.txt")
            try:
                missing.result()
            except FileNotFoundError:
                print("Missing file surfaced: True")


if __name__ == "__main__":
    main()

将它保存为文件,然后运行 python3 reading_digest.py:

Reading digest
- a.txt: 5 words, python=2
- b.txt: 5 words, python=1
Missing file surfaced: True

这个示例有意设计得过小,因此无法证明存在性能提升。它证明的是进程启动、序列化、异常传递、工作线程或进程数量限制,以及稳定输出。对于很小的输入,顺序分析器可能反而更快。

三个小任务与第 23 章检查清单

  1. 消除共享。 将使用锁保护余额的方式改为让工作线程返回各自的局部计数,然后集中合并它们。
  2. 保持顺序。 使用 submit() 和 as_completed(),保存输入索引,然后按照该索引输出。
  3. 如实测量。 使用更大的本地文本比较顺序分析和进程分析,并记录工作进程数量和数据大小。

当你能够做到以下几点时,就已经为第 23 章做好准备:

  • 你能区分并发与并行;
  • 你能说明线程共享哪些内容,以及标准 CPython 3.12 的 GIL 边界;
  • 你能使用锁保护一个复合不变量;
  • 你会检查每一个 Future,并让其中的异常显现出来;
  • 你会在 main 保护条件之后使用顶层的、可序列化的进程任务;
  • 你会根据资源限制来约束工作线程或进程的数量以及提交的任务数量;
  • 你知道 timeout、cancel 和 shutdown 都不会终止任意一个正在运行的任务;
  • 你可以运行这个摘要程序,并复现经过排序的输出。

接下来,一个事件循环将使用协程、任务、截止时间和协作式取消来协调许多 I/O 等待。