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

解决aioredis中RuntimeError: await wasn't used with future错误

问题描述

多个运行在独立线程中的Worker(如MyFirstWorker),其中同步方法_process需要执行异步Redis的hget操作,出现RuntimeError: await wasn't used with future错误。直接用await会因_process是同步方法报错,用asyncio.run也无法解决问题。

问题根源
  1. 异步Redis连接初始化错误:Worker的__init__中通过asyncio.run(get_async_redis_connection(...))创建的连接,绑定的事件循环在asyncio.run执行完毕后会被关闭,后续使用该连接时会引用已关闭的循环,触发错误。
  2. 冗余的连接调用:async_get_redis_data中使用async with self._redis_async_connection.client()是多余操作,aioredis的Redis实例可直接调用异步方法。
  3. 线程循环冲突:多线程环境下,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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 08:05:56