异步函数使用字典缓存请求结果未生效,重复请求始终触发下载的问题求助
异步函数使用字典缓存请求结果未生效,重复请求始终触发下载的问题求助
看起来你遇到的是并发场景下的缓存竞争问题,咱们一步步来分析和解决:
问题根源
你用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
额外注意事项
- 下载失败的处理:一定要确保不把无效的空DataFrame存入缓存,或者后续检查缓存时判断数据有效性,避免重复使用失败结果。
- 你原来的
download_close_prices函数里,await asyncio.gather(*tasks)之后没有合并结果的逻辑,记得补上pd.concat(results, axis=1),否则返回的是空DataFrame。
备注:内容来源于stack exchange,提问作者SAK
相关产品推荐
相关产品推荐

