Введение в Apache Airflow на Python
Mike Metzger
Data Engineer




all_success — все предыдущие задачи выполнены успешно, без пропусковall_failed — все предыдущие задачи завершились с ошибкойall_done — все предыдущие задачи завершены, независимо от результатаone_failed — хотя бы одна задача завершилась с ошибкойone_success — хотя бы одна задача выполнена успешноnone_failed — все предыдущие задачи выполнены успешно или пропущеныfrom airflow.utils.trigger_rule import TriggerRule@tasktrigger_rule=TriggerRule.<TriggerRule enum>@task(trigger_rule=TriggerRule.ALL_SUCCESS)
def run_if_everything_succeeds:
print('All previous tasks succeeded!')
retries — параметр задачи, определяющий количество повторных попыток Airflow перед тем, как пометить задачу как неудавшуюся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_idwait_for_completion — блокирует дальнейшие задачи до завершенияpoke_interval — частота проверкиconf — данные, передаваемые дочернему DAGtrigger_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 # Run first task, then kick off child Dag
trigger_child >> cleanup # Run a cleanup task after child Dag completes
Введение в Apache Airflow на Python