Introduction à Apache Airflow en Python
Mike Metzger
Data Engineer




all_success – Toutes les tâches précédentes ont réussi, aucune omissionall_failed – Toutes les tâches précédentes ont échouéall_done – Toutes les tâches précédentes sont terminées, peu importe le résultatone_failed – Au moins une tâche a échouéone_success – Au moins une tâche a réussinone_failed – Toutes les tâches précédentes ont réussi ou ont été omisesfrom 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 – paramètre de tâche qui définit combien de fois Airflow doit relancer une tâche en échec avant de la marquer échouéeretry_delay – paramètre de tâche acceptant un timedelta définissant l'attente entre les relances@task(
trigger_rule=TriggerRule.ONE_SUCCESS,
retries=3,
retry_delay=timedelta(minutes=5)
)
from airflow.providers.standard.operators.trigger_dagrun
import TriggerDagRunOperatortrigger_dag_id – Doit correspondre à dag_idwait_for_completion – Bloque les tâches suivantes jusqu'à la finpoke_interval – Fréquence de vérificationconf – Données à transmettre au Dag enfanttrigger_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 # Exécuter la première tâche, puis déclencher le Dag enfant
trigger_child >> cleanup # Exécuter une tâche de nettoyage après la fin du Dag enfant
Introduction à Apache Airflow en Python