从线程代码转asyncio:如何无等待尝试获取asyncio.Lock?
异步非阻塞获取asyncio.Lock的实现方案
核心思路
asyncio.Lock没有提供直接的blocking=False参数,但可以通过asyncio.wait_for配合超时时间0来模拟非阻塞获取锁的行为——当锁被占用时,wait_for会立即抛出TimeoutError,我们通过捕获这个错误来判断是否成功获取锁。
异步版try_acquire_lock辅助函数
替换原来的同步上下文管理器,使用@asynccontextmanager实现异步版本:
from contextlib import asynccontextmanager from typing import AsyncIterator import asyncio @asynccontextmanager async def try_acquire_async_lock(lock: asyncio.Lock) -> AsyncIterator[bool]: locked = False try: # 用timeout=0模拟非阻塞获取 await asyncio.wait_for(lock.acquire(), timeout=0) locked = True yield locked except asyncio.TimeoutError: # 锁被占用,返回False yield locked finally: if locked: lock.release()
适配后的异步generate方法
把原来的同步方法改成异步函数,配合上面的上下文管理器使用:
async def generate(self, cti: Freeswitch_acd_api) -> bool: log = logger.getChild('Cached_data._generate') if self.data_expiration and tz.utcnow() < self.data_expiration: log.debug(f'{type(self).__name__} data not expired yet') else: async with try_acquire_async_lock(self.lock) as locked: if locked: log.debug(f'{type(self).__name__} regenerating') try: # 若self._generate是同步阻塞方法,需用asyncio.to_thread包装 # new_data = await asyncio.to_thread(self._generate, cti) # 若是异步方法,直接await即可 new_data = await self._generate(cti) except Freeswitch_error as e: log.exception('FS error trying to generate data: %r', e) return False else: self.data = new_data self.data_expiration = tz.utcnow() + tz.timedelta(seconds=self.max_cache_seconds) return True
关键注意事项
- 确保锁实例是
asyncio.Lock(或asyncio.RLock,如需可重入),而非线程锁。 - 若原
self._generate是同步阻塞方法,必须用asyncio.to_thread包装为异步任务,避免阻塞事件循环。 - 异步上下文管理器必须用
async with调用,不能使用普通with。
场景适配说明
针对你提到的三个异步任务连接不同服务器的场景:
- 当其中一个任务成功获取锁并更新缓存时,其他任务会立即得到
locked=False的结果,直接跳过更新逻辑,避免重复操作。 - 若持有锁的任务因服务器维护失败,后续任务检查缓存过期时,会重新尝试获取锁并从其他服务器拉取数据,保留了多服务器的容错能力。
内容的提问来源于stack exchange,提问作者royce3
相关产品推荐
相关产品推荐

