Python 基础体系 · 第 55/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。

Python 线程:生命周期、锁、条件变量、竞态和死锁诊断

线程(thread)是进程内的一条执行路径。多个线程属于同一个进程,因此通常共享堆内存、模块全局变量、文件描述符和其他进程级资源;线程之间的主要隔离边界是各自的调用栈、寄存器状态以及线程局部数据。Python 的 threading 模块建立在更底层的 _thread 模块之上,适合处理需要等待网络、文件、数据库或其他外部资源的并发任务。(docs.python.org)

线程带来的核心问题不是“如何同时运行几段代码”,而是:

  1. 一个线程何时真正开始、结束以及退出;
  2. 多个线程如何安全地访问共享状态;
  3. 一个线程如何等待另一个线程改变状态;
  4. 为什么代码有时正确、有时错误;
  5. 当线程不再推进时,如何判断是慢、阻塞、活锁还是死锁。

一、先区分并发、并行与线程

并发描述的是多个任务在时间上存在重叠。任务可以交替执行,也可以真正同时执行。

并行描述的是多个任务在同一时刻使用多个执行资源运行。并行通常需要多个 CPU 核心,或者由底层 C、Rust、Fortran 等代码释放解释器限制后使用多个核心。

线程通常提供的是进程内并发模型:

进程 P
├── 主线程
├── 工作线程 A
├── 工作线程 B
└── 工作线程 C

所有线程共享进程地址空间:

线程 A ─┐
线程 B ─┼── 共享堆、全局变量、文件描述符、连接池
线程 C ─┘

共享内存使线程之间传递数据很便宜,但也使“读到什么、改成什么、改动是否被覆盖”成为程序正确性的一部分。

1. GIL 不等于线程安全

在常见的 CPython 构建中,GIL(Global Interpreter Lock,全局解释器锁)限制同一时刻只有一个线程执行 Python 字节码。它会限制纯 Python CPU 密集型代码从线程中获得的并行收益,但不会自动把一个业务操作变成不可分割的事务。(docs.python.org)

例如:

counter += 1

从业务语义看,它是一次“加一”;从执行过程看,至少可以抽象为:

读取 counter
计算 counter + 1
写回 counter

如果两个线程交错执行:

初始 counter = 0

线程 A:读取 0
线程 B:读取 0
线程 A:计算 1
线程 B:计算 1
线程 A:写回 1
线程 B:写回 1

最终结果是 1,而不是期望的 2

GIL 可能让某些底层操作在某个 CPython 版本中表现出“通常不会被打断”的现象,但这不是业务层面的复合操作保证,也不应作为同步协议。代码不能依赖:

  • 某个字典操作在当前实现中恰好是原子的;
  • 某个字节码序列当前没有发生线程切换;
  • GIL 会保护多个共享对象之间的一致性;
  • 未来更换 Python 构建、解释器或扩展后仍然保持同样行为。

从 Python 3.13 开始,CPython 提供了可以禁用 GIL 的自由线程构建,但这类构建不是默认构建;Python 3.14 的线程代码仍必须明确建立自己的互斥边界。(docs.python.org)


二、线程生命周期:从对象到终止

一个 Thread 对象的生命周期可以分为以下阶段:

stateDiagram-v2
    [*] --> New: Thread(...)
    New --> Started: start()
    Started --> Running: 调度并进入 run()
    Running --> Waiting: 阻塞 I/O / Lock / Condition / Event
    Waiting --> Running: 条件满足或超时
    Running --> Terminated: run() 正常返回
    Running --> Terminated: run() 未捕获异常
    Terminated --> Joined: 其他线程 join()
    Terminated --> [*]

需要注意,Thread 对象和操作系统线程不是同一个概念:

  • 创建 Thread(...) 只创建了线程对象;
  • 调用 start() 才会安排新的执行线程运行;
  • 直接调用 run() 不会创建新线程,只是在当前线程同步执行;
  • run() 返回或抛出未捕获异常后,线程结束;
  • join() 是等待线程结束,不是启动线程;
  • 同一个 Thread 对象最多只能调用一次 start()。(docs.python.org)

1. start()run() 的区别

from threading import Thread
import threading


def work():
    print("current:", threading.current_thread().name)


t = Thread(target=work, name="worker")

t.run()    # 在当前线程执行
print("after run")

t2 = Thread(target=work, name="worker-2")
t2.start() # 创建并启动新线程
t2.join()

典型输出类似:

current: MainThread
after run
current: worker-2

第一处 run() 没有并发发生。第二处 start() 才让 work() 在名为 worker-2 的新线程中执行。

2. join() 的语义

from threading import Thread
import time


def work():
    time.sleep(1)
    print("worker finished")


t = Thread(target=work, name="worker")
t.start()

print("main waits")
t.join()
print("main continues")

执行路径是:

主线程:start()
工作线程:sleep(1)
主线程:join(),等待
工作线程:打印并返回
主线程:join() 返回,继续执行

带超时的 join() 不会通过返回值告诉调用者是否真的结束,因为 join() 始终返回 None。应在之后调用 is_alive()

t.join(timeout=0.2)

if t.is_alive():
    print("worker did not finish within 0.2 seconds")
else:
    print("worker finished")

join() 只能等待结束,不能强制终止线程。Python 的 Thread 没有通用的安全接口来杀死、暂停、恢复或中断任意线程。(docs.python.org)

3. 用 Event 实现可控退出

线程不能被安全地“外部杀死”,因此长时间运行的线程应主动检查停止信号:

import threading
import time


stop_event = threading.Event()


def worker():
    while not stop_event.is_set():
        print("processing")
        # wait() 比 sleep() 更适合,因为收到停止信号后可以提前返回
        stop_event.wait(0.5)

    print("worker exits gracefully")


t = threading.Thread(target=worker, name="worker")
t.start()

time.sleep(1.2)
stop_event.set()
t.join(timeout=2)

if t.is_alive():
    raise RuntimeError("worker failed to stop")

Event 内部维护一个布尔标志:

  • 初始为 False
  • set() 将其变为 True,并唤醒等待者;
  • clear() 恢复为 False
  • wait() 在标志为 False 时阻塞;
  • wait(timeout) 在超时后返回 False,标志被设置时返回 True。(docs.python.org)

Event 适合表达“是否应该停止”“服务是否已初始化”“配置是否已就绪”这类单向或广播式状态。它不适合直接表示“队列中当前有多少个元素”。

4. 守护线程不是资源管理方案

守护线程 daemon=True 的含义是:当程序中只剩守护线程时,解释器可以退出。解释器退出时,守护线程可能被直接停止,因此文件、数据库事务、网络连接等资源不一定能正常清理。(docs.python.org)

t = Thread(target=worker, daemon=True)

这适合:

  • 不重要的后台观察任务;
  • 进程退出时可以丢弃的缓存刷新;
  • 没有必须提交的事务或必须关闭的资源。

这不适合:

  • 数据落盘;
  • 消息确认;
  • 数据库事务提交;
  • 释放锁、连接、文件等必须清理的资源。

生产代码通常应使用非守护线程、停止事件和 join() 形成明确的关闭协议。


三、锁:把共享状态划定为临界区

是一种同步原语,用来限制同一时刻能够进入某段代码的线程数量。

临界区是访问共享状态、且必须保持一致性的代码区域。例如:

balance = balance - amount

如果多个线程可以同时执行这段操作,就必须保证“读取余额、检查余额、写回余额”作为一个整体完成。

1. Lock 的状态机

普通锁只有两个状态:

unlocked --acquire()--> locked
locked   --release()--> unlocked

当锁已被占用时,其他线程调用 acquire() 会阻塞,直到某个线程释放锁。普通 Lock 不记录“所有者”;根据 Python 文档,获取锁的线程和释放锁的线程可以不同,但这种跨线程释放通常会让代码协议难以理解,除非有非常明确的设计。(docs.python.org)

import threading

lock = threading.Lock()
counter = 0


def increment():
    global counter
    with lock:
        counter += 1

with lock: 等价于:

lock.acquire()
try:
    counter += 1
finally:
    lock.release()

finally 的意义在于,即使临界区抛出异常,也必须释放锁;否则其他线程可能永久等待。Python 的锁、可重入锁、条件变量、信号量和有界信号量都支持上下文管理器协议。(docs.python.org)

2. 完整算例:修复计数器竞态

下面的代码故意在读写之间插入让出执行的机会,以放大竞态:

import threading
import time


counter = 0
N = 10_000
THREADS = 4


def unsafe_increment():
    global counter

    for _ in range(N):
        value = counter
        time.sleep(0)  # 仅用于放大交错,不是同步手段
        counter = value + 1


workers = [
    threading.Thread(target=unsafe_increment)
    for _ in range(THREADS)
]

for t in workers:
    t.start()

for t in workers:
    t.join()

print(counter)
print("expected:", N * THREADS)

理论上期望结果为:

expected: 40000

但实际结果可能小于 40000。原因不是某个线程“忘了加一”,而是多个线程基于同一个旧值计算并互相覆盖。

修复版本如下:

import threading
import time


counter = 0
lock = threading.Lock()
N = 10_000
THREADS = 4


def safe_increment():
    global counter

    for _ in range(N):
        with lock:
            value = counter
            time.sleep(0)  # 仍然安全,但会延长持锁时间
            counter = value + 1


workers = [
    threading.Thread(target=safe_increment)
    for _ in range(THREADS)
]

for t in workers:
    t.start()

for t in workers:
    t.join()

print(counter)
print("expected:", N * THREADS)

这里正确性的关键不是 sleep(0),而是:

  1. 读取 counter 前获取锁;
  2. 计算和写回仍在锁内;
  3. 每次获取都有对应释放;
  4. 主线程 join() 所有工作线程后再读取最终结果。

从性能角度看,应尽量缩短临界区:

def better_increment():
    global counter

    for _ in range(N):
        with lock:
            counter += 1
        # 慢操作放在锁外

但“缩短锁范围”不能破坏不变量。如果业务要求“检查余额后扣款”不可分割,那么检查和扣款必须处于同一个临界区。

3. LockRLock

RLock 是可重入锁。同一个线程可以多次获取它,但必须进行相同次数的释放。它额外维护当前所有者和递归层数。(docs.python.org)

import threading


lock = threading.RLock()


def outer():
    with lock:
        inner()


def inner():
    with lock:
        print("re-entered safely")


outer()

如果这里使用普通 Lockouter() 获取锁后调用 inner(),而 inner() 再次获取同一把锁,就会等待自己释放锁,形成自死锁。

RLock 解决的是“同一线程嵌套获取同一把锁”的问题,不是通用的死锁修复工具。若代码层级过深、多个方法隐式共享同一把锁,使用 RLock 可能掩盖糟糕的锁边界设计。


四、条件变量:等待“状态成立”,而不是等待“有人通知”

条件变量(condition variable)是“锁 + 等待队列”的同步原语。它解决的问题是:

当前线程需要等待某个共享状态满足条件,另一个线程改变状态后负责通知。

例如生产者—消费者模型中,消费者等待“队列非空”,生产者在放入元素后通知消费者。

条件变量本身不保存业务状态。业务状态必须由程序自己维护:

共享状态:items
同步工具:Condition
业务条件:len(items) > 0

1. wait() 的真实过程

调用:

with condition:
    condition.wait()

不是简单地“睡眠”。它的过程是:

  1. 当前线程已经持有条件变量关联的锁;
  2. 当前线程进入等待队列;
  3. 原子地释放底层锁;
  4. 当前线程阻塞;
  5. 其他线程修改共享状态并调用 notify()notify_all()
  6. 当前线程被唤醒;
  7. 当前线程重新获取底层锁;
  8. wait() 返回。

Python 文档明确要求:调用 wait() 前必须持有底层锁;被唤醒后,线程要先重新获取锁,才会从 wait() 返回。(docs.python.org)

2. 为什么必须使用 while

错误写法:

with cv:
    if not items:
        cv.wait()
    item = items.pop(0)

正确写法:

with cv:
    while not items:
        cv.wait()
    item = items.pop(0)

原因有两个:

原因一:通知不等于条件成立

假设消费者 A 和消费者 B 都等待队列非空:

队列为空
A 等待
B 等待

生产者放入一个元素
生产者 notify_all()

A 被唤醒
B 被唤醒

A 重新获得锁并取走唯一元素后,B 才能重新获得锁。此时 B 被通知过,但队列已经为空。如果使用 if,B 会直接执行 pop(),导致异常;使用 while,B 会重新检查条件并继续等待。

原因二:等待可能因超时返回

wait(timeout=...) 超时后也会返回。返回并不代表业务条件成立,只代表“等待结束”。

wait_for(predicate) 可以自动重复检查谓词:

with cv:
    ready = cv.wait_for(lambda: len(items) > 0, timeout=2)

    if not ready:
        raise TimeoutError("no item available")

    item = items.pop(0)

wait_for() 的谓词会在持有锁时执行,本质上近似于:

while not predicate():
    cv.wait()

(docs.python.org)

3. 完整生产者—消费者实现

from collections import deque
import threading
import time


class BoundedBuffer:
    def __init__(self, capacity: int):
        if capacity <= 0:
            raise ValueError("capacity must be positive")

        self._items = deque()
        self._capacity = capacity
        self._cv = threading.Condition()

    def put(self, item) -> None:
        with self._cv:
            self._cv.wait_for(
                lambda: len(self._items) < self._capacity
            )
            self._items.append(item)

            # 队列从可能为空变成非空,至少一个消费者可能有用
            self._cv.notify()

    def get(self, timeout: float | None = None):
        with self._cv:
            available = self._cv.wait_for(
                lambda: bool(self._items),
                timeout=timeout,
            )

            if not available:
                raise TimeoutError("buffer is empty")

            item = self._items.popleft()

            # 队列从满变成未满,至少一个生产者可能有用
            self._cv.notify()
            return item


buffer = BoundedBuffer(capacity=2)
stop = object()


def producer():
    for i in range(5):
        buffer.put(i)
        print("produced", i)
        time.sleep(0.05)

    buffer.put(stop)


def consumer():
    while True:
        item = buffer.get(timeout=2)
        if item is stop:
            print("consumer exits")
            return

        print("consumed", item)
        time.sleep(0.1)


p = threading.Thread(target=producer, name="producer")
c = threading.Thread(target=consumer, name="consumer")

p.start()
c.start()

p.join()
c.join()

这里有两个不变量:

0 <= len(_items) <= capacity

put() 只有在 len(_items) < capacity 时才追加元素;get() 只有在队列非空时才删除元素。检查条件和修改队列必须使用同一把锁,否则即使两个方法各自“看起来正确”,组合起来仍然会发生竞态。

4. notify()notify_all()

notify() 默认唤醒一个等待线程;notify_all() 唤醒所有等待线程。被唤醒的线程不会立即执行,因为通知者仍然持有底层锁;通知者退出 with 块释放锁后,被唤醒线程才有机会继续。(docs.python.org)

一般可以这样判断:

  • 一个状态变化最多只对一个线程有用:使用 notify()
  • 状态变化可能让多个线程都满足条件:使用 notify_all()
  • 无法准确判断,或者谓词复杂:使用 notify_all(),让每个线程自行通过 whilewait_for() 再检查。

错误模式是先释放锁,再通知:

# 不推荐
with cv:
    items.append(item)

cv.notify()  # 没有持有 cv,直接 RuntimeError

notify() 必须在持有条件变量底层锁时调用。(docs.python.org)


五、竞态条件:结果依赖不可控的执行顺序

竞态条件(race condition)是指程序结果依赖多个并发执行步骤的相对顺序,而这个顺序没有被程序显式约束。

可以将一个共享状态操作表示为:

S0 --线程 A 的操作--> S1
S0 --线程 B 的操作--> S2

如果 A、B 的操作可以安全交换,称它们在该状态下具有可交换性:

A(B(S0)) == B(A(S0))

如果不相等,执行顺序就会影响结果,必须同步。

1. 典型竞态:检查—执行分离

if key not in cache:
    cache[key] = load_from_database(key)

两个线程可能同时观察到 key not in cache,然后同时加载并写入:

线程 A:检查,未命中
线程 B:检查,未命中
线程 A:加载数据库
线程 B:加载数据库
线程 A:写缓存
线程 B:写缓存

即使最终缓存值没有损坏,也可能产生:

  • 重复昂贵查询;
  • 重复创建资源;
  • 两次发送通知;
  • 两次扣款或提交;
  • 先完成的结果被后完成的旧结果覆盖。

如果业务要求只允许一次初始化,检查和写入必须放在同一个同步协议中:

import threading


class LazyValue:
    def __init__(self):
        self._lock = threading.Lock()
        self._loaded = False
        self._value = None

    def get(self):
        with self._lock:
            if not self._loaded:
                self._value = self._load()
                self._loaded = True
            return self._value

    def _load(self):
        return "loaded"

这里的锁保证了初始化只有一个执行者,但也意味着 _load() 执行期间其他线程会等待。如果加载过程很慢,应根据业务设计更细致的状态机,例如:

EMPTY
  │ 一个线程开始加载
  ▼
LOADING
  │ 成功
  ▼
READY

LOADING
  │ 失败
  ▼
FAILED

只用一个布尔值无法区分“尚未开始”“正在加载”和“加载失败”。

2. 共享可变状态与不可变消息

降低竞态的一种方式不是增加更多锁,而是减少共享可变对象:

# 多个线程共同修改同一个列表
shared = []

# 更清晰的方式:线程产生结果,通过 Queue 交给单一消费者
from queue import Queue

results = Queue()

queue.Queue 提供线程安全的数据交换接口,适合将“共享状态修改”转变为“消息传递”。(docs.python.org)

这并不意味着队列自动解决所有问题。消费者仍需处理:

  • 生产者异常;
  • 结束标记;
  • 队列阻塞;
  • 消费失败后的重试;
  • task_done()join() 的配对;
  • 关闭阶段仍有任务未处理。

六、死锁:所有线程都在等待系统无法满足的条件

死锁(deadlock)是一个并发系统进入这样的状态:

每个相关线程都在等待某个事件,而这个事件只能由另一个同样处于等待状态的线程触发。

最经典的两个锁死锁:

import threading
import time


lock_a = threading.Lock()
lock_b = threading.Lock()


def worker_1():
    with lock_a:
        print("worker_1 owns A")
        time.sleep(0.1)

        with lock_b:
            print("worker_1 owns B")


def worker_2():
    with lock_b:
        print("worker_2 owns B")
        time.sleep(0.1)

        with lock_a:
            print("worker_2 owns A")


t1 = threading.Thread(target=worker_1)
t2 = threading.Thread(target=worker_2)

t1.start()
t2.start()

t1.join()
t2.join()  # 可能永久等待

可能的执行路径:

T1:获取 A
T2:获取 B
T1:等待 B
T2:等待 A

此时:

T1 等待 T2 释放 B
T2 等待 T1 释放 A

但两个线程都不会释放自己已经持有的锁,因为释放动作位于第二把锁获取之后。

1. 用等待图形式化死锁

将线程和锁建模为有向图:

T1 ──等待──> B
B  ──持有──> T2

T2 ──等待──> A
A  ──持有──> T1

存在环:

T1 → B → T2 → A → T1

等待图中的环是死锁的重要证据。

2. 死锁的四个必要条件

经典死锁理论通常使用四个必要条件:

  1. 互斥:资源同一时刻只能被一个线程占用;
  2. 占有且等待:线程持有已有资源,同时等待其他资源;
  3. 不可剥夺:资源不能被外部强制收回;
  4. 循环等待:线程之间形成环形等待。

普通锁天然提供互斥和不可剥夺;嵌套获取多个锁时容易产生“占有且等待”和循环等待。

3. 通过全局锁顺序破坏循环等待

让所有代码遵守同一个锁顺序:

def worker_1():
    with lock_a:
        with lock_b:
            pass


def worker_2():
    with lock_a:
        with lock_b:
            pass

此时:

所有线程都先请求 A,再请求 B

线程 T2 可能等待 A,但不会在持有 B 的同时等待 A,因此不会形成上述环。

如果锁的业务顺序无法固定,可以按稳定键排序:

def lock_two(first_name, first_lock, second_name, second_lock):
    locks = sorted(
        [(first_name, first_lock), (second_name, second_lock)],
        key=lambda pair: pair[0],
    )

    with locks[0][1]:
        with locks[1][1]:
            pass

实际项目中更推荐封装成清晰的资源管理抽象,避免调用者自行组合锁。

4. 自死锁与锁泄漏

自死锁:

lock = threading.Lock()


def broken():
    with lock:
        with lock:  # 当前线程等待自己释放 lock
            pass

锁泄漏:

lock.acquire()
do_work()  # 如果这里抛异常,release 永远不会执行
lock.release()

正确写法:

with lock:
    do_work()

对于 RLock,另一种常见错误是获取次数和释放次数不匹配:

rlock.acquire()
rlock.acquire()
rlock.release()
# 仍然被当前线程持有

RLock 必须每获取一次就释放一次;未配对的释放可能让其他线程永久等待。(docs.python.org)


七、死锁与其他“卡住”现象的区别

线程长时间没有完成,不一定是死锁。

1. 阻塞 I/O

线程可能在:

  • 网络连接;
  • socket 读取;
  • 文件读取;
  • 数据库查询;
  • 外部命令;
  • DNS 解析。

这种情况下,线程通常在等待外部系统,不一定存在环形依赖。应检查底层 API 是否设置了连接、读取和整体操作超时。

2. 锁竞争

多个线程都等待一把锁,但持锁线程仍在执行:

T1:持有 lock,执行慢查询
T2:等待 lock
T3:等待 lock

这可能是性能问题,也可能最终导致超时,但还不是死锁。诊断重点是:持锁线程为什么没有尽快退出临界区。

3. 活锁

活锁中线程没有阻塞,而是不断重试、让步或回滚,却始终没有完成:

T1:发现冲突,释放
T2:发现冲突,释放
T1:再次获取,发现冲突,释放
T2:再次获取,发现冲突,释放

日志表现为线程仍然活跃,但业务进度为零。需要加入退避、随机抖动或重新设计冲突协议。

4. 饥饿

某个线程长期得不到资源,但其他线程仍然不断完成。普通锁唤醒哪个等待线程并没有定义的公平顺序,不能依赖严格 FIFO 调度。(docs.python.org)


八、死锁诊断:先获取现场,再解释现场

死锁诊断的关键不是盲目增加日志,而是回答三个问题:

  1. 哪些线程仍然存活;
  2. 每个线程当前停在哪里;
  3. 哪些线程正在等待哪些锁、条件或外部资源。

1. 使用 faulthandler 导出所有线程栈

可以在程序中注册信号:

import faulthandler
import signal

faulthandler.register(signal.SIGUSR1)

在 Unix 系统中,进程卡住时执行:

kill -USR1 <pid>

Python 会将当前线程的栈信息输出到标准错误。若程序启动时无法修改代码,也可以使用:

python -X faulthandler app.py

或者在代码中:

import faulthandler

faulthandler.dump_traceback_later(
    30,
    repeat=True,
)

这表示如果主线程或进程在指定时间内没有完成某个阶段,就周期性输出线程栈。诊断结束后应取消:

faulthandler.cancel_dump_traceback_later()

线程栈能回答“卡在哪里”,但不一定直接显示“等待的是哪一把锁”。因此应结合锁获取日志。

2. 使用 sys._current_frames() 主动打印现场

import sys
import threading
import traceback


def dump_thread_stacks():
    frames = sys._current_frames()

    for thread in threading.enumerate():
        frame = frames.get(thread.ident)

        print(
            f"\nthread name={thread.name!r} "
            f"ident={thread.ident} "
            f"native_id={thread.native_id} "
            f"alive={thread.is_alive()} "
            f"daemon={thread.daemon}"
        )

        if frame is not None:
            traceback.print_stack(frame)

threading.enumerate() 返回当前活动线程,ident 是 Python 线程标识,native_id 是操作系统分配的线程 ID。标识在执行期间可用于关联日志,但线程结束后可能被复用,不能把它当作永久业务身份。(docs.python.org)

3. 给锁建立可观测的等待日志

直接写:

with lock:
    update()

代码简洁,但出现长时间等待时,日志可能无法说明“谁在等、等了多久、哪一把锁”。

可以封装带日志的获取函数:

import logging
import threading
import time


logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(threadName)s %(message)s",
)


def acquire_with_log(lock: threading.Lock, name: str, timeout: float):
    start = time.monotonic()
    logging.info("waiting for lock=%s", name)

    acquired = lock.acquire(timeout=timeout)

    elapsed = time.monotonic() - start

    if acquired:
        logging.info(
            "acquired lock=%s wait_seconds=%.3f",
            name,
            elapsed,
        )
    else:
        logging.error(
            "timeout lock=%s wait_seconds=%.3f",
            name,
            elapsed,
        )

    return acquired

使用时必须保证获取成功才释放:

if not acquire_with_log(lock_a, "A", timeout=2):
    raise TimeoutError("failed to acquire A")

try:
    update()
finally:
    lock_a.release()

更稳妥的方式仍然是上下文管理器;如果需要监测持锁时间,可以封装一个上下文管理器,在进入和退出时记录时间。

4. 使用超时把永久等待变成可诊断失败

无限等待:

lock.acquire()

在未知故障下可能让整个服务永久卡住。

有界等待:

if not lock.acquire(timeout=2):
    dump_thread_stacks()
    raise TimeoutError("lock acquisition timeout")

条件等待也应设置业务合理的超时:

with cv:
    if not cv.wait_for(is_ready, timeout=5):
        raise TimeoutError("service did not become ready")

超时不是死锁修复。它做的是:

永久等待
    ↓
有限等待
    ↓
记录现场、失败、重试或降级

如果超时后直接重试而不改变状态,可能把死锁变成高频重试风暴。超时处理必须明确说明资源是否已部分获取、操作是否已经提交以及是否可以安全重放。


九、异常传播与关闭协议

直接创建线程时,线程函数中的未捕获异常不会自动在主线程中重新抛出。Python 会调用 threading.excepthook() 处理线程中的未捕获异常,默认通常输出到标准错误。(docs.python.org)

因此,不能仅靠:

t.start()
t.join()
print("all good")

判断工作是否成功。join() 只说明线程结束,不说明任务成功。

一种简单的异常传递方式:

import threading
import traceback


errors = []
errors_lock = threading.Lock()


def worker():
    try:
        do_work()
    except BaseException as exc:
        with errors_lock:
            errors.append((threading.current_thread().name, exc))
        traceback.print_exc()


def do_work():
    raise RuntimeError("failed")


t = threading.Thread(target=worker, name="worker")
t.start()
t.join()

if errors:
    name, exc = errors[0]
    raise RuntimeError(f"{name} failed") from exc

如果任务模型更适合“提交任务—获取结果—传播异常”,应考虑 concurrent.futures.ThreadPoolExecutor。Python 并发执行文档将它列为较高层的线程任务接口;线程池会管理工作线程,Future.result() 可以在调用方获取返回值或重新抛出任务异常。(docs.python.org)

线程池并不会消除竞态和死锁。尤其要警惕线程池内任务互相等待:

from concurrent.futures import ThreadPoolExecutor


def task(executor):
    future = executor.submit(lambda: 42)
    return future.result()  # 可能占满线程池后相互等待


with ThreadPoolExecutor(max_workers=1) as executor:
    executor.submit(task, executor).result()

当唯一工作线程执行 task() 时,它又提交了一个新任务并等待结果;但新任务没有空闲线程执行,形成线程池内部的等待死锁。


十、线程与 asyncio、进程的边界

Python 并发工具的选择至少要考虑两个维度:

任务类型:CPU 密集型 / I/O 密集型
编程方式:抢占式多线程 / 事件驱动协作式并发

threading 使用操作系统线程,适合已有阻塞式库、阻塞 I/O 或需要并发访问外部资源的代码。asyncio 则通过事件循环和协作式任务,在不必为每个任务创建操作系统线程的情况下处理异步 I/O;官方并发文档也将二者列为不同的并发模型。(docs.python.org)

可以粗略地这样区分:

场景 通常更自然的模型
阻塞式 HTTP、文件、数据库客户端 线程
原生支持异步的高并发网络连接 asyncio
纯 Python CPU 密集计算 进程或其他并行方案
底层扩展释放 GIL 的计算 可能使用线程
必须共享大量进程内状态 线程,但需要明确同步协议

线程和异步代码也可以组合,但边界必须明确。例如,在事件循环中直接调用阻塞函数会阻塞整个事件循环;反过来,在线程中直接操作只允许事件循环线程访问的对象,也可能破坏异步框架的约束。组合时应使用框架提供的线程桥接机制,而不是随意跨线程调用对象。


十一、Python 3.14 下需要特别注意的版本边界

1. 自由线程构建不是默认语义

Python 3.14 文档仍区分普通构建和自由线程构建。自由线程构建允许多个线程真正并行执行 Python 代码,但其可用性、扩展兼容性和性能特征不能简单等同于普通 GIL 构建。无论是否存在 GIL,业务共享状态都需要明确的锁、条件变量、队列或其他同步协议。(docs.python.org)

2. Threadcontext 参数

Python 3.14 为 Thread 增加了 context 参数,用于指定线程启动时使用的 contextvars.Context。如果不显式传递,行为受 sys.flags.thread_inherit_context 控制;文档说明该标志在自由线程构建中的默认值与普通构建不同。(docs.python.org)

因此,依赖请求上下文、租户标识、追踪 ID 的代码不应假设所有线程都自动获得相同上下文。需要传播时,应明确使用:

from contextvars import copy_context
from threading import Thread


ctx = copy_context()

t = Thread(
    target=ctx.run,
    args=(worker,),
    name="worker",
)
t.start()
t.join()

3. 线程名可以进入操作系统诊断工具

Python 3.14 中,调用 Thread.start() 时,在支持的平台上会尝试设置操作系统线程名。线程名可能受操作系统长度限制,因此应使用短而稳定的名字,例如:

Thread(
    target=worker,
    name="cache-refresh",
)

这样在线程栈、系统监控或调试器输出中,更容易将线程与业务角色关联起来。(docs.python.org)


十二、建立一套可验证的线程设计

一个线程组件至少应该能够回答以下问题:

生命周期

  • 谁创建线程?
  • 谁调用 start()
  • 谁负责停止?
  • 线程是否必须完成清理?
  • 谁调用 join()
  • 线程退出后如何报告成功或失败?

共享状态

  • 哪些变量会被多个线程访问?
  • 哪些操作必须作为一个不可分割的临界区?
  • 保护这些变量的锁是哪一把?
  • 是否存在检查—执行分离?
  • 是否可以改为不可变数据或消息队列?

等待关系

  • 线程等待的是锁、条件、事件、队列还是外部 I/O?
  • 等待是否有超时?
  • 条件变量的谓词是什么?
  • 被唤醒后是否重新检查谓词?
  • 通知发生时是否持有正确的底层锁?

死锁风险

  • 是否同时持有多把锁?
  • 全局锁顺序是什么?
  • 是否可能嵌套获取同一把普通锁?
  • 是否有线程池任务等待同一线程池中的其他任务?
  • 是否存在主线程等待工作线程、工作线程又等待主线程的环?

诊断能力

  • 每个线程是否有稳定名称?
  • 是否记录线程 ID、任务 ID 和资源名称?
  • 锁等待是否有超时和日志?
  • 进程卡住时能否导出所有线程栈?
  • 超时之后是重试、失败、降级还是退出?

线程并发正确性的基本结构可以归纳为:

明确的共享状态
    +
明确的保护边界
    +
明确的等待谓词
    +
明确的退出信号
    +
明确的超时与诊断路径

锁解决互斥,条件变量解决状态等待,事件解决信号广播,队列解决线程间消息传递;它们各自解决不同问题,不能用一个原语替代全部并发协议。真正可靠的线程代码,不是“加一把锁就不报错”,而是能够从生命周期、状态变化和等待关系上解释每一次并发行为。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。