如何获取Celery Chord所有成员的task_id并跟踪子任务执行进度?
我来帮你解决这个问题——要跟踪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}")
关键改动说明
直接获取Group ID:
在process_and_group里,我们先显式创建group对象,然后直接通过item_tasks.id拿到子任务组的ID,不需要通过主任务的info传递,这是最可靠的方式,避免了状态传递的不确定性。子任务只维护自身进度:
去掉了process_item里更新chord主任务状态的代码,每个子任务只更新自己的PROGRESS状态,携带进度百分比和当前处理的item,这样我们可以直接从子任务的info里拿到进度信息,逻辑更清晰。正确使用GroupResult:
通过app.GroupResult(group_id)初始化子任务组的结果对象后,我们可以直接遍历它来获取每个子任务的AsyncResult,进而拿到每个任务的实时状态和进度数据。直观的进度统计:
遍历所有子任务,累加每个任务的进度值计算平均进度,同时统计已完成的任务数量,这样就能直观展示整体的执行进度。
关于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

