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

如何实现RQ任务的暂停与恢复?多任务轮询执行需求求助

实现RQ任务的分段执行与切换恢复

要实现你描述的“每个任务执行5分钟后暂停、切换下一个,最终恢复所有任务完成”的需求,由于RQ本身没有原生的任务暂停/恢复机制,我们需要通过断点续传+Redis状态跟踪+调度控制的方式来实现。以下是具体方案:

核心思路

  1. 拆分长任务:将原本10分钟的任务拆分为可中断的分段执行逻辑,任务内部定期检查暂停信号。
  2. 状态跟踪:用Redis存储每个任务的执行进度和暂停标记,确保暂停后能从断点恢复。
  3. 调度控制:实现一个调度任务,负责依次启动任务、定时暂停、最后统一恢复所有任务。

代码实现

1. 修改任务函数,支持断点续传与暂停检查

import time
import redis
from rq import Queue

# 初始化Redis连接
r = redis.Redis(host='localhost')
q = Queue(connection=r)

def count_words(url, task_id):
    # Redis键定义:跟踪进度和暂停状态
    progress_key = f"task:{task_id}:progress"
    pause_key = f"task:{task_id}:pause"
    
    # 读取已保存的执行进度(如果有)
    saved_progress = r.get(progress_key)
    if saved_progress:
        elapsed_time = time.time() - float(saved_progress)
    else:
        elapsed_time = 0.0
        r.set(progress_key, time.time())  # 记录任务开始时间
    
    total_duration = 600  # 总执行时长:10分钟
    chunk_duration = 300   # 单次最大执行时长:5分钟
    remaining_time = total_duration - elapsed_time

    while remaining_time > 0:
        # 检查是否收到暂停信号
        if r.exists(pause_key):
            return f"Task {task_id} paused. Remaining time: {remaining_time:.0f}s"
        
        # 执行本次任务片段(用sleep模拟实际业务逻辑)
        run_duration = min(chunk_duration, remaining_time)
        time.sleep(run_duration)
        
        # 更新进度
        elapsed_time += run_duration
        remaining_time = total_duration - elapsed_time
        r.set(progress_key, time.time() - elapsed_time)  # 保存当前进度基准时间
    
    # 任务完成,清理Redis中的状态键
    r.delete(progress_key)
    r.delete(pause_key)
    return f"Task {task_id} completed. Processed URL: {url}"

2. 实现任务调度逻辑

这个调度任务负责控制三个任务的执行、暂停与恢复流程:

def task_scheduler(task_list):
    """
    task_list: 列表,元素为(任务URL, 任务ID)的元组
    """
    # 依次启动并暂停每个任务
    for url, task_id in task_list:
        # 启动当前任务
        q.enqueue(
            count_words,
            args=(url, task_id),
            job_timeout='2h',
            result_ttl=1000
        )
        # 等待5分钟,让任务执行满时长
        time.sleep(300)
        # 标记任务暂停
        r.set(f"task:{task_id}:pause", "1")
    
    # 恢复所有暂停的任务:删除暂停标记,重新加入队列续执行
    for url, task_id in task_list:
        r.delete(f"task:{task_id}:pause")
        q.enqueue(
            count_words,
            args=(url, task_id),
            job_timeout='2h',
            result_ttl=1000
        )

3. 修改FastAPI接口,收集任务并触发调度

from fastapi import FastAPI
from fastapi.responses import JSONResponse
import uuid

app = FastAPI()
pending_tasks = []  # 收集待调度的任务

def success_return(data):
    return {"status": "success", "data": data}

@app.get("/add")
async def add_task(url: str):
    # 生成唯一任务ID
    task_id = str(uuid.uuid4())
    pending_tasks.append((url, task_id))
    
    # 当收集到3个任务时,触发调度
    if len(pending_tasks) == 3:
        q.enqueue(task_scheduler, args=(pending_tasks,))
        pending_tasks.clear()  # 清空待调度列表
    
    return JSONResponse(content=success_return({
        "length_queue": len(q),
        "task_id": task_id
    }))

关键注意事项

  • RQ Worker的容错:当任务因暂停主动返回时,RQ会将其标记为finished,我们通过返回值区分“暂停”和“完成”,后续恢复时重新加入队列即可。
  • Redis状态一致性:任务执行过程中进度和暂停标记的读写操作要保证原子性,示例中简单的set/get即可满足需求,复杂场景可使用Redis事务。
  • 调度任务的可靠性:调度任务本身需设置足够长的超时时间,比如job_timeout='30m',确保能完成整个调度流程。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 09:53:11