如何实现RQ任务的暂停与恢复?多任务轮询执行需求求助
实现RQ任务的分段执行与切换恢复
要实现你描述的“每个任务执行5分钟后暂停、切换下一个,最终恢复所有任务完成”的需求,由于RQ本身没有原生的任务暂停/恢复机制,我们需要通过断点续传+Redis状态跟踪+调度控制的方式来实现。以下是具体方案:
核心思路
- 拆分长任务:将原本10分钟的任务拆分为可中断的分段执行逻辑,任务内部定期检查暂停信号。
- 状态跟踪:用Redis存储每个任务的执行进度和暂停标记,确保暂停后能从断点恢复。
- 调度控制:实现一个调度任务,负责依次启动任务、定时暂停、最后统一恢复所有任务。
代码实现
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
相关产品推荐
相关产品推荐

