如何使用Airflow TaskFlow API实现多任务共享子任务并配置依赖
解决方案
问题根因
你当前的代码存在两个核心问题导致依赖不符合预期:
last_task定义在data_name的循环内部,会被重复覆盖,且每轮循环无论data_name是否为Data1,都执行了last_task(success),导致所有taskb的输出都被传入last_task,自然所有任务都和last_task建立了依赖关系- Airflow要求同一DAG下task_id全局唯一,你在循环内重复定义task_id为
last_task的任务会导致任务识别异常
修改后代码
def etl(): # 把last_task定义移到循环外层,保证全局唯一 @task(task_id='last_task') def last_task(success): dim_experiments.main() return for item in ['FIRST','SECCOND','THIRD']: # 这里原代码循环的item值和判断条件不匹配,你可以根据实际业务调整 if item == 'FIRST': requests = ['Data1','Data3'] else: requests = ['Data1'] for data_name in requests: @task(task_id=f'{item}_{data_name}_task_a') def taska(): a,b = some_func() vars_dict = {'a': a, 'b': b} return vars_dict # 这里原代码用了未定义的account变量,替换为循环变量item @task(task_id=f'{item}_{data_name}_get_liveops_data') def taskb(vars_dict): some_other_func() return True vars_dict = taska() success = taskb(vars_dict) # 仅在data_name为Data1时调用last_task,传入当前分支的success if data_name=='Data1': last_task(success) myc_dag = etl()
修改说明
last_task移到循环外层后不会被重复定义,task_id保持全局唯一- 仅在
data_name == 'Data1'的分支内调用last_task,仅传入当前Data1分支下taskb的返回值,因此last_task只会和Data1对应的taska、taskb建立依赖,不会关联Data3的任务 - 顺带修正了原代码的两处笔误:未定义的
account变量、循环item值和判断条件不匹配的问题,你可以根据实际业务逻辑调整
效果对比
当前错误DAG结构:
修改后期望DAG结构:
内容的提问来源于stack exchange,提问作者mrc
相关产品推荐
相关产品推荐

