AsyncReadWriteLock.release_write报错:NoneType无法用于await表达式
问题排查与解决
直接错误原因
你碰到的TypeError: object NoneType can't be used in 'await' expression是因为asyncio.Lock的release()是同步方法,不是协程,不能用await调用。当你执行await self._writer_lock.release()时,release()返回None,await None就触发了这个错误。
快速修复
修改release_write方法,去掉await关键字:
async def release_write(self): """Release the lock for writing.""" # Release the writer lock, only if locked if self._writer_lock.locked(): self._writer_lock.release() # 去掉await,直接调用同步方法
额外逻辑修复(实现真正的并发读)
你的AsyncReadWriteLock还有逻辑错误,会导致无法实现“多线程并发读”的预期:
- 当前
acquire_read方法里用async with self._writer_lock,会让所有读操作串行执行(因为写锁是互斥的),完全违背了读写锁的设计初衷。
修正后的acquire_read方法:
async def acquire_read(self): """Acquire the lock for reading.""" # 先确保当前没有活跃的写者 await self._writer_lock.acquire() self._writer_lock.release() # 安全增加读者计数(避免多线程竞争) async with self._readers_lock: self._readers += 1
另外,load_checkpoint函数里有个笔误:return max_completion应该改成return latest_checkpoint,否则会抛出NameError。
修复后的完整代码
修正后的AsyncReadWriteLock类
import asyncio # Read-Write asynchronous & multithreaded locking. class AsyncReadWriteLock: def __init__(self): self._readers = 0 # Number of active readers self._writer_lock = asyncio.Lock() # Lock for writers self._readers_lock = asyncio.Lock() # Lock to protect the readers counter self._readers_wait = asyncio.Condition(self._readers_lock) # Condition for readers async def acquire_read(self): """Acquire the lock for reading.""" # 确保当前无活跃写者 await self._writer_lock.acquire() self._writer_lock.release() # 安全更新读者计数 async with self._readers_lock: self._readers += 1 async def release_read(self): """Release the lock for reading.""" async with self._readers_lock: self._readers -= 1 # 无读者时通知等待的写者 if self._readers == 0: self._readers_wait.notify_all() async def acquire_write(self): """Acquire the lock for writing.""" # 先获取写锁,阻止其他写者 await self._writer_lock.acquire() # 等待所有读者完成 async with self._readers_lock: while self._readers > 0: await self._readers_wait.wait() async def release_write(self): """Release the lock for writing.""" if self._writer_lock.locked(): self._writer_lock.release() # 同步方法,无需await
修正后的load_checkpoint函数
async def load_checkpoint() -> int: """ Load the last completed batch index from checkpoint file Returns the index of the last completed batch, or -1 if no checkpoint exists """ global latest_checkpoint await checkpoint_lock.acquire_write() try: with open(CHECKPOINT_FILE, 'r') as f: line = f.readline().strip() if line and line.isdigit(): latest_checkpoint = int(line) logging.info(f"Loaded checkpoint: last completed index = {latest_checkpoint}") return latest_checkpoint # 修正笔误 else: logging.info("Checkpoint file exists but contains no valid index") return -1 except FileNotFoundError: logging.info("No checkpoint file found, starting from the beginning") return -1 except Exception as e: logging.error(f"Error loading checkpoint: {str(e)}") return -1 finally: await checkpoint_lock.release_write()
内容的提问来源于stack exchange,提问作者Jacob Malland
相关产品推荐
相关产品推荐

