前言
线程共享内存,所以只要两个线程碰同一个可变对象,就必须想清楚「怎么同步」。但同步不等于「到处加锁」——threading模块提供了五六种原语,每种解决的是不同形状的协调问题:有的保护一段临界区,有的负责线程之间传消息,有的只是发个信号,还有的是限制同时干活的人数。
新手最常见的误区是「一把大锁走天下」:所有共享访问都套同一个Lock,结果是锁范围过大、线程全都排队,多线程退化成单线程,还容易死锁。另一类误区是反过来,用普通列表当队列、两边一读一写,偶发RuntimeError或数据丢失。
本文不讲语法表,而是给五个真实形状的场景:计数器、生产者-消费者、优雅停止、限流、分阶段汇合,每个场景配一个可运行的例子,并说明为什么选这种原语而不是别的。示例以 CPython 3.8 及以上为基准。
一、场景:共享计数器 —— 用Lock
问题形状:多个线程要给同一个计数累加,读-改-写不能被打断。
# 适用于 Python 3.8+
import threading
counter = 0
counter_lock = threading.Lock()
def increase(times):
global counter
for _ in range(times):
with counter_lock: # 临界区只包住真正需要保护的三步
counter += 1
threads = [threading.Thread(target=increase, args=(50_000,)) for _ in range(4)]
for t in threads:
t.start()
for t in threads:
t.join()
print(counter) # 400000为什么用Lock:这里需要的是互斥——同一时刻只允许一个线程进入临界区。Lock是最轻的选择。要点是把临界区缩到最小:把print、文件读写这些不需要保护的操作挪到锁外面。
什么时候改RLock:如果同一个线程可能嵌套获取同一把锁(例如一个被锁保护的函数内部又调用了另一个同样加锁的函数),Lock会让它自己把自己锁死,这时必须用RLock。注意RLock只能由加锁的那个线程释放。
二、场景:生产者-消费者 —— 用queue.Queue
问题形状:一批线程产出数据,另一批线程消费数据,两边速度不匹配。
这是最经典也最容易写错的场景。手写list+Condition虽然可以,但queue.Queue已经把锁、阻塞、唤醒、计数全封装好了。
# 适用于 Python 3.8+
import queue
import threading
import time
q = queue.Queue(maxsize=5) # 队列上限 5,满了生产者会阻塞
def producer(n):
for i in range(n):
item = f"数据-{i}"
q.put(item) # 队列满时自动阻塞
print(f"生产 {item}")
time.sleep(0.02)
def consumer(name):
while True:
item = q.get() # 队列空时自动阻塞
if item is None: # 约定:None 表示收工
q.task_done()
break
print(f"[{name}] 消费 {item}")
time.sleep(0.05)
q.task_done() # 告诉队列这一条处理完了
p = threading.Thread(target=producer, args=(10,))
c = threading.Thread(target=consumer, args=("C1",), name="C1")
p.start()
c.start()
p.join()
q.put(None) # 发送结束信号
c.join()
q.join() # 等所有 task_done
print("全部完成")为什么用Queue:queue模块文档明确写着,Queue已经「实现了所有必需的加锁语义」,并且可以安全地在多个生产者和多个消费者之间传递。task_done()与join()配对,用于等「所有取出的任务都被处理完」。
两个关键细节:一是用maxsize做背压,防止生产太快把内存吃光;二是结束信号要选一个业务数据中不可能出现的值,示例里用的None就是常见约定。
三、场景:优雅停止与启动栅栏 —— 用Event
问题形状:一个线程干活,另一个线程要能让它停下来;或者多个线程要一起「等口令」再开始。
# 适用于 Python 3.8+
import threading
import time
stop_event = threading.Event()
start_event = threading.Event()
def worker(name):
print(f"{name} 等待启动口令")
start_event.wait() # 阻塞直到被别人 set()
while not stop_event.is_set():
time.sleep(0.05)
print(f"{name} 工作中……")
print(f"{name} 已停止")
ts = [threading.Thread(target=worker, args=(f"W{i}",)) for i in range(3)]
for t in ts:
t.start()
time.sleep(0.1)
start_event.set() # 一声令下,三个线程同时开始
time.sleep(0.2)
stop_event.set() # 通知大家收工
for t in ts:
t.join()为什么用Event:Event管的是一个布尔标志,set()会唤醒所有等待的线程,wait(timeout=None)返回True(被置位)或False(超时)。它比Condition更简单,因为它不需要先持有锁。这类「广播一个状态」的场景,用Event最贴切。
常用技巧:把while not stop.wait(timeout=1)写进循环,就同时获得了「定时轮询」和「立即响应停止信号」两个效果,比while not stop: sleep(1)强得多。
四、场景:限制并发数 —— 用Semaphore
问题形状:有一批任务,但同一时刻最多只允许若干个同时访问某个有限资源(比如数据库连接数、目标站点的并发上限)。
Lock的计数只有 0 和 1,Semaphore把计数扩展到 N。
# 适用于 Python 3.8+
import threading
import time
pool_limit = threading.Semaphore(3) # 同时最多 3 个线程进入
def access_db(i):
with pool_limit: # 计数为 0 时在这里阻塞
print(f"任务 {i} 拿到名额")
time.sleep(0.2) # 模拟一次数据库访问
threads = [threading.Thread(target=access_db, args=(i,)) for i in range(8)]
for t in threads:
t.start()
for t in threads:
t.join()
print("全部结束")为什么用Semaphore:它管理的计数器表示「release 次数减 acquire 次数再加上初值」,acquire()在必要时阻塞,保证计数不会被减成负数。8 个任务会分批放行,每批最多 3 个。
一个更严格的选择是BoundedSemaphore:它会在release()次数超过初值时抛ValueError。因为「释放多了」几乎一定是代码 bug,用有界信号量能把这种 bug 尽早暴露出来——官方文档也是这么建议的。
五、场景:分阶段汇合 —— 用Barrier
问题形状:多个线程分成若干阶段执行,必须所有线程都完成上一阶段,才能一起进入下一阶段。
# 适用于 Python 3.8+
import threading
import time
barrier = threading.Barrier(3)
def stage_worker(name):
print(f"{name} 完成第一阶段数据准备")
time.sleep(0.05)
index = barrier.wait() # 等齐 3 个线程才一起通过
print(f"{name} 通过栅栏(序号 {index}),开始第二阶段")
ts = [threading.Thread(target=stage_worker, args=(f"W{i}",)) for i in range(3)]
for t in ts:
t.start()
for t in ts:
t.join()为什么用Barrier:它专为「N 个线程必须全部到齐才放行」设计。wait()返回 0 到parties-1之间的整数,每个线程不同,可以用来指定一个线程做收尾工作(比如if index == 0:打印汇总)。如果有一个线程超时或abort(),栅栏进入 broken 状态,其他等待的线程会收到BrokenBarrierError。
六、怎么选:一张对照表
| 场景形状 | 首选原语 | 为什么 |
|---|
| 保护一段共享读改 | Lock | 最简单、开销最小 |
| 同一线程嵌套加锁 | RLock | 可重入,避免自锁死 |
| 多生产多消费传数据 | queue.Queue | 自带锁与阻塞语义 |
| 等待某个条件成立 | Condition | 可等待复杂谓词(wait_for) |
| 广播一个状态/优雅停止 | Event | 简单布尔标志,可唤醒全部 |
| 限制同时访问的线程数 | Semaphore/BoundedSemaphore | 计数可大于 1 |
| 全部到齐才继续 | Barrier | 分阶段同步 |
一个通用原则:能用Queue传数据,就不要用共享变量加锁。队列把「同步」这件事收敛到一个被反复测试过的实现里,比手写锁安全得多。
常见坑点
1. 用普通列表当队列,两头并发读写
❌ 一个线程lst.append(x),另一个线程while lst: lst.pop(0),偶发异常或数据错乱。 ✅ 改用queue.Queue,它内部自带锁。
2. 锁的范围包太大
❌ 把print、网络请求、文件读写全都塞进with lock:。 ✅ 只把真正的读-改-写放进临界区,其余挪到锁外。
3. 用Lock却写了嵌套加锁
❌ 持锁期间又去acquire()同一把Lock,线程永久卡住。 ✅ 需要重入就用RLock;更好的做法是理顺调用层次,避免嵌套加锁。
4. 忘记task_done(),q.join()永远不返回
❌ 消费者取走数据却没调用task_done(),主线程q.join()挂死。 ✅ 每处理完一条就task_done();包括收到结束信号那一次。
5. 用time.sleep轮询停止信号
❌while not stop.is_set(): time.sleep(1),最多延迟 1 秒才停。 ✅while not stop.wait(timeout=1):,set()时立刻返回。
6.Semaphore释放次数多于获取次数
❌ 用Semaphore却多调了release(),计数虚高,限流形同虚设。 ✅ 需要严格保护资源上限时改用BoundedSemaphore,让超额释放直接报错。
7. 用Event代替Lock做互斥
❌ 两个线程用Event互相等待对方set(),写成复杂易错的交替逻辑。 ✅Event是状态广播,不是互斥量;互斥请用Lock。
8. 多把锁以不同顺序获取,造成死锁
❌ 线程 A 先拿锁 1 再拿锁 2,线程 B 反过来,互相等对方释放。 ✅ 全局约定统一的加锁顺序;能不嵌套就不嵌套。
总结
| 原语 | 解决的形状 | 一句话记忆 |
|---|
Lock | 互斥 | 同一时刻只许一个进 |
RLock | 可重入互斥 | 同一线程可以反复进 |
Condition | 等条件 | 等谓词成立再走 |
Event | 广播状态 | 一声令下全体响应 |
Semaphore | 限流 | 同时最多 N 个 |
Barrier | 汇合 | 到齐了才放行 |
queue.Queue | 传递数据 | 自带锁的安全管道 |
同步机制的选型,本质是先看清「问题形状」,再挑对应的原语。互斥、传值、广播、限流、汇合,是五种完全不同的形状;用错了形状,即使语法正确,也会写出又慢又容易死锁的代码。实在拿不准时,优先退回到queue.Queue加线程池的组合——它覆盖了大多数真实需求,且比手写锁更难出错。