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

如何将Redis Queue失败任务重试至不同资源队列(Python RQ)

问题分析

你的核心需求是让RQ任务失败后自动升级到资源更高的队列重试,同时保留原Job ID,但现有代码存在两个关键问题:

  1. 重试队列由固定环境变量RETRY_QUEUE决定,无法根据当前任务所在队列动态升级(比如small队列的任务不能自动转到medium)
  2. 队列判断逻辑硬编码,没有形成灵活的升级映射,扩展性差

注意:RQ的job.retry()方法本身会复用原Job ID,完全符合你保留ID的需求,不需要重新创建任务。

解决方案
  1. 建立队列升级映射表,明确每个队列对应的下一级重试队列
  2. 从任务本身获取当前所在队列名称,动态匹配目标重试队列
  3. 修正重试逻辑,确保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对象,避免潜在的连接问题
验证步骤
  1. 分别启动三个Worker实例,设置对应的环境变量:
    • Small Worker:REDIS_QUEUE=small
    • Medium Worker:REDIS_QUEUE=medium
    • Large Worker:REDIS_QUEUE=large
  2. 提交一个任务到small队列,人为触发失败(比如模拟内存不足、超时)
  3. 检查Redis中的队列数据,确认任务是否从small移到medium队列
  4. 任务在medium队列再次失败后,确认是否移到large队列
  5. 查看任务的Job ID,验证重试过程中ID保持不变

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 04:46:35