CPython asyncio 同步原语全解析:Lock、Event、Condition、Semaphore 与 Barrier 的用法与源码实现
本指南以 CPython 官方文档 asyncio 同步原语 为骨架,结合 locks.py 源码与 test_locks.py 测试,系统讲解 asyncio 模块中用于协程(Task)间同步互斥的六大原语——Lock、Event、Condition、Semaphore、BoundedSemaphore 与 Barrier。读完你不仅能正确选用并编写基于 async with 的并发代码,还能理解公平调度、取消传播、状态机等底层机制,并掌握用 asyncio.wait_for 为同步操作施加超时的标准做法。
一、总体设计:与 threading 的异同
asyncio 同步原语的设计目标是尽量模仿 threading 模块中同名原语的语义,便于熟悉多线程编程的开发者平滑迁移。官方文档 asyncio-sync.rst 明确指出两个关键差异:
- 并非线程安全:asyncio 原语只应在单线程事件循环内的多个 Task 之间使用,不能用来做操作系统线程同步。若需要在真实线程间互斥,应使用 threading 模块;
- 方法不接受
timeout参数:需要限时等待时,应借助 asyncio.wait_for 函数包装相应操作。
从源码结构看,全部实现集中在 Lib/asyncio/locks.py,文件顶部定义了导出清单:
__all__ = ('Lock', 'Event', 'Condition', 'Semaphore',
'BoundedSemaphore', 'Barrier')
并在 Lib/asyncio/init.py 通过 from .locks import * 统一暴露为 asyncio.Lock 等公共 API。所有原语都继承自 Lib/asyncio/mixins.py 中的 _LoopBoundMixin,通过 _get_loop() 在首次使用时刻拿到当前运行的事件循环,这也是 3.10 版本起所有原语移除 loop 构造参数的实现基础——test_lock_doesnt_accept_loop_parameter 专门验证了传入 loop= 会抛出 TypeError。
二、Lock:互斥锁
2.1 基本语义与推荐用法
Lock() 实现一个供 asyncio 任务使用的互斥锁,保证共享资源的独占访问。它只有 locked / unlocked 两种状态,初始为 unlocked。推荐使用 async with 上下文管理器(源码见 _ContextManagerMixin.__aenter__/__aexit__,进入时 acquire()、退出时 release()):
lock = asyncio.Lock()
# ... later
async with lock:
# access shared state
它与下述显式 acquire()/release() 写法完全等价:
lock = asyncio.Lock()
# ... later
await lock.acquire()
try:
# access shared state
finally:
lock.release()
async with 写法能保证即使临界区内抛出异常,锁也一定会被释放,因此是官方推荐首选。
2.2 方法详解
| 方法 | 说明 |
|---|---|
acquire()(async) |
等待锁变为 unlocked,然后置为 locked 并返回 True。多个协程同时阻塞等待时,只有一个能继续执行 |
release() |
把锁从 locked 复位为 unlocked;若锁本就 unlocked,抛出 RuntimeError |
locked() |
锁处于 locked 时返回 True |
2.3 源码级的"公平"保证
文档强调 Lock 的获取是公平的:后来者不能插队,最终获得锁的一定是最先开始等待的那个协程。查看 locks.py 中 Lock.acquire 的实现(L90 起)可印证其机制:
async def acquire(self):
# Implement fair scheduling, where thread always waits
# its turn. Jumping the queue if all are cancelled is an optimization.
if (not self._locked and (self._waiters is None or
all(w.cancelled() for w in self._waiters))):
self._locked = True
return True
if self._waiters is None:
self._waiters = collections.deque()
fut = self._get_loop().create_future()
self._waiters.append(fut)
...
- 锁空闲且没有(或仅剩全部已取消的)等待者时,立即抢占成功;
- 否则创建一个
Future追加到deque等待队列尾部,由release()通过_wake_up_first()唤醒队首 Future——这正是 FIFO 顺序的来源; - 若等待期间任务被取消,代码会检查"锁尚未被他人获取但队列中仍有等待者"的边界情况并主动唤醒队首任务,从而保证锁不变量:不会出现无人持有锁、却没人被唤醒的死等状态。
release() 在 _locked 为假时抛出 RuntimeError('Lock is not acquired.')(对应文档"对未上锁的锁调用 release 会抛异常")。相关测试包括 test_lock、test_release_not_acquired 与验证取消竞争场景的 test_cancel_release_race 等。
2.4 加超时
由于 Lock.acquire() 不接受 timeout,需要限时获取时用 asyncio.wait_for:
try:
await asyncio.wait_for(lock.acquire(), timeout=1.0)
except asyncio.TimeoutError:
print("没能在一秒内获取锁")
else:
try:
# access shared state
finally:
lock.release()
注意:超时本质是取消内部等待的 Future,上文提到的取消处理逻辑会保证锁状态不被破坏,这是可以安全使用 wait_for 的前提。
三、Event:事件通知
3.1 基本语义
Event() 用于向多个任务广播"某个事件已经发生"。它管理一个内部标志位:初始为 False,用 set() 置为 True,用 clear() 复位为 False;wait() 会一直阻塞直到标志位为 True。
3.10 起移除构造参数
loop。
3.2 官方示例
文档 asyncio-sync.rst 给出的完整可运行示例:
import asyncio
async def waiter(event):
print('waiting for it ...')
await event.wait()
print('... got it!')
async def main():
# Create an Event object.
event = asyncio.Event()
# Spawn a Task to wait until 'event' is set.
waiter_task = asyncio.create_task(waiter(event))
# Sleep for 1 second and set the event.
await asyncio.sleep(1)
event.set()
# Wait until the waiter task is finished.
await waiter_task
asyncio.run(main())
运行输出依次为 waiting for it ...,约 1 秒后 ... got it!。Event 适合"一对多"的一次性放行语义:所有正在 wait() 的任务都会被 set() 同时唤醒。
3.3 方法详解
| 方法 | 说明 |
|---|---|
wait()(async) |
若事件已 set 立即返回 True,否则阻塞直到其他任务调用 set() |
set() |
设置事件;所有正在等待的任务被立即全部唤醒 |
clear() |
复位事件;之后调用 wait() 的任务将再次阻塞,直到下一次 set() |
is_set() |
事件已 set 时返回 True |
从 Event 的源码实现(L155 起)看,set() 会遍历内部 _waiters deque,对每个未完成的 Future 调用 set_result(True);wait() 则在标志位已为真时直接返回、否则挂起一个新的等待 Future。这就是"set 唤醒全部、已 set 后 wait 不再阻塞"语义的来源。对应测试如 test_wait_on_set 与 test_clear_with_waiters。
3.4 一个典型竞态提醒
Event 没有"一次性消费"概念:clear() 必须由业务代码显式调用。若多个消费者各自等待不同轮次的信号,推荐每次轮次使用独立的 Event,或在收到信号后主动 clear() 并接受"信号可能丢失/合并"的语义。
四、Condition:条件变量
4.1 基本语义:Event + Lock 的组合
Condition 允许一个任务等待某个事件发生,随后获得共享资源的独占访问权。本质上它结合了 Event 与 Lock 的功能。最灵活的一点是:多个 Condition 对象可以共享同一个 Lock,从而让关注同一共享资源不同状态的多个任务,能够围绕同一把锁协调互斥访问。
构造函数签名为 Condition(lock=None):可选参数 lock 必须是一个 Lock 对象,传 None 时内部自动创建一个新的 Lock。从 源码(L226 起)可见其实现——把底层锁的 locked()、acquire()、release() 三个方法直接"导出"为自身方法:
def __init__(self, lock=None):
if lock is None:
lock = Lock()
self._lock = lock
# Export the lock's locked(), acquire() and release() methods.
self.locked = lock.locked
self.acquire = lock.acquire
self.release = lock.release
self._waiters = collections.deque()
因此 Condition 自身的 acquire()/release()/locked() 语义与 Lock 完全一致(底层锁被解锁时 release() 会抛 RuntimeError)。
3.10 起移除构造参数
loop。
4.2 推荐用法与等价写法
cond = asyncio.Condition()
# ... later
async with cond:
await cond.wait()
等价于:
cond = asyncio.Condition()
# ... later
await cond.acquire()
try:
await cond.wait()
finally:
cond.release()
4.3 方法详解
| 方法 | 说明 |
|---|---|
wait()(async) |
必须先持有锁,否则抛 RuntimeError。执行时先释放底层锁并阻塞,直到被 notify()/notify_all() 唤醒;被唤醒后重新获取锁并返回 True |
wait_for(predicate)(async) |
反复调用 wait() 直到 predicate(无参可调用对象、返回值按布尔解释)为真,返回其最终值。推荐优先用它,避免处理假唤醒 |
notify(n=1) |
唤醒 n 个等待任务(默认 1);不足 n 个时全部唤醒。必须先持有锁,否则抛 RuntimeError |
notify_all() |
唤醒全部等待任务,约束同 notify() |
acquire() / release() / locked() |
委托给底层锁,见上文 |
4.4 假唤醒(spurious wakeup)与源码佐证
文档特别提示:任务可能从 wait() 中"无原因地"被唤醒返回,因此调用方应始终重新检查状态并准备再次 wait()——这也是推荐 wait_for 的原因。查看 Condition.wait 源码(L245 起)会发现,除了常规流程,源码还刻意实现了两类保护:
- 取消也必须重新持锁:即使
wait()因任务被取消而中断,外层finally仍会循环执行await self.acquire()以确保重新拿到锁后才抛出CancelledError; - 通知接力:若
wait()结束后因任何异常退出,而本任务可能已"收到过通知",会调用self._notify(1)把通知转交给下一个等待任务——这正是文档所述"允许假唤醒"这一 Condition 协议的一部分。
wait_for 的实现则是在循环中反复求值谓词:
result = predicate()
while not result:
await self.wait()
result = predicate()
return result
因此它能天然容忍假唤醒——谓词不为真就继续等。对应测试有 test_wait、test_wait_cancel_contested、test_wait_for 与 test_notify_all 等。
4.5 典型使用模式
async def consumer(cond, buffer):
async with cond:
await cond.wait_for(lambda: len(buffer) > 0)
item = buffer.pop(0)
return item
async def producer(cond, buffer, item):
async with cond:
buffer.append(item)
cond.notify() # 也可用 notify_all() 唤醒多个消费者
五、Semaphore 与 BoundedSemaphore:信号量
5.1 基本语义
Semaphore(value=1) 管理一个内部计数器:每次 acquire() 减一,每次 release() 加一。计数器永远不会小于 0——当 acquire() 发现计数器为 0 时就会阻塞,直到有任务调用 release()。可选参数 value 为计数器初始值(默认 1),传入小于 0 的值会抛出 ValueError(源码中错误消息为 "Semaphore initial value must be >= 0")。
3.10 起移除构造参数
loop。
5.2 推荐用法(限制并发数为 10)
sem = asyncio.Semaphore(10)
# ... later
async with sem:
# work with shared resource
等价于:
sem = asyncio.Semaphore(10)
# ... later
await sem.acquire()
try:
# work with shared resource
finally:
sem.release()
这种"固定并发上限"的限流模式是信号量最经典的应用,常用于限制并发下载、并发请求连接池等场景。
5.3 方法详解
| 方法 | 说明 |
|---|---|
acquire()(async) |
计数器大于 0 时立即减一并返回 True;为 0 时阻塞到 release() 后被唤醒再返回 True |
locked() |
当信号量无法被立即获取时返回 True |
release() |
计数器加一;若有任务正等待获取,唤醒之 |
与 BoundedSemaphore 不同,普通 Semaphore 允许 release() 的次数多于 acquire(),计数器可以无限增长,因此多调 release() 并不会报错。
5.4 源码细节:FIFO 与取消回补
从 Semaphore.acquire 源码(L383 起)可见,其等待队列同样使用 deque 维护,注释明确写着 "Maintain FIFO, wait for others to start even if _value > 0"——即使计数器为正,新到者也会排在老等待者之后,与 Lock 一致地保持公平。若一个已收到唤醒信号的等待任务在真正获得许可前被取消,取消处理分支会执行 self._value += 1 回补计数并把机会让给下一位,不会丢失许可额度。test_acquire_fifo_order 用三个协程各抢两次、验证结果严格按 c1/c2/c3 顺序交错输出,正是这一行为的测试证据。
5.5 BoundedSemaphore:防止过度释放
BoundedSemaphore(value=1) 是 Semaphore 的受限版本:若某次 release() 会让内部计数器超过初始 value,则抛出 ValueError。其实现仅覆盖了 release():
class BoundedSemaphore(Semaphore):
def __init__(self, value=1):
self._bound_value = value
super().__init__(value)
def release(self):
if self._value >= self._bound_value:
raise ValueError('BoundedSemaphore released too many times')
super().release()
3.10 起移除构造参数
loop。
这在工程上是更安全的默认选择:能在开发期尽早暴露"释放次数超过获取次数"这类资源泄漏式 bug。
六、Barrier:屏障
6.1 基本语义
Barrier(parties) 让调用方阻塞,直到有 parties 个任务在屏障上等待,届时所有等待任务会同时被解除阻塞。屏障可以任意次数重复使用,非常适合"多任务在已知同步点会合后一起继续"的并行模式(如分阶段计算的阶段栅栏)。
- 新增于 3.11 版本;
- 构造时若
parties < 1会抛ValueError; async with barrier可作为await barrier.wait()的替代写法(其上下文管理器会把wait()的返回值赋给as变量,见下文 6.3)。
6.2 官方示例与预期输出
文档 asyncio-sync.rst 给出的示例:
import asyncio
async def example_barrier():
# barrier with 3 parties
b = asyncio.Barrier(3)
# create 2 new waiting tasks
asyncio.create_task(b.wait())
asyncio.create_task(b.wait())
await asyncio.sleep(0)
print(b)
# The third .wait() call passes the barrier
await b.wait()
print(b)
print("barrier passed")
await asyncio.sleep(0)
print(b)
asyncio.run(example_barrier())
运行结果为:
<asyncio.locks.Barrier object at 0x... [filling, waiters:2/3]>
<asyncio.locks.Barrier object at 0x... [draining, waiters:0/3]>
barrier passed
<asyncio.locks.Barrier object at 0x... [filling, waiters:0/3]>
从输出可见 Barrier 的完整生命周期:前两个任务进入后处于 filling(填充) 态、waiters:2/3;第三个任务补齐后放行,屏障进入 draining(排空) 态;参与者陆续离开后屏障自动回到 [filling, waiters:0/3],等待下一轮复用。repr 中的状态名与源码中定义的 _BarrierState 枚举一一对应(见 locks.py 的枚举定义 L466 起:FILLING / DRAINING / RESETTING / BROKEN)。
6.3 方法、属性与异常
| 成员 | 说明 |
|---|---|
wait()(async) |
在屏障上等待;当 parties 个任务都调用后,大家同时被放行。返回值是 0 到 parties-1 区间内的整数,且对每个任务互不相同,可用于指派某个任务做特殊收尾工作。若等待期间屏障被 reset/broken 会抛 BrokenBarrierError,任务被取消则抛 CancelledError。被取消的任务会离开屏障(filling 态下等待计数减 1),屏障状态保持不变 |
reset()(async) |
把屏障复位为默认空状态;所有正在等待的任务会收到 BrokenBarrierError。屏障已损坏时,与其 reset 不如直接弃用并新建一个 |
abort()(async) |
将屏障置为 broken 状态;当前及未来的所有 wait() 调用都会以 BrokenBarrierError 失败。当某个参与方需要中止、以避免其余任务无限等待时使用 |
parties(属性) |
通过屏障所需的任务数 |
n_waiting(属性) |
当前处于 filling 态下正在等待的任务数 |
broken(属性) |
屏障处于 broken 状态时为 True |
BrokenBarrierError(异常) |
RuntimeError 的子类,在 Barrier 被 reset 或 broken 时抛出。其定义位于 Lib/asyncio/exceptions.py 的 BrokenBarrierError(RuntimeError) |
6.4 利用返回值做"队首任务"收尾
wait() 的返回值可用来挑出唯一任务执行特殊处理,文档示例:
...
async with barrier as position:
if position == 0:
# Only one task prints this
print('End of *draining phase*')
内部实现上,Barrier.wait(L508 起)维护 _count 作为当前到达序号,把该序号作为返回值,并在 index + 1 == self._parties 时由最后到达者调用 _release() 唤醒全体;整个流程包在一把内部 Condition(即 Barrier 的 _cond)里,状态迁移由 _BarrierState 状态机驱动。测试覆盖非常细,包括正常放行、filling/draining 各阶段与 reset/abort 的组合场景,如 test_barrier、test_reset_barrier_while_tasks_waiting、test_abort_barrier 与 test_blocking_tasks_while_draining 等。
七、版本演进与兼容性注意事项
- 3.9:
await lock、yield from lock以及with await lock、with (yield from lock)等旧式用法被彻底移除,必须改用async with lock; - 3.10:
Lock、Event、Condition、Semaphore、BoundedSemaphore全部移除构造参数loop(原语不再绑定外部循环,改用运行中循环); - 3.11:新增
Barrier与BrokenBarrierError。
需要特别留意:所有原语对象本身不可被直接 await(Lock 已不实现 __await__)。测试 test_lock_by_with_statement 明确断言 await lock 或 with await lock: 会抛出 TypeError: 'Lock' object can't be awaited,把这种误用挡在编译/运行早期。
八、选型速查与实践建议
| 场景 | 首选原语 |
|---|---|
| 保护一段临界区,只允许一个任务进入 | Lock |
| 广播"某事件已发生",多个等待者全部放行 | Event |
| 需要"等条件成立 + 独占访问共享资源" | Condition(可用 wait_for 简化) |
| 限制并发数为 N(N ≥ 1) | Semaphore(N) |
| 同前,且想尽早暴露过度释放 bug | BoundedSemaphore(N) |
| 等待固定数量任务在同步点会合、可反复使用 | Barrier(parties) |
工程实践上的三条建议:
- 优先使用
async with,让锁/条件/信号量的获取与释放自动配对,杜绝finally遗漏导致的死锁; - 所有需要超时的等待一律包
asyncio.wait_for,不要尝试自行拼接"轮询 + sleep"; - 牢记非线程安全约束:涉及
loop.run_in_executor或真实多线程共享时,应使用 threading 体系的Lock、Event、Semaphore,而不是本模块对象。
若想更深入验证或扩展理解,可在当前仓库中直接阅读三类一手资料:实现全文 Lib/asyncio/locks.py、API 定义 Doc/library/asyncio-sync.rst、以及覆盖取消竞争、FIFO 公平性、Barrier 状态机等边界场景的 Lib/test/test_asyncio/test_locks.py。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0626
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00