一、为什么需要RLock

RLock(可重入锁)提供了以下好处:

1. 允许可重入性(Reentrancy)

在异步编程(如使用 asyncio)中,如果一个协程需要多次获取同一个锁(例如在递归调用或嵌套调用中),普通的 Lock 会阻塞自身,导致死锁。而 RLock 允许同一个协程多次获取锁,只需确保最终释放的次数与获取次数一致。

示例场景
如果某个函数内部多次调用自身(递归),或调用其他需要同一锁的函数,使用 RLock 可以避免阻塞。

async def nested_function(lock):
    await lock.acquire()
    try:
        # 递归调用
        await nested_function(lock)
    finally:
        lock.release()

如果使用普通 Lock,上述代码会导致死锁,因为同一个协程再次调用 acquire 时会等待自己释放锁。而 RLock 会允许同一协程多次获取锁。

2. 避免死锁风险

在复杂的异步逻辑中,若多个锁的获取顺序不当,可能导致死锁。RLock 通过允许可重入性,减少了因同一协程多次请求锁而引发的死锁可能性。

3. 灵活控制资源访问

在需要分层次或分阶段访问共享资源的场景中,RLock 可以更灵活地管理锁的层级结构。例如,在异步数据库操作或文件读写中,确保同一协程的不同步骤安全地操作资源。

4. 兼容同步代码的习惯

熟悉多线程编程的开发者在迁移到异步模式时,RLock 提供了与同步代码(如 threading.RLock)类似的语义,降低了学习成本。


二、手动实现一个 async.RLock (AsyncRLock)


class AsyncRLock(object):
    def __init__(self):
        self._lock = asyncio.Lock()  # 底层互斥锁: 非可重入
        self._owner: Optional[asyncio.Task] = None  # 当前持有锁的协程(Task)
        self._count = 0  # 可重入计数

    async def acquire(self) -> None:
        curr_task = asyncio.current_task()

        if self._owner is curr_task:
            # 递归调用,计数+1
            self._count += 1
            return

        # 有2种情况:
        # - 没有人占用锁,直接获取锁
        # - 被锁占用, 等待被释放后重新获取锁
        await self._lock.acquire()
        self._owner = curr_task
        self._count = 1

    def release(self) -> None:
        '''
        注意这里没有定义为 async 方法 !!!
        因为: 释放锁只是修改内部状态, 不涉及等待、阻塞或 I/O 没必要用 async def
        '''
        curr_task = asyncio.current_task()

        # 非锁持有者, 不允许调用 release
        if self._owner is not curr_task:
            raise RuntimeError("Cannot release un-acquired lock")

        # 引用次数减一
        self._count -= 1

        # 如果已经没有持有数量, 强制释放锁
        if self._count <= 0:
            self._count = 0
            self._owner = None
            self._lock.release()   # 必须最后释放锁, 否则会有 meta 混乱的问题

    def locked(self) -> bool:
        '''
        同样, 不涉及 等待、阻塞或 I/O, 没有设置为 async def
        '''

        return self._owner is not None

    ##########################################################################
    # 支持被 async with 调用
    async def __aenter__(self):
        await self.acquire()
        return self

    async def __aexit__(self, exc_type, exc, tb):
        self.release()

    ##########################################################################

    ##########################################################################
    # 防止被 with AsyncRLock 同步误调用
    def __enter__(self):
        raise RuntimeError(
            "Use 'async with' instead of 'with' for AsyncRLock"
        )

    def __exit__(self, exc_type, exc, tb):
        pass

    ##########################################################################


三、UT设计:如何验证 AsyncRLock的正确性(完整的代码,参考第四部分)

RLock 的三个关键特性, 我们就按照3个case 分别验证:

  1. 可重入性:允许同一协程多次获取锁。
  2. 释放后可获取性:锁释放后,其他协程可以立即获取。
  3. 互斥性:同一时间只有一个协程能持有锁。

Case 1:测试可重入性(Reentrancy)

测试函数
async def test_for_rlock():
    # 测试互斥:另一个协程必须等待
    task_other = asyncio.create_task(test_for_get_rock_other(rlock))
    task_outer = asyncio.create_task(test_for_get_rock_outer(rlock))
    result1, result2 = await asyncio.gather(task_other, task_outer)
测试逻辑
  1. 流程

    • test_for_get_rock_outer 先获取 RLock
    • 它调用 test_for_get_rock_inner,后者再次尝试获取同一 RLock
    • 如果 RLock 是可重入的,内层函数将成功获取锁(Inner acquired)。
  2. 验证的特性

    • 可重入性:同一协程(outer)可以多次嵌套获取锁,而不会阻塞自己。
  3. 输出验证

    --  测试互斥:另一个协程必须等待
    $$ [other task] sleep 0.1s
    ** Outer acquired
        ** Inner acquired
    ** Back in outer
    $$ Other task acquired, cost time: ~0.200s (外层任务执行时间)
    result1: other
    result2: inner
    

    ✅ 输出显示 Inner acquired 成功,且未出现死锁,证明 RLock 允许可重入。


Case 2:测试锁释放后的可获取性

测试函数
async def test_for_rlock():
    ....
    # 之前已经获取 & 释放了

    # 目前锁已经被释放了,可以再次获取
    print('--  已经被释放了, 测试可以再次获取')
    await test_for_get_rock_outer(rlock)
    await test_for_get_rock_other(rlock)
    print('\n\n')
测试逻辑
  1. 流程

    • 先运行 test_for_get_rock_outer,释放锁后。
    • 启动 test_for_get_rock_other,它等待 0.1s 后尝试获取锁。
    • 如果锁已释放,Other task 应立即获取锁(耗时极短)。
  2. 验证的特性

    • 锁的正确释放与可用性:锁释放后,其他协程可以立即获取。
  3. 输出验证

    --  已经被释放了, 测试可以再次获取
    ** Outer acquired
        ** Inner acquired
    ** Back in outer
    $$ [other task] sleep 0.1s
    
    $$ Other task acquired, cost time: 4.0531158447265625e-06
    
    ✅ 耗时接近 0,说明锁被释放后可立即获取,验证了 RLock 的状态管理正确。

Case 3:测试互斥性(Mutual Exclusion)

测试函数
async def test_for_rlock():
    rlock = AsyncRLock()
    .......

    # 锁已经被再次释放!

    # 测试互斥:另一个协程必须等待
    print('--  测试互斥:另一个协程必须等待')
    task_other = asyncio.create_task(test_for_get_rock_other(rlock))
    task_outer = asyncio.create_task(test_for_get_rock_outer(rlock))

    result1, result2 = await asyncio.gather(task_other, task_outer)
    print(f'result1: {result1}')
    print(f'result2: {result2}')
测试逻辑
  1. 流程

    • 同时启动 task_outer(持有锁)和 task_other(尝试获取锁)。
    • task_outer 先获取锁,task_other 必须等待直到锁被释放。
    • task_outer 执行完成后释放锁,task_other 才能获取锁。
  2. 验证的特性

    • 互斥性:同一时间只有一个协程能持有锁,其他协程必须等待。
  3. 输出验证

    --  测试互斥:另一个协程必须等待
    $$ [other task] sleep 0.1s
    ** Outer acquired
        ** Inner acquired
    ** Back in outer
    
    $$ Other task acquired, cost time: 0.1001734733581543
    result1: other
    result2: inner
    
    ✅ Other task 的耗时接近 0.2s(外层任务的执行时间),说明它被阻塞直到锁释放,证明互斥性成立。

关键验证点总结

Case验证的特性如何验证
可重入性允许同一协程多次加锁outer 调用 inner 时,两者均能成功获取锁,且无死锁。
释放后可获取锁释放后可被其他协程立即获取other task 在锁释放后几乎立即获取,耗时接近 0
互斥性同一时间仅一个协程持有锁other task 必须等待 outer 释放锁后才能获取,耗时等于 outer 的执行时间。

代码可靠性验证的关键细节

  1. 可重入性

    • test_for_get_rock_inner 在 test_for_get_rock_outer 内部调用,且两者共享同一 RLock
    • async with lock 的嵌套使用不会阻塞,证明 RLock 允许可重入。
  2. 释放后可获取

    • test_for_get_rock_other 在锁释放后立即获取锁,耗时极短,说明锁状态正确。
  3. 互斥性

    • 两个协程同时运行时,task_other 被阻塞直到 task_outer 释放锁,耗时等于 task_outer 的执行时间(0.2s),证明互斥性。

四、完整的测试Case

async def test_for_get_rock_inner(lock):
    '''
    测试被递归调用获取lock
    '''

    async with lock:
        print("    ** Inner acquired")
        await asyncio.sleep(0.2)
        return "inner"


async def test_for_get_rock_outer(lock):
    '''
    测试被 outer 调用获取lock
    '''
    async with lock:
        print("** Outer acquired")
        result = await test_for_get_rock_inner(lock)
        print("** Back in outer")
        return result

async def test_for_get_rock_other(lock):
    '''
    测试被其他协程调用获取lock, 被动等待结束后获得锁
    '''

    # 确保不先获得锁
    print('$$ [other task] sleep 0.1s')
    await asyncio.sleep(0.1)

    # 预期被其他 task 获得锁, 等待到其他 task 释放后获得
    btime = time.time()
    async with lock:
        print(f"\n$$ Other task acquired, cost time: {time.time() - btime}")

    return 'other'

async def test_for_rlock():
    rlock = AsyncRLock()

    # 测试可重入: outer → inner
    print('--  测试可重入: outer → inner')
    result = await test_for_get_rock_outer(rlock)
    print("Result:", result)
    print('\n\n')

    # 目前锁已经被释放了,可以再次获取
    print('--  已经被释放了, 测试可以再次获取')
    await test_for_get_rock_outer(rlock)
    await test_for_get_rock_other(rlock)
    print('\n\n')

    # 锁已经被再次释放!

    # 测试互斥:另一个协程必须等待
    print('--  测试互斥:另一个协程必须等待')
    task_other = asyncio.create_task(test_for_get_rock_other(rlock))
    task_outer = asyncio.create_task(test_for_get_rock_outer(rlock))

    result1, result2 = await asyncio.gather(task_other, task_outer)
    print(f'result1: {result1}')
    print(f'result2: {result2}')


asyncio.run(test_for_rlock())
预期输出:

--  测试可重入: outer → inner
** Outer acquired
    ** Inner acquired
** Back in outer
Result: inner

--  已经被释放了, 测试可以再次获取
** Outer acquired
    ** Inner acquired
** Back in outer
$$ [other task] sleep 0.1s

$$ Other task acquired, cost time: 4.0531158447265625e-06

--  测试互斥:另一个协程必须等待
$$ [other task] sleep 0.1s
** Outer acquired
    ** Inner acquired
** Back in outer

$$ Other task acquired, cost time: 0.1001734733581543
result1: other
result2: inner


五、关键注意事项与潜在问题

1. 必须绑定到 asyncio.Task

  • 使用 asyncio.current_task() 获取当前协程身份(而非线程 ID)
  • ❌ 错误做法:用 id(asyncio.current_task()) 或其他标识(Task 对象本身可哈希且唯一)

2. 异常安全:确保 release() 总是被调用

  • 必须通过 __aexit__ 自动释放(即只允许 async with
  • 手动调用 acquire/release 容易漏掉 release(尤其在异常路径)

3. 不要混用 with 和 async with

  • 实现中 __enter__ 主动报错,防止误用同步上下文管理器

4. 不支持跨事件循环

  • current_task() 依赖当前 loop,不能在多个 loop 间共享

5. 性能开销

  • 每次 acquire 都需检查 current_task()(但开销极小)
  • 比 asyncio.Lock 略慢,但远优于死锁

6. 不能用于 await 表达式直接返回

  • 设计为上下文管理器或显式 acquire/release,不是 awaitable

六、常见错误实现(避坑指南)

❌ 错误 1:用计数器但不绑定协程

# 危险!多个协程可能交错导致计数错误
class BadRLock:
    def __init__(self):
        self._count = 0
        self._lock = asyncio.Lock()
    
    async def acquire(self):
        if self._count > 0:
            self._count += 1  # ❌ 无协程绑定,多协程会混乱
        else:
            await self._lock.acquire()
            self._count = 1

❌ 错误 2:忘记异常安全

async def bad_usage(lock):
    await lock.acquire()
    raise ValueError("Oops!")  # release() 永远不会被调用!
    lock.release()

✅ 正确:始终用 async with


七、总结

手动实现 AsyncRLock 的核心是:

  1. 用 asyncio.Lock 保证跨协程互斥
  2. 用 current_task() 绑定持有者
  3. 用计数器支持可重入
  4. 通过 __aexit__ 保证异常安全

更多推荐