Python asyncio并发写入任务频繁取消,如何实现交替执行?
问题分析与修复
问题原因
- 删除任务被取消的本质是缓存值被覆盖:两个写入任务都操作同一个
testkey,每次调用setex时都会覆盖之前的缓存值,并取消旧的删除任务(旧值已被替换,对应的删除任务失去意义)。由于两个任务交替写入,后执行的写入操作会取消前一次写入对应的删除任务,导致日志中频繁出现删除任务被取消的信息。 - 写入操作已实现交替执行:从输出日志可以看到
test1->0、test2->0、test1->1、test2->1交替获取锁并完成写入,说明写入任务已经实现交替执行,被取消的是删除任务而非写入任务。
修复方案
如果希望避免频繁取消删除任务,可以修改删除任务的逻辑:让删除任务在执行时先检查当前缓存值是否与任务创建时的值一致,仅在值未被修改时执行删除操作。这样无需在每次setex时取消旧任务,旧任务会自然结束但不会误删新值。同时移除多余的全局_write_lock(per-key锁已足够保证线程安全)。
修改后的代码如下:
import asyncio from collections import defaultdict from contextlib import asynccontextmanager from copy import deepcopy from datetime import datetime class LocalCacheUserToken: _g: dict[str, dict] = {} _tasks: dict[str, asyncio.Task] = {} _locks: defaultdict[str, asyncio.Lock] = defaultdict(asyncio.Lock) @classmethod @asynccontextmanager async def acquire_lock(cls, key: str): async with cls._locks[key]: yield @classmethod async def get(cls, key: str): async with cls.acquire_lock(key): value = deepcopy(cls._g.get(key, None)) return value @classmethod async def setex(cls, key: str, time: int, value: dict): async with cls.acquire_lock(key): print(f"[{datetime.now()}] {value=} get lock") cls._g[key] = value # 不再强制取消旧任务,让其自行判断是否执行删除 task = asyncio.create_task(cls._delay_delete(key, time, deepcopy(value))) cls._tasks[key] = task @classmethod async def _delay_delete(cls, key, time, original_value): try: await asyncio.sleep(time) async with cls.acquire_lock(key): current_value = cls._g.get(key) # 仅当值未被修改时执行删除 if current_value == original_value: del cls._g[key] cls._tasks.pop(key, None) except asyncio.CancelledError: print(f"[{datetime.now()}] Deletion task for {key=}, {original_value=} was cancelled") async def p(a: LocalCacheUserToken): while True: print(f"[{datetime.now()}] {await a.get('test')=}") await asyncio.sleep(1) async def write(a: LocalCacheUserToken, k: str): for i in range(20): await a.setex("test", 5, {'value': f'{k}->{i}'}) await asyncio.sleep(0.01) # 示例使用 async def main(): asyncio.create_task(p(LocalCacheUserToken)) asyncio.create_task(write(LocalCacheUserToken, k='test1')) asyncio.create_task(write(LocalCacheUserToken, k='test2')) await asyncio.sleep(10) asyncio.run(main())
如果期望两个写入任务的结果都被保留,需要让它们操作不同的key(如test1和test2),这样每个key的删除任务都能正常执行,不会被其他写入操作干扰。
内容的提问来源于stack exchange,提问作者user25772377
相关产品推荐
相关产品推荐

