解决aioredis中RuntimeError: await wasn't used with future错误
问题描述
多个运行在独立线程中的Worker(如MyFirstWorker),其中同步方法_process需要执行异步Redis的hget操作,出现RuntimeError: await wasn't used with future错误。直接用await会因_process是同步方法报错,用asyncio.run也无法解决问题。
问题根源
- 异步Redis连接初始化错误:Worker的
__init__中通过asyncio.run(get_async_redis_connection(...))创建的连接,绑定的事件循环在asyncio.run执行完毕后会被关闭,后续使用该连接时会引用已关闭的循环,触发错误。 - 冗余的连接调用:
async_get_redis_data中使用async with self._redis_async_connection.client()是多余操作,aioredis的Redis实例可直接调用异步方法。 - 线程循环冲突:多线程环境下,asyncio事件循环是线程绑定的,跨线程共享循环或连接会导致异步任务执行异常。
修复方案
1. 修正异步Redis连接的线程安全初始化
使用线程本地存储为每个线程创建独立的事件循环和Redis连接,避免跨线程共享资源:
import asyncio import threading from aioredis import Redis, ConnectionPool class Worker(): def __init__(self, azure_connection_string: str, redis_connection_string: str): super().__init__() self._redis_connection = get_redis_connection(redis_connection_string) self._redis_connection_string = redis_connection_string self._connection_string = azure_connection_string # 线程本地存储,保存当前线程的异步Redis连接和事件循环 self._local = threading.local() def _get_async_redis_connection(self): if not hasattr(self._local, "redis"): # 为当前线程创建独立事件循环 loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) # 异步初始化Redis连接 async def init_conn(): hostport, *options = self._redis_connection_string.split(",") host, _, port = hostport.partition(":") _, _, password = options[0].partition("=") redis_url = f'rediss://:{password}@{host}:{port}' pool = ConnectionPool.from_url(redis_url, ssl=True) return Redis(connection_pool=pool) self._local.redis = loop.run_until_complete(init_conn()) self._local.loop = loop return self._local.redis
2. 简化异步Redis操作方法
去掉冗余的client()调用,直接使用线程专属的异步连接:
class MyFirstWorker(Worker): def __init__(self, azure_connection_string: str, redis_connection_string: str): super().__init__(azure_connection_string, redis_connection_string) async def async_get_redis_data(self, key, opt): redis = self._get_async_redis_connection() windows = await redis.hget(key, "windows") doors = await redis.hget(key, "doors") return windows, doors
3. 在同步方法中正确执行异步任务
优先复用当前线程的事件循环,避免重复创建循环导致冲突:
def _process(self, opt): key = str(os.getenv("KEY")) # 获取或创建当前线程的事件循环 try: loop = asyncio.get_running_loop() except RuntimeError: loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) # 运行异步任务并获取结果 windows, doors = loop.run_until_complete(self.async_get_redis_data(key, opt)) matching_dataclass = MatchingPayload(json.loads(windows)["PayloadFields"]) logger.info(f"Matching completed") put_redis_data( self._redis_connection, key, "match:results", matching_dataclass )
4. 修复同步Redis操作的参数错误
put_redis_data函数存在参数名不匹配问题,修正如下:
def put_redis_data(redis_connection, hash_key, hash_name, payload): ttl = int(os.environ.get("TTL", 12)) redis_connection.hset(hash_key, hash_name, json.dumps(payload)) redis_connection.expire(hash_key, timedelta(hours=ttl))
关键注意事项
- asyncio事件循环是线程绑定的,必须为每个线程创建独立的循环和Redis连接。
- aioredis的
Redis实例线程不安全,不能跨线程共享使用。 - 避免在同步方法中频繁创建新的事件循环,优先复用当前线程的循环以提升性能。
内容的提问来源于stack exchange,提问作者Evgeniy
相关产品推荐
相关产品推荐

