首页
/ CPython asyncio 同步原语全解析:Lock、Event、Condition、Semaphore 与 Barrier 的用法与源码实现

CPython asyncio 同步原语全解析:Lock、Event、Condition、Semaphore 与 Barrier 的用法与源码实现

2026-09-06 18:53:45作者:魏献源Searcher

本指南以 CPython 官方文档 asyncio 同步原语 为骨架,结合 locks.py 源码与 test_locks.py 测试,系统讲解 asyncio 模块中用于协程(Task)间同步互斥的六大原语——LockEventConditionSemaphoreBoundedSemaphoreBarrier。读完你不仅能正确选用并编写基于 async with 的并发代码,还能理解公平调度、取消传播、状态机等底层机制,并掌握用 asyncio.wait_for 为同步操作施加超时的标准做法。

一、总体设计:与 threading 的异同

asyncio 同步原语的设计目标是尽量模仿 threading 模块中同名原语的语义,便于熟悉多线程编程的开发者平滑迁移。官方文档 asyncio-sync.rst 明确指出两个关键差异:

  1. 并非线程安全:asyncio 原语只应在单线程事件循环内的多个 Task 之间使用,不能用来做操作系统线程同步。若需要在真实线程间互斥,应使用 threading 模块;
  2. 方法不接受 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_locktest_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() 复位为 Falsewait() 会一直阻塞直到标志位为 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_settest_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 起)会发现,除了常规流程,源码还刻意实现了两类保护:

  1. 取消也必须重新持锁:即使 wait() 因任务被取消而中断,外层 finally 仍会循环执行 await self.acquire() 以确保重新拿到锁后才抛出 CancelledError
  2. 通知接力:若 wait() 结束后因任何异常退出,而本任务可能已"收到过通知",会调用 self._notify(1) 把通知转交给下一个等待任务——这正是文档所述"允许假唤醒"这一 Condition 协议的一部分。

wait_for 的实现则是在循环中反复求值谓词:

result = predicate()
while not result:
    await self.wait()
    result = predicate()
return result

因此它能天然容忍假唤醒——谓词不为真就继续等。对应测试有 test_waittest_wait_cancel_contestedtest_wait_fortest_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 个任务都调用后,大家同时被放行。返回值是 0parties-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.pyBrokenBarrierError(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_barriertest_reset_barrier_while_tasks_waitingtest_abort_barriertest_blocking_tasks_while_draining 等。

七、版本演进与兼容性注意事项

  • 3.9await lockyield from lock 以及 with await lockwith (yield from lock) 等旧式用法被彻底移除,必须改用 async with lock
  • 3.10LockEventConditionSemaphoreBoundedSemaphore 全部移除构造参数 loop(原语不再绑定外部循环,改用运行中循环);
  • 3.11:新增 BarrierBrokenBarrierError

需要特别留意:所有原语对象本身不可被直接 awaitLock 已不实现 __await__)。测试 test_lock_by_with_statement 明确断言 await lockwith await lock: 会抛出 TypeError: 'Lock' object can't be awaited,把这种误用挡在编译/运行早期。

八、选型速查与实践建议

场景 首选原语
保护一段临界区,只允许一个任务进入 Lock
广播"某事件已发生",多个等待者全部放行 Event
需要"等条件成立 + 独占访问共享资源" Condition(可用 wait_for 简化)
限制并发数为 N(N ≥ 1) Semaphore(N)
同前,且想尽早暴露过度释放 bug BoundedSemaphore(N)
等待固定数量任务在同步点会合、可反复使用 Barrier(parties)

工程实践上的三条建议:

  1. 优先使用 async with,让锁/条件/信号量的获取与释放自动配对,杜绝 finally 遗漏导致的死锁;
  2. 所有需要超时的等待一律包 asyncio.wait_for,不要尝试自行拼接"轮询 + sleep";
  3. 牢记非线程安全约束:涉及 loop.run_in_executor 或真实多线程共享时,应使用 threading 体系的 LockEventSemaphore,而不是本模块对象。

若想更深入验证或扩展理解,可在当前仓库中直接阅读三类一手资料:实现全文 Lib/asyncio/locks.py、API 定义 Doc/library/asyncio-sync.rst、以及覆盖取消竞争、FIFO 公平性、Barrier 状态机等边界场景的 Lib/test/test_asyncio/test_locks.py

登录后查看全文
热门项目推荐
相关项目推荐