Airflow中能否为不同上游任务单独配置触发规则?
Airflow 是否支持为单个上游任务单独设置触发规则
Airflow 原生不支持针对单个上游任务配置差异化的触发规则,所有内置触发规则(all_success/all_done/one_success等)都是对任务的全部上游依赖统一生效。
你描述的多上游差异化依赖需求,不需要修改框架源码,通过新增一个无业务逻辑的中间适配任务即可实现,完全匹配你提出的三个触发条件。
具体实现思路
核心逻辑是把「仅要求执行完成、不要求成功」的上游任务单独拆分链路,通过中间任务把它的终态(无论成功/失败)转换成下游任务可识别的成功信号,最终下游任务仍然使用默认的all_success触发规则即可:
- 要求必须成功的
check_fact_table_1、check_dim_table_3直接作为produce_bq_table的上游,只要其中任意一个失败,下游任务就不会触发 - 仅要求跑完所有重试的
check_fact_table_2,先接入一个配置了all_done触发规则的空任务,这个空任务没有实际业务逻辑,只要上游的check_fact_table_2跑完所有流程(成功/失败都算)就会标记自身为成功,再把这个空任务作为produce_bq_table的上游 - 最终
produce_bq_table用默认的all_success规则,只有三个上游(两个强制成功的校验任务、一个承接fact_table2状态的空任务)都成功时才会启动,完全匹配你的需求
可直接复用的代码示例
# 导入依赖,Airflow 1.x版本请将EmptyOperator替换为DummyOperator from airflow.providers.google.cloud.operators.bigquery import BigQueryOperator, BigQueryCheckOperator from airflow.operators.empty import EmptyOperator from airflow.models.baseoperator import chain # 原有三个校验任务保持原有定义即可,注意给check_fact_table_2配置符合预期的重试次数 check_fact_table_1 = BigQueryCheckOperator( task_id='check_fact_table_1', sql='-- 校验fact_table_1数据的SQL逻辑', dag=dag ) check_fact_table_2 = BigQueryCheckOperator( task_id='check_fact_table_2', retries=3, # 配置你需要的重试次数,任务会跑完所有重试才会进入终态 sql='-- 校验fact_table_2数据的SQL逻辑', dag=dag ) check_dim_table_3 = BigQueryCheckOperator( task_id='check_dim_table_3', sql='-- 校验dim_table_3数据的SQL逻辑', dag=dag ) # 新增中间适配空任务,仅做状态转换,无实际业务逻辑 fact2_check_finish_flag = EmptyOperator( task_id='fact2_check_finish_flag', trigger_rule='all_done', dag=dag ) # 原有生产BQ表的任务保持定义不变,使用默认all_success触发规则即可 produce_bq_table = BigQueryOperator( dag=dag, table_name='bq_table$20220601', task_id='produce_bq_table', sql='sql/produce_bq_table.sql' ) # 配置依赖链路 chain(check_fact_table_2, fact2_check_finish_flag) produce_bq_table.set_upstream([ check_fact_table_1, fact2_check_finish_flag, check_dim_table_3 ])
注意:如果后续有其他差异化依赖需求,比如某个上游任务允许被跳过、仅要求不失败即可,可以调整中间适配任务的触发规则,比如用
none_failed、none_failed_min_one_success等内置规则完成适配,不需要写自定义代码。
内容的提问来源于stack exchange,提问作者khusnanadia
相关产品推荐
相关产品推荐

