Celery Chord 动态任务编排与回调
利用 Celery Chord 实现动态任务依赖图,支持运行时添加子任务并等待所有完成后再执行回调,适用于数据管道。 · 难度:入门 · +10XP
动态任务编排
传统 Celery chain/chord 的任务列表在定义时固定。本教程展示如何利用 group 与 chord 的动态构造能力:根据用户输入生成 N 个子任务,所有子任务完成后触发一个聚合回调。同时处理任务失败重试与超时场景。你将学习使用 AsyncResult 实时监控任务进度,并通过 Redis 或数据库持久化任务状态。
from celery import chord, group
@shared_task
def process_chunk(data_chunk):
# 处理数据块
return transformed
@shared_task
def aggregate_results(results):
# 合并所有结果
return sum(results)
def run_pipeline(data_parts):
tasks = [process_chunk.s(part) for part in data_parts]
callback = aggregate_results.s()
workflow = chord(tasks, callback)
result = workflow.delay()
return result.id