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

Python asyncio协程读取Redis时待处理任务被销毁错误如何解决?

问题根因分析

报错含义

Task was destroyed but it is pending! 是Python asyncio的典型报错,含义是仍在等待IO完成的任务没有被正确持有,还没执行完就被垃圾回收器销毁了,你的两个报错分别对应消费者消费任务被意外销毁、Redis连接底层的IO任务被意外销毁两个问题。

具体问题定位

消费者代码问题

  1. Redis连接释放逻辑错误:你每次循环都新建Redis连接,调用的ds_handle.close()是异步方法,你没有加await就直接进入下一轮循环,连接释放的逻辑还没执行完,连接对象就被销毁,底层的IO任务还处于pending状态就被回收,触发第二个报错。
  2. 异常处理逻辑缺失:你只捕获了异常打日志,没有做异常后的重试和资源清理,一旦Redis断连抛出未处理的异常,会直接冒泡到事件循环层,导致consume任务被直接取消,触发第一个报错,对应队列的处理逻辑自然就停止运行。
  3. 频繁新建连接的性能问题:每次循环新建销毁连接会给Redis带来不必要的压力,也大幅提升了断连报错的概率。

订阅者代码问题

  1. 无重连逻辑:Redis的pubsub是长连接,服务端默认会主动断开长时间空闲的连接,你没有加断连后的重连逻辑,连接一断整个订阅任务就直接崩溃。
  2. 任务无持有引用:你用asyncio.ensure_future创建的reader任务只存在于局部变量中,如果外层subscriber协程因为异常退出,reader任务没有被其他地方持有,会直接被GC回收,触发pending任务销毁的报错。
  3. 无异常捕获逻辑:订阅逻辑没有任何异常捕获,连接超时、网络波动都会导致整个订阅协程直接退出。
修复方案

消费者代码修复

  • 复用Redis连接池,不要每次循环新建连接,连接池本身会处理断连重试、连接复用的逻辑
  • 释放连接时添加await,或者改用异步上下文管理器管理连接生命周期
  • 异常捕获后添加重试逻辑,避免异常冒泡导致消费任务被取消,示例修改逻辑:
while True:
    try:
        async with ds.get_datastore_handle(ds.get_uri(conf=configuration)) as ds_handle:
            # 用blpop替代lpop,避免空轮询,减少无效请求
            token = await ds_handle.blpop(queue, timeout=MAX)
            if token is not None:
                token = token[1] # blpop返回的是(队列名, 值)的元组
                result = await processor.consume(json.loads(token), ds_handle)
                status = await processor.relay(result, ds_handle)
                logger.debug(status)
    except Exception as e:
        logger.error(f"消费出错: {e}", exc_info=True) # exc_info会自动打印堆栈,不用单独打with_traceback
        await asyncio.sleep(randint(MIN, MAX)) # 出错后休眠再重试,避免疯狂报错

订阅者代码修复

  • 加外层重连循环,断连后自动重新订阅
  • 持有创建的reader任务引用,避免被GC回收
  • 添加异常捕获,出错后主动清理资源再重试,示例修改逻辑:
async def subscriber(conf: dict, channel: str, processor: processor.Strategy) -> None:
    logger = Log().get_logger(f"subscriber_{channel}", conf['logFolder'], conf['logFormat'], conf['USE'])
    # 外层循环实现重连
    while True:
        tsk = None
        ds = None
        try:
            ds_uri = ds_linker.get_uri(conf=conf)
            ds = await ds_linker.get_datastore_handle(ds_uri)
            pub = await ds.subscribe(channel)
            ch = pub[0]
            async def reader(ch):
                while await ch.wait_message():
                    msg = await ch.get_json()
                    await processor.handle_message(msg=msg)
            tsk = asyncio.ensure_future(reader(ch))
            await tsk
        except Exception as e:
            logger.error(f"订阅出错: {e}", exc_info=True)
            # 出错后主动清理资源
            if tsk and not tsk.done():
                tsk.cancel()
                try:
                    await tsk
                except asyncio.CancelledError:
                    pass
            if ds:
                await ds.close()
            # 休眠后重连
            await asyncio.sleep(randint(1,5))

通用优化

  • 所有创建的asyncio任务都要在全局或者上层作用域持有引用,不要让任务仅存在于局部变量中
  • 给Redis连接配置心跳参数,避免服务端主动断开空闲连接,比如创建连接时添加heartbeat=30参数
  • 如果要提升可靠性,可以后续替换为Redis Streams实现,天然支持消费者组、消息确认、死信队列等能力,比当前list+pubsub的方案更稳定,能避免消息丢失,也不需要自己实现很多重试逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 15:27:04