[Python3高阶编程] - 异步编程深度学习指南三:跟我一起实现一个 async.RLock
一、为什么需要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 分别验证:
- 可重入性:允许同一协程多次获取锁。
- 释放后可获取性:锁释放后,其他协程可以立即获取。
- 互斥性:同一时间只有一个协程能持有锁。
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)
测试逻辑:
-
流程:
test_for_get_rock_outer先获取RLock。- 它调用
test_for_get_rock_inner,后者再次尝试获取同一RLock。 - 如果
RLock是可重入的,内层函数将成功获取锁(Inner acquired)。
-
验证的特性:
- 可重入性:同一协程(
outer)可以多次嵌套获取锁,而不会阻塞自己。
- 可重入性:同一协程(
-
输出验证:
-- 测试互斥:另一个协程必须等待 $$ [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')
测试逻辑:
-
流程:
- 先运行
test_for_get_rock_outer,释放锁后。 - 启动
test_for_get_rock_other,它等待0.1s后尝试获取锁。 - 如果锁已释放,
Other task应立即获取锁(耗时极短)。
- 先运行
-
验证的特性:
- 锁的正确释放与可用性:锁释放后,其他协程可以立即获取。
-
输出验证:
-- 已经被释放了, 测试可以再次获取 ** 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}')
测试逻辑:
-
流程:
- 同时启动
task_outer(持有锁)和task_other(尝试获取锁)。 task_outer先获取锁,task_other必须等待直到锁被释放。task_outer执行完成后释放锁,task_other才能获取锁。
- 同时启动
-
验证的特性:
- 互斥性:同一时间只有一个协程能持有锁,其他协程必须等待。
-
输出验证:
-- 测试互斥:另一个协程必须等待 $$ [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 的执行时间。 |
代码可靠性验证的关键细节
-
可重入性:
test_for_get_rock_inner在test_for_get_rock_outer内部调用,且两者共享同一RLock。async with lock的嵌套使用不会阻塞,证明RLock允许可重入。
-
释放后可获取:
test_for_get_rock_other在锁释放后立即获取锁,耗时极短,说明锁状态正确。
-
互斥性:
- 两个协程同时运行时,
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的核心是:
- 用
asyncio.Lock保证跨协程互斥- 用
current_task()绑定持有者- 用计数器支持可重入
- 通过
__aexit__保证异常安全
更多推荐
所有评论(0)