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

异步函数使用字典缓存请求结果未生效,重复请求始终触发下载的问题求助

异步函数使用字典缓存请求结果未生效,重复请求始终触发下载的问题求助

看起来你遇到的是并发场景下的缓存竞争问题,咱们一步步来分析和解决:

问题根源

你用asyncio.gather同时启动了所有ticker的请求任务,当存在多个相同ticker时,这些协程几乎是同时去检查ticker_price_dict缓存。这时候第一个协程还没来得及把下载好的结果存入字典,其他协程就已经判定ticker不在缓存里了,所以都会触发重复的下载请求。

另外你提到的「下载失败不缓存」的需求,咱们也可以在解决方案里一起处理。


解决方案

方案1:用锁保护缓存操作,避免并发竞争

给缓存字典加一个asyncio.Lock,确保同一时间只有一个协程能检查和更新缓存。这样当第一个协程在处理某个ticker时,其他相同ticker的协程会等待锁释放,再去检查缓存,这时候结果已经存入字典,就不会重复下载了。

修改后的代码片段如下:

async def fetch_close_price(aiosession: aiohttp.ClientSession,
                            ticker: str, start: datetime, end: datetime, ticker_price_dict, cache_lock):
    # 先加锁检查缓存
    async with cache_lock:
        print(f'Checking cache for ticker {ticker}')
        if ticker in ticker_price_dict:
            print(f'Found cache for ticker {ticker}')
            return ticker_price_dict[ticker]

    # 缓存未命中,发起请求
    params = {
        'period1': int(to_midnight(start, naive=False).timestamp()),
        'period2': int(to_midnight(end, naive=False).timestamp()),
        'interval': '1d',
        'includeAdjustedClose': 'false'
    }

    async with aiosession.get(DATA_URL_PART.format(ticker=ticker), params=params) as response:
        if response.status == 200:
            print(f"Downloading {ticker}", response.status)
            data = await response.json()
            # ... 省略你处理数据生成close_prices和dates的代码 ...
            df = pd.DataFrame({ticker: close_prices}, index=dates)
            # 写入缓存时再次加锁
            async with cache_lock:
                ticker_price_dict[ticker] = df
            print(f"Downloaded and setting {ticker}", response.status)
            return df
        else:
            print(f"Failed to fetch data for {ticker}", response.status)
            # 下载失败,不存入缓存
            return pd.DataFrame()

async def download_close_prices(aiosession: aiohttp.ClientSession,
                                tickers: list[str], start: datetime, end: datetime):
    all_ticker_close = pd.DataFrame()
    ticker_price_dict = {}
    # 创建缓存锁
    cache_lock = asyncio.Lock()
    tasks = [fetch_close_price(aiosession, ticker, start, end, ticker_price_dict, cache_lock) for ticker in tickers]
    results = await asyncio.gather(*tasks)
    # 合并所有结果(你原来的代码漏了这一步)
    all_ticker_close = pd.concat(results, axis=1)
    return all_ticker_close

方案2:缓存「正在进行的任务」,更高效避免重复请求

这种方式会在缓存里存储两种值:要么是已经下载好的DataFrame,要么是一个asyncio.Future对象(代表正在进行的请求任务)。当有重复ticker请求时,直接等待正在进行的任务完成,而不是重新发起请求,性能会更优。

示例代码思路:

async def fetch_close_price(aiosession: aiohttp.ClientSession,
                            ticker: str, start: datetime, end: datetime, ticker_price_dict, cache_lock):
    async with cache_lock:
        if ticker in ticker_price_dict:
            value = ticker_price_dict[ticker]
            if isinstance(value, asyncio.Future):
                # 有正在进行的请求,等待它完成
                print(f'Waiting for in-progress download of {ticker}')
                return await value
            else:
                print(f'Found cache for ticker {ticker}')
                return value
        # 创建Future标记该ticker正在下载
        future = asyncio.get_event_loop().create_future()
        ticker_price_dict[ticker] = future

    try:
        params = {
            'period1': int(to_midnight(start, naive=False).timestamp()),
            'period2': int(to_midnight(end, naive=False).timestamp()),
            'interval': '1d',
            'includeAdjustedClose': 'false'
        }

        async with aiosession.get(DATA_URL_PART.format(ticker=ticker), params=params) as response:
            if response.status == 200:
                print(f"Downloading {ticker}", response.status)
                data = await response.json()
                # ... 处理数据 ...
                df = pd.DataFrame({ticker: close_prices}, index=dates)
                async with cache_lock:
                    ticker_price_dict[ticker] = df
                future.set_result(df)
                return df
            else:
                print(f"Failed to fetch data for {ticker}", response.status)
                async with cache_lock:
                    del ticker_price_dict[ticker]  # 移除正在下载的标记
                future.set_result(pd.DataFrame())
                return pd.DataFrame()
    except Exception as e:
        async with cache_lock:
            del ticker_price_dict[ticker]
        future.set_exception(e)
        raise

额外注意事项

  1. 下载失败的处理:一定要确保不把无效的空DataFrame存入缓存,或者后续检查缓存时判断数据有效性,避免重复使用失败结果。
  2. 你原来的download_close_prices函数里,await asyncio.gather(*tasks)之后没有合并结果的逻辑,记得补上pd.concat(results, axis=1),否则返回的是空DataFrame。

备注:内容来源于stack exchange,提问作者SAK

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 08:53:05