Python asyncio协程读取Redis时待处理任务被销毁错误如何解决?
问题根因分析
报错含义
Task was destroyed but it is pending! 是Python asyncio的典型报错,含义是仍在等待IO完成的任务没有被正确持有,还没执行完就被垃圾回收器销毁了,你的两个报错分别对应消费者消费任务被意外销毁、Redis连接底层的IO任务被意外销毁两个问题。
具体问题定位
消费者代码问题
- Redis连接释放逻辑错误:你每次循环都新建Redis连接,调用的
ds_handle.close()是异步方法,你没有加await就直接进入下一轮循环,连接释放的逻辑还没执行完,连接对象就被销毁,底层的IO任务还处于pending状态就被回收,触发第二个报错。 - 异常处理逻辑缺失:你只捕获了异常打日志,没有做异常后的重试和资源清理,一旦Redis断连抛出未处理的异常,会直接冒泡到事件循环层,导致
consume任务被直接取消,触发第一个报错,对应队列的处理逻辑自然就停止运行。 - 频繁新建连接的性能问题:每次循环新建销毁连接会给Redis带来不必要的压力,也大幅提升了断连报错的概率。
订阅者代码问题
- 无重连逻辑:Redis的pubsub是长连接,服务端默认会主动断开长时间空闲的连接,你没有加断连后的重连逻辑,连接一断整个订阅任务就直接崩溃。
- 任务无持有引用:你用
asyncio.ensure_future创建的reader任务只存在于局部变量中,如果外层subscriber协程因为异常退出,reader任务没有被其他地方持有,会直接被GC回收,触发pending任务销毁的报错。 - 无异常捕获逻辑:订阅逻辑没有任何异常捕获,连接超时、网络波动都会导致整个订阅协程直接退出。
修复方案
消费者代码修复
- 复用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
相关产品推荐
相关产品推荐

