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

Python中Redis连接池未释放致连接耗尽超时问题排查

问题分析与解决

核心问题:连接池重复创建导致连接泄漏

你的get_redis_connection函数每次被调用时,都会创建一个新的ConnectionPool实例。这意味着每次执行async with cache.get_redis_connection(BAR_DB) as redis,都会生成一个独立的连接池,而非复用已有连接池内的连接。随着定时任务反复执行,大量独立连接池被创建,每个池都持有专属的Redis连接,最终导致connected_clients持续增长,直到触及系统或Redis的连接上限,触发超时错误。

另外,业务代码里存在笔误:grouped_data[key]应为grouped_bars[key],这会引发运行时错误,需优先修正。

修复方案

1. 全局复用连接池(关键修复)

修改get_redis_connection,让每个数据库对应唯一的连接池实例,避免重复创建:

from redis import asyncio as aioredis
from typing import Dict

# 全局存储已创建的连接池,key为db编号
_connection_pools: Dict[int, aioredis.ConnectionPool] = {}

def get_redis_connection(db_name: int) -> aioredis.Redis:
    redis_protocol = "rediss" if redis_settings.redis_ssl else "redis"
    
    # 检查该db的连接池是否已存在,不存在则创建
    if db_name not in _connection_pools:
        pool = aioredis.ConnectionPool.from_url(
            f"{redis_protocol}://:{redis_settings.redis_password}@{redis_settings.redis_hostname}:{redis_settings.redis_port}/{db_name}",
            connection_class=aioredis.Connection,
            max_connections=redis_settings.redis_pool_size,
            health_check_interval=HEALTH_CHECK_INTERVAL,
        )
        _connection_pools[db_name] = pool
    
    # 复用已有的连接池创建Redis客户端实例
    return aioredis.Redis(
        connection_pool=_connection_pools[db_name],
        auto_close_connection_pool=False,  # 禁止自动关闭全局连接池
        host=redis_settings.redis_hostname,
        port=redis_settings.redis_port,
        db=db_name,
        password=redis_settings.redis_password,
        ssl=redis_settings.redis_ssl,
        health_check_interval=HEALTH_CHECK_INTERVAL,
    )

2. 优化Pipeline使用效率

目前每个key都单独创建Pipeline,可优化为在一个Pipeline中批量处理所有key,减少网络交互次数:

async def populate_data_to_cache(bars: Iterator[BarsUnique]):
    grouped_bars = {}
    for bar in bars:
        key = f"{bar.code}:{bar.dir}"
        encoded = encode_data(bar)
        grouped_bars.setdefault(key, []).append(encoded)

    async with cache.get_redis_connection(BAR_DB) as redis:
        async with redis.pipeline(transaction=True) as p:
            for key, data_list in grouped_bars.items():
                await p.delete(key)
                await p.rpush(key, *data_list)
                await p.expire(key, timedelta(days=cache.DEFAULT_EXPIRATION_DAYS))
            # 一次性执行所有命令,降低网络往返开销
            await p.execute()

3. 验证连接池配置合理性

确保redis_settings.redis_pool_size设置合理(建议200-500之间,根据业务并发调整),不要超过Redis的maxclients(你的配置为7500)。同时检查Redis的连接超时配置,避免闲置连接长期未被回收。

额外建议

  • 使用redis-cli info clients监控连接状态,修复后观察connected_clients是否稳定在连接池大小范围内。
  • 控制定时任务的并发数,不要超过连接池的max_connections,避免出现连接等待超时。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 02:37:24