如何将Redis Queue失败任务重试至不同资源队列(Python RQ)
问题分析
你的核心需求是让RQ任务失败后自动升级到资源更高的队列重试,同时保留原Job ID,但现有代码存在两个关键问题:
- 重试队列由固定环境变量
RETRY_QUEUE决定,无法根据当前任务所在队列动态升级(比如small队列的任务不能自动转到medium) - 队列判断逻辑硬编码,没有形成灵活的升级映射,扩展性差
注意:RQ的job.retry()方法本身会复用原Job ID,完全符合你保留ID的需求,不需要重新创建任务。
解决方案
- 建立队列升级映射表,明确每个队列对应的下一级重试队列
- 从任务本身获取当前所在队列名称,动态匹配目标重试队列
- 修正重试逻辑,确保
retries_left正确递减,同时覆盖异常抛出和Worker被杀死两种场景
修改后代码
import os import struct from rq import Connection, Worker, Queue, Job # 假设你已经初始化好redis_connection变量 # redis_connection = Redis(...) # 队列升级映射:key是当前队列,value是下一级重试队列 QUEUE_UPGRADE_MAP = { 'small': 'medium', 'medium': 'large', 'large': 'large' # large队列任务失败后不再升级,可根据需求调整 } def get_target_retry_queue(job: Job) -> str: """根据当前任务所在队列获取目标重试队列""" current_queue = job.queue_name # 获取任务当前所在队列名称 return QUEUE_UPGRADE_MAP.get(current_queue, current_queue) def retry_job(job: Job): target_queue_name = get_target_retry_queue(job) target_queue = Queue(target_queue_name, connection=job.connection) # 调用retry方法,指定目标队列,复用原Job ID job.retry(queue=target_queue) def retry_handler(job, exc_type, exception, traceback): if job.retries_left > 0: retry_job(job) return False # 告诉RQ不再执行后续异常处理器 return True # 允许默认处理(比如将任务移到failed队列) def work_horse_killed_handler(job: Job, retpid: int, ret_val: int, rusage: struct.struct_rusage): # ECS终止Worker时,同样升级队列重试 if job.retries_left > 0: retry_job(job) return False if __name__ == '__main__': listen_queue = os.getenv('REDIS_QUEUE', 'small') with Connection(redis_connection): worker = Worker( queues=[Queue(listen_queue, connection=redis_connection)], exception_handlers=[retry_handler], work_horse_killed_handler=work_horse_killed_handler ) worker.work(with_scheduler=True) print('Worker has exited')
关键修改说明
- QUEUE_UPGRADE_MAP:灵活定义队列升级规则,后续新增队列只需修改映射表,无需改动业务逻辑
- get_target_retry_queue:从
job.queue_name获取任务当前所在队列,动态匹配目标重试队列,解决了原代码固定重试队列的问题 - 重试逻辑:确保只有当
retries_left > 0时才执行重试,避免无限重试;job.retry()会自动递减retries_left,无需手动修改 - Worker初始化:将
map(Queue, listen)改为显式创建Queue对象,避免潜在的连接问题
验证步骤
- 分别启动三个Worker实例,设置对应的环境变量:
- Small Worker:
REDIS_QUEUE=small - Medium Worker:
REDIS_QUEUE=medium - Large Worker:
REDIS_QUEUE=large
- Small Worker:
- 提交一个任务到small队列,人为触发失败(比如模拟内存不足、超时)
- 检查Redis中的队列数据,确认任务是否从small移到medium队列
- 任务在medium队列再次失败后,确认是否移到large队列
- 查看任务的Job ID,验证重试过程中ID保持不变
内容的提问来源于stack exchange,提问作者cary
相关产品推荐
相关产品推荐

