
CPython asyncio 同步原语全解析Lock、Event、Condition、Semaphore 与 Barrier 的用法与源码实现【免费下载链接】cpythonThe Python programming language项目地址: https://gitcode.com/GitHub_Trending/cp/cpython本指南以 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抛出RuntimeErrorlocked()锁处于locked时返回True2.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_fortry: await asyncio.wait_for(lock.acquire(), timeout1.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()会遍历内部_waitersdeque对每个未完成的 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(lockNone)可选参数lock必须是一个Lock对象传None时内部自动创建一个新的Lock。从 源码L226 起可见其实现——把底层锁的locked()、acquire()、release()三个方法直接导出为自身方法def __init__(self, lockNone): if lock is None: lock Lock() self._lock lock # Export the locks 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()唤醒被唤醒后重新获取锁并返回Truewait_for(predicate)async反复调用wait()直到predicate无参可调用对象、返回值按布尔解释为真返回其最终值。推荐优先用它避免处理假唤醒notify(n1)唤醒n个等待任务默认 1不足n个时全部唤醒。必须先持有锁否则抛RuntimeErrornotify_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(value1)管理一个内部计数器每次acquire()减一每次release()加一。计数器永远不会小于 0——当acquire()发现计数器为 0 时就会阻塞直到有任务调用release()。可选参数value为计数器初始值默认 1传入小于 0 的值会抛出ValueError源码中错误消息为Semaphore initial value must be 0。3.10 起移除构造参数loop。5.2 推荐用法限制并发数为 10sem 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()后被唤醒再返回Truelocked()当信号量无法被立即获取时返回Truerelease()计数器加一若有任务正等待获取唤醒之与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(value1)是Semaphore的受限版本若某次release()会让内部计数器超过初始 value则抛出ValueError。其实现仅覆盖了release()class BoundedSemaphore(Semaphore): def __init__(self, value1): 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会抛ValueErrorasync 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 状态时为TrueBrokenBarrierError异常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.waitL508 起维护_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.9await lock、yield from lock以及with await lock、with (yield from lock)等旧式用法被彻底移除必须改用async with lock3.10Lock、Event、Condition、Semaphore、BoundedSemaphore全部移除构造参数loop原语不再绑定外部循环改用运行中循环3.11新增Barrier与BrokenBarrierError。需要特别留意所有原语对象本身不可被直接awaitLock已不实现__await__。测试 test_lock_by_with_statement 明确断言await lock或with await lock:会抛出TypeError: Lock object cant be awaited把这种误用挡在编译/运行早期。八、选型速查与实践建议场景首选原语保护一段临界区只允许一个任务进入Lock广播某事件已发生多个等待者全部放行Event需要等条件成立 独占访问共享资源Condition可用wait_for简化限制并发数为 NN ≥ 1Semaphore(N)同前且想尽早暴露过度释放 bugBoundedSemaphore(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。【免费下载链接】cpythonThe Python programming language项目地址: https://gitcode.com/GitHub_Trending/cp/cpython创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考