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

如何获取Celery Chord所有成员的task_id并跟踪子任务执行进度?

跟踪Celery Chord中子任务的执行进度

我来帮你解决这个问题——要跟踪Chord里所有子任务的进度,核心是要利用Chord底层依赖的Group结构,正确获取子任务组的ID,然后逐个获取每个子任务的状态和进度信息。咱们一步步来修改你的代码:

首先,先分析你现有代码里的几个问题:

  • 你在process_item里尝试更新chord主任务的状态,这没必要,反而会干扰主任务的状态跟踪
  • 你通过res1.result获取chord的result ID后,试图从主任务的info里拿group_id,但其实可以更直接地在创建chord时就拿到group的ID
  • group_res.results返回None是因为你没有正确初始化GroupResult,或者没有等待子任务开始执行

修改后的完整代码

import random
import time
from typing import List
import celery
from celery import Celery, chord, group
from celery.utils.log import get_task_logger

app = Celery( "chaintasks", backend="redis://localhost:6379", broker="pyamqp://guest@localhost//", )
app.conf.update(
    task_track_started=True,
    result_persistent=True,
    # 确保GroupResult能正确存储
    result_expires=3600,
)

@app.task(bind=True)
def process_item(self: celery.Task, item: str) -> str:
    t = random.randint(0, 10)
    for i in range(t + 1):  # +1是为了让进度能走到100%
        # 只更新当前子任务的进度状态
        self.update_state(
            state="PROGRESS",
            meta={"progress": round(i / t * 100, 2), "item": item}
        )
        time.sleep(1)
    return item + "_processed"

@app.task()
def group_items(items: List[str]) -> str:
    return ":".join(items)

@app.task
def process_and_group(items: List[str]):
    # 先创建子任务组
    item_tasks = group(process_item.s(i) for i in items)
    # 基于组创建chord
    chord_task = chord(item_tasks, group_items.s())
    # 执行chord任务
    chord_result = chord_task.delay()
    # 返回chord的结果ID和子任务组的ID,方便后续跟踪
    return {
        "chord_result_id": chord_result.id,
        "group_id": item_tasks.id  # 直接从group对象获取ID,最可靠
    }

random.seed(42)

if __name__ == "__main__":
    all_items = ["item1", "item2", "item3", "item4"]
    # 启动任务并获取chord和group的ID
    init_result = process_and_group.delay(all_items)
    
    # 等待初始化完成,拿到关键ID
    while init_result.state != "SUCCESS":
        time.sleep(0.5)
    task_ids = init_result.result
    chord_result_id = task_ids["chord_result_id"]
    group_id = task_ids["group_id"]
    
    # 初始化子任务组和最终chord任务的结果对象
    group_result = app.GroupResult(group_id)
    chord_result = app.AsyncResult(chord_result_id)
    
    # 持续跟踪进度直到所有任务完成
    while chord_result.state != "SUCCESS":
        print(f"\n当前Chord整体状态: {chord_result.state}")
        print("子任务进度详情:")
        
        total_progress = 0
        completed_tasks = 0
        # 遍历每个子任务的AsyncResult,获取实时状态
        for subtask in group_result:
            if subtask.state == "SUCCESS":
                completed_tasks += 1
                total_progress += 100
                print(f"- {subtask.info.split('_')[0]}: 已完成")
            elif subtask.state == "PROGRESS":
                progress = subtask.info["progress"]
                total_progress += progress
                print(f"- {subtask.info['item']}: {progress}%")
            else:
                item_name = subtask.info["item"] if isinstance(subtask.info, dict) else "未知任务"
                print(f"- {item_name}: {subtask.state}")
        
        # 计算整体平均进度
        avg_progress = round(total_progress / len(group_result), 2) if group_result else 0
        print(f"整体平均进度: {avg_progress}% ({completed_tasks}/{len(group_result)} 任务完成)")
        
        time.sleep(1)
    
    print(f"\n所有任务完成!最终汇总结果: {chord_result.result}")

关键改动说明

  1. 直接获取Group ID:
    在process_and_group里,我们先显式创建group对象,然后直接通过item_tasks.id拿到子任务组的ID,不需要通过主任务的info传递,这是最可靠的方式,避免了状态传递的不确定性。

  2. 子任务只维护自身进度:
    去掉了process_item里更新chord主任务状态的代码,每个子任务只更新自己的PROGRESS状态,携带进度百分比和当前处理的item,这样我们可以直接从子任务的info里拿到进度信息,逻辑更清晰。

  3. 正确使用GroupResult:
    通过app.GroupResult(group_id)初始化子任务组的结果对象后,我们可以直接遍历它来获取每个子任务的AsyncResult,进而拿到每个任务的实时状态和进度数据。

  4. 直观的进度统计:
    遍历所有子任务,累加每个任务的进度值计算平均进度,同时统计已完成的任务数量,这样就能直观展示整体的执行进度。

关于self.request.group的疑问

你代码里的self.request.group其实就是当前子任务所属的Group的ID,和我们通过item_tasks.id拿到的是同一个值。你可以在process_item里添加打印来验证:

print(f"当前子任务所属Group ID: {self.request.group}")

用这个ID初始化GroupResult就能获取所有子任务的信息,之前你拿到group_res.results为None,是因为当时子任务还没开始执行,或者你没有等待任务初始化完成就去获取结果。

这样修改后,你就能完美跟踪Chord里每个子任务的执行进度,以及整体的平均进度了!

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 19:12:55