如何在Celery Chain中查询未被Worker拾取的待执行任务?
如何查询Celery任务链中未被Worker拾取的待执行任务
核心原因说明
Celery任务链(Chain)默认采用惰性触发机制:只有前一个任务执行完成后,才会将下一个任务的消息发送到Redis Broker中。因此,未执行的后续任务并不会提前出现在Broker队列或Worker的scheduled/reserved列表里,这就是你之前通过celery.control.inspect()和Redis查询不到的原因。
解决方案:自定义任务链追踪机制
要在chain变量作用域外查询待执行任务,需要在创建链时主动记录所有子任务的ID,并结合结果后端和Worker状态进行筛选。
步骤1:创建任务链时记录所有子任务ID
通过提前生成每个子任务的ID,并将链的任务ID列表存入Redis(或django-celery-results的扩展字段),实现全局追踪:
from celery import chain import redis from your_celery_app import app # 导入你的Celery实例 # 初始化Redis连接 r = redis.Redis(host='localhost', port=6379, db=0) @app.task def example(): from time import sleep sleep(10) # 任务完成后,从链的追踪列表中移除当前任务ID task_id = example.request.id root_task_id = example.request.headers.get('root_task_id') if root_task_id: r.lrem(f"chain_tasks:{root_task_id}", 0, task_id) # 创建带追踪的任务链 def create_tracked_chain(): sub_tasks = [example.s() for _ in range(4)] chain_task_ids = [] prev_task = None root_task_id = app.gen_task_id() for idx, task in enumerate(sub_tasks): task_id = app.gen_task_id() chain_task_ids.append(task_id) # 为任务添加根任务ID标识,便于后续更新追踪列表 task = task.set(task_id=task_id, headers={'root_task_id': root_task_id}) if prev_task: prev_task = prev_task.set(link=task) else: root_task = task.set(task_id=root_task_id) chain_task_ids[0] = root_task_id prev_task = task # 执行链并将任务ID列表存入Redis root_task.apply_async() r.rpush(f"chain_tasks:{root_task_id}", *chain_task_ids) return root_task_id # 创建链并获取根任务ID root_task_id = create_tracked_chain()
步骤2:查询待执行任务
通过Redis获取链的所有任务ID,结合django-celery-results的已完成任务和Worker的活跃任务,筛选出待执行任务:
from django_celery_results.models import TaskResult from your_celery_app import app from celery.result import AsyncResult def get_pending_chain_tasks(root_task_id): # 从Redis获取链的所有任务ID chain_key = f"chain_tasks:{root_task_id}" all_task_ids = [tid.decode() for tid in r.lrange(chain_key, 0, -1)] if not all_task_ids: return [] # 查询已完成的任务ID集合 completed_ids = set(TaskResult.objects.filter( task_id__in=all_task_ids, status='SUCCESS' ).values_list('task_id', flat=True)) # 查询Worker当前活跃的任务ID集合 inspect = app.control.inspect() active_ids = set() active_tasks = inspect.active() if active_tasks: for worker_tasks in active_tasks.values(): active_ids.update([t['id'] for t in worker_tasks]) # 计算待执行任务ID pending_ids = [tid for tid in all_task_ids if tid not in completed_ids and tid not in active_ids] # 返回待执行任务的AsyncResult列表 return [AsyncResult(tid) for tid in pending_ids] # 查询并打印待执行任务 pending_tasks = get_pending_chain_tasks(root_task_id) for task in pending_tasks: print(f"待执行任务ID: {task.id}, 当前状态: {task.status}")
补充说明
- 若不需要持久化追踪,也可以将链的任务ID列表存入django-celery-results的
TaskResult模型扩展字段中。 - Celery 5.x+版本可使用
Task.on_success钩子替代任务内部的Redis更新逻辑,代码结构更简洁。
内容的提问来源于stack exchange,提问作者Jonesus
相关产品推荐
相关产品推荐

