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

如何实现async aiohttp线程安全类?Tornado缓存线程安全排查与优化

1. 如何实现基于asyncio+aiohttp的线程安全类

要实现线程安全的asyncio+aiohttp类,核心需解决两个问题,同时遵循asyncio运行规则:

  • 单例初始化安全:用线程锁保护单例创建逻辑,避免多线程同时初始化导致重复创建或状态重置。
  • 共享状态安全:对字典、列表等非线程安全的共享数据,使用线程锁(跨线程场景)或协程锁(单线程协程场景)保护所有读写操作。
  • aiohttp对象隔离:ClientSession属于特定事件循环,不能跨线程调用;若在多线程环境下使用,需为每个线程独立创建事件循环和Session,或通过loop.run_in_executor将阻塞操作委托给线程池。
2. 你的代码中cls._instance._cache = {}的线程安全性及优化方案

结论:该操作在Tornado多线程环境下不安全,存在两个核心风险:

  1. 单例初始化竞争:__new__方法未加锁,多线程同时调用时可能多次进入if cls._instance is None分支,导致_cache被重复重置为空字典。
  2. 缓存读写竞争:Python原生字典不是线程安全的数据结构,多线程/多协程并发读写时会引发数据竞争,可能出现KeyError、数据覆盖或缓存状态不一致的问题。

优化方案(修改后的代码):

import logging
import aiohttp
import time
import threading

# Constants
DEFAULT_TIMEOUT = 20
MAX_ERRORS = 3
HTTP_READ_TIMEOUT = 1

class HTTPRequestCache:
    _instance = None
    _init_lock = threading.Lock()  # 保护单例初始化的线程锁

    def __new__(cls):
        with cls._init_lock:
            if cls._instance is None:
                cls._instance = super().__new__(cls)
                cls._instance._cache = {}
                cls._instance._time_out = DEFAULT_TIMEOUT
                cls._instance._http_read_timeout = HTTP_READ_TIMEOUT
                cls._instance._cache_lock = threading.Lock()  # 保护缓存读写的线程锁
        return cls._instance

    async def _fetch_update(self, url):
        try:
            async with aiohttp.ClientSession() as session:
                logging.info(f"Fetching {url}")
                async with session.get(url, timeout=self._http_read_timeout) as resp:
                    resp.raise_for_status()
                    resp_data = await resp.json()

                    cached_at = time.time()
                    cache_entry = {
                        "cached_at": cached_at,
                        "config": resp_data,
                        "errors": 0
                    }
                    with self._cache_lock:
                        self._cache[url] = cache_entry
                    logging.info(f"Updated cache for {url}")
        except aiohttp.ClientError as e:
            logging.error(f"Error occurred while updating cache for {url}: {e}")

    async def get(self, url):
        need_update = False
        # 加锁判断缓存状态,避免读写竞争
        with self._cache_lock:
            if url not in self._cache or self._cache[url]["cached_at"] < time.time() - self._time_out:
                need_update = True
        if need_update:
            await self._fetch_update(url)
        # 加锁读取缓存,确保数据一致性
        with self._cache_lock:
            return self._cache.get(url, {}).get("config")

关键修改说明:

  1. 单例初始化锁:新增_init_lock,确保多线程环境下只会初始化一次单例,避免_cache被重复重置。
  2. 缓存读写锁:新增_cache_lock,所有对_cache的读写操作都在锁的保护下执行,彻底解决数据竞争问题。
  3. 逻辑拆分:将缓存状态判断与更新操作拆分,减少锁的持有时间,兼顾安全性和性能。

补充场景适配:

如果你的Tornado环境是纯单线程事件循环(默认运行模式),可以将_cache_lock替换为asyncio.Lock,用协程锁替代线程锁,性能会更轻量;但如果开启了多线程worker,threading.Lock是更稳妥的选择,支持跨线程同步。

内容的提问来源于stack exchange,提问作者Mini

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 23:00:39