Python RQ多任务依赖实现:多任务完成后触发WebSocket通知任务
实现RQ嵌套任务流:多依赖任务+最终WebSocket通知
看起来你已经搭好了基本的任务依赖框架,我来帮你完善这个流程,实现先执行拷贝任务→并行执行三个注册任务→全部完成后发WebSocket通知的完整逻辑。
完整代码实现
首先假设你的任务函数(raw_nifti_copymachine、modality_registrator、send_websocket_notification)已经定义好,下面是完整的任务链构建代码:
from rq import Queue from redis import Redis # 替换成你实际的任务模块路径 from your_tasks_module import raw_nifti_copymachine, modality_registrator, send_websocket_notification # 初始化Redis连接和RQ队列 redis_conn = Redis(host="localhost", port=6379, db=0) redisqueue = Queue(connection=redis_conn) # -------------------------- # 1. 启动核心拷贝任务 # -------------------------- rncopy = redisqueue.enqueue( raw_nifti_copymachine, patientid, imagepath, timeout=6000 # 根据实际耗时调整超时时间 ) # -------------------------- # 2. 启动三个并行的注册任务(依赖拷贝任务完成) # -------------------------- t1c_reg = redisqueue.enqueue( modality_registrator, patientid, "t1c", timeout=6000, depends_on=rncopy # 只有rncopy完成后才会执行 ) t2_reg = redisqueue.enqueue( modality_registrator, patientid, "t2", timeout=6000, depends_on=rncopy ) # 假设第三个模态是flair,你可以替换成自己的实际模态类型 flair_reg = redisqueue.enqueue( modality_registrator, patientid, "flair", timeout=6000, depends_on=rncopy ) # -------------------------- # 3. 启动WebSocket通知任务(依赖三个注册任务全部完成) # -------------------------- notify_task = redisqueue.enqueue( send_websocket_notification, patientid, message="所有影像注册任务已完成!", timeout=30, # 新版本RQ支持传入任务列表,所有任务完成后才触发 depends_on=[t1c_reg, t2_reg, flair_reg] )
关键细节说明
多任务依赖的版本要求
上面的代码依赖RQ 1.10.0及以上版本,这个版本开始支持depends_on传入任务列表(所有任务完成后才执行当前任务)。如果你的RQ版本较低,先执行升级:pip install --upgrade rqWebSocket通知任务的示例实现
如果你还没写send_websocket_notification,这里给一个基于Flask-SocketIO的示例(适配多进程场景):from flask_socketio import SocketIO # 使用Redis作为消息队列,实现RQ worker和Web服务的跨进程通信 socketio = SocketIO(message_queue="redis://localhost:6379/0") def send_websocket_notification(patientid, message): # 发送消息到指定患者的WebSocket房间 socketio.emit( "task_completed", {"patient_id": patientid, "content": message}, room=f"patient_{patientid}" )错误处理建议
- 如果某个任务失败,依赖它的后续任务会被标记为
failed,你可以通过rq-dashboard监控任务状态 - 可以给任务添加
failure_callback参数,自定义失败后的处理逻辑(比如发送错误通知):def handle_task_failure(job, exc_type, exc_value, traceback): # 这里可以写失败后的通知逻辑 print(f"任务 {job.id} 失败: {exc_value}") # 在enqueue时添加回调 rncopy = redisqueue.enqueue( raw_nifti_copymachine, patientid, imagepath, timeout=6000, failure_callback=handle_task_failure )
- 如果某个任务失败,依赖它的后续任务会被标记为
旧版本RQ的兼容方案(如果无法升级)
如果你的RQ版本不支持多任务依赖,可以用一个中间任务来轮询等待所有注册任务完成:
import time from rq import Job from redis import Redis def wait_for_all_registrations(patientid, job_ids): redis_conn = Redis(host="localhost", port=6379, db=0) # 轮询检查所有任务是否完成 while True: all_finished = all( Job.fetch(job_id, connection=redis_conn).is_finished for job_id in job_ids ) if all_finished: send_websocket_notification(patientid, "所有注册任务已完成!") break time.sleep(5) # 每隔5秒检查一次 # 启动中间任务,依赖拷贝任务完成 wait_task = redisqueue.enqueue( wait_for_all_registrations, patientid, [t1c_reg.id, t2_reg.id, flair_reg.id], depends_on=rncopy )
这种方式虽然可行,但轮询会消耗一定资源,优先推荐升级RQ使用原生的多任务依赖功能。
内容的提问来源于stack exchange,提问作者florian
相关产品推荐
相关产品推荐

