Python 中的 Apache Airflow 入门
Mike Metzger
Data Engineer




all_success:所有先前任务成功,且无跳过all_failed:所有先前任务失败all_done:所有先前任务完成,结果不限one_failed:至少一个任务失败one_success:至少一个任务成功none_failed:所有先前任务成功或被跳过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!')
retries:任务参数,定义任务失败后在标记为失败前的重试次数retry_delay:任务参数,接受 timedelta,定义重试间隔@task(
trigger_rule=TriggerRule.ONE_SUCCESS,
retries=3,
retry_delay=timedelta(minutes=5)
)
from airflow.providers.standard.operators.trigger_dagrun
import TriggerDagRunOperatortrigger_dag_id:必须与 dag_id 匹配wait_for_completion:在完成前阻塞后续任务poke_interval:检查频率conf:传递给子 Dag 的数据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 入门