
本文详解为何简单用 dict 缓存异步函数结果会失效,并提供一种基于 asyncio.future 的线程安全、协程友好的共享缓存方案,确保重复 ticker 仅发起一次 http 请求,同时规避失败响应被误缓存的风险。
本文详解为何简单用 dict 缓存异步函数结果会失效,并提供一种基于 asyncio.future 的线程安全、协程友好的共享缓存方案,确保重复 ticker 仅发起一次 http 请求,同时规避失败响应被误缓存的风险。
在使用 asyncio 和 aiohttp 批量获取金融数据(如 Yahoo Finance 的股价)时,一个常见需求是:当输入列表中存在重复 ticker(例如 ["AAPL", "XCN18679-USD", "XCN18679-USD"]),我们希望只对每个唯一 ticker 发起一次网络请求,后续调用直接复用结果——即实现“请求级去重”(request-level deduplication)。但若仅用普通 dict(如 ticker_price_dict = {})进行缓存,你会发现缓存始终未命中,所有重复 ticker 都触发了独立请求。根本原因在于:多个协程并发执行时,它们几乎同时检查 ticker in ticker_price_dict,此时缓存为空;于是全部进入下载逻辑,导致重复请求——这并非缓存失效,而是典型的竞态条件(race condition)。
要解决这个问题,不能依赖“事后写入”的朴素缓存,而需采用 “预占位 + 异步等待”策略:首次遇到某个 ticker 时,立即在字典中存入一个 asyncio.Future 对象作为占位符(placeholder),表示“该 ticker 正在加载中”;后续协程再查到这个 Future,就直接 await 它,自动同步到最终结果。这样既保证了单次请求,又天然支持并发等待。
以下是改造后的核心缓存逻辑(已适配你的 fetch_close_price 结构):
import asyncio
import pandas as pd
from datetime import datetime
# 假设这些辅助函数已定义
# def to_midnight(dt: datetime, naive: bool) -> datetime: ...
# BASE_URL, DATA_URL_PART 已定义
async def fetch_close_price(
aiosession: aiohttp.ClientSession,
ticker: str,
start: datetime,
end: datetime,
ticker_price_dict: dict
):
print(f'Checking cache for ticker {ticker}')
# 尝试从缓存获取:可能是 Future(正在加载)或 DataFrame(已就绪)
cached = ticker_price_dict.get(ticker)
if cached is not None:
if isinstance(cached, asyncio.Future):
# 其他协程已在请求中,等待其完成
print(f'ticker {ticker} is already being fetched — awaiting result...')
df = await cached
return df
else:
# 已缓存有效 DataFrame(注意:不缓存空/失败结果!)
print(f'Found cached result for ticker {ticker}')
return cached
# 缓存未命中 → 创建 Future 占位,并存入字典
loop = asyncio.get_running_loop()
future = loop.create_future()
ticker_price_dict[ticker] = future
try:
# 执行实际 HTTP 请求(原逻辑)
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:
data = await response.json()
# ... 解析 close_prices 和 dates(略)...
df = pd.DataFrame({ticker: close_prices}, index=dates)
# ✅ 成功:将结果写入缓存(替换 Future),并标记 future 完成
ticker_price_dict[ticker] = df
future.set_result(df)
print(f"Successfully cached {ticker}")
return df
else:
# ❌ 失败:不缓存!清除占位符,避免后续协程等待失败
ticker_price_dict.pop(ticker, None)
future.cancel()
print(f"Failed to fetch {ticker}, status {response.status}")
return pd.DataFrame() # 或抛出异常,按需处理
except Exception as e:
# 清理失败状态
ticker_price_dict.pop(ticker, None)
future.cancel()
raise e
# 其余函数保持不变(download_close_prices / download_close_prices_all)
关键设计要点说明:
- ✅ Future 占位:
future = loop.create_future()是轻量级同步对象,无性能开销;ticker_price_dict[ticker] = future立即阻断其他协程的重复请求路径。 - ✅ 失败不缓存:HTTP 错误或异常发生时,主动
pop字典并cancel()Future,确保下次重试不受干扰。 - ✅ 类型安全判断:通过
isinstance(cached, asyncio.Future)区分“加载中”与“已就绪”,比检查类名更健壮。 - ⚠️ 注意缓存粒度:此方案缓存的是 单次请求结果(对应特定
ticker+start+end组合)。若需支持时间范围参数化缓存,应将(ticker, start, end)构造为复合 key(如f"{ticker}_{start.date()}_{end.date()}")。
最后提醒:如果你的业务场景中重复 ticker 出现频率不高,或可预处理去重(如 list(set(tickers))),那最简方案仍是前置 dedupe —— 但当输入来自不同来源、无法提前合并时,上述 Future 占位法就是高并发异步环境下的标准解法。










