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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 16:43:16