触发器与容错

Python 中的 Apache Airflow 入门

Mike Metzger

Data Engineer

构建健壮的流水线

  • 触发器:控制任务何时运行
    • 定义在特定条件下执行的任务
    • 默认触发器:全部成功

说明触发器如何控制任务何时运行的插图

  • 容错:在任务内恢复
    • 重试单个任务

说明容错如何重试失败任务的插图

Python 中的 Apache Airflow 入门

触发规则

  • 在任务开始前检查上游任务状态
  • 让您定义任务是否应继续
  • 默认:所有上游任务成功后再继续
  • 触发规则可修改此行为

展示触发规则检查上游任务状态的插图

Python 中的 Apache Airflow 入门

触发规则用法

  • 通知任务
  • 清理任务
  • 条件执行

展示通知与清理等触发规则用例的插图

Python 中的 Apache Airflow 入门

关键触发规则

  • all_success:所有先前任务成功,且无跳过
  • all_failed:所有先前任务失败
  • all_done:所有先前任务完成,结果不限
  • one_failed:至少一个任务失败
  • one_success:至少一个任务成功
  • none_failed:所有先前任务成功或被跳过
Python 中的 Apache Airflow 入门

触发规则实现

  • from airflow.utils.trigger_rule import TriggerRule
  • 触发规则应用于 @task 装饰器
  • trigger_rule=TriggerRule.<TriggerRule enum>
@task(trigger_rule=TriggerRule.ALL_SUCCESS)
def run_if_everything_succeeds:
  print('All previous tasks succeeded!')
Python 中的 Apache Airflow 入门

容错属性

  • retries:任务参数,定义任务失败后在标记为失败前的重试次数
  • retry_delay:任务参数,接受 timedelta,定义重试间隔
@task(
  trigger_rule=TriggerRule.ONE_SUCCESS,
  retries=3,
  retry_delay=timedelta(minutes=5)
)
Python 中的 Apache Airflow 入门

TriggerDagRunOperator

  • 允许一个 Dag 触发另一个 Dag
  • 提供代码与 Dags
  • from airflow.providers.standard.operators.trigger_dagrun import TriggerDagRunOperator
  • trigger_dag_id:必须与 dag_id 匹配
  • wait_for_completion:在完成前阻塞后续任务
  • poke_interval:检查频率
  • conf:传递给子 Dag 的数据
Python 中的 Apache Airflow 入门

TriggerDagRun 示例

trigger_child = TriggerDagRunOperator(
        task_id="trigger_child_pipeline",
        trigger_dag_id="child_pipeline_dag",   
        wait_for_completion=True,              
        poke_interval=30,                      
        conf={                                 
            "source": "s3://my-bucket/raw/",
        },
    )

task1 >> trigger_child    # 先运行第一个任务,然后触发子 Dag
trigger_child >> cleanup  # 子 Dag 完成后运行清理任务
Python 中的 Apache Airflow 入门

Passons à la pratique !

Python 中的 Apache Airflow 入门

Preparing Video For Download...