You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python asyncio并发写入任务频繁取消,如何实现交替执行?

问题分析与修复

问题原因

  1. 删除任务被取消的本质是缓存值被覆盖:两个写入任务都操作同一个test key,每次调用setex时都会覆盖之前的缓存值,并取消旧的删除任务(旧值已被替换,对应的删除任务失去意义)。由于两个任务交替写入,后执行的写入操作会取消前一次写入对应的删除任务,导致日志中频繁出现删除任务被取消的信息。
  2. 写入操作已实现交替执行:从输出日志可以看到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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.21 19:22:02