Apache Airflow เบื้องต้นด้วย Python
Mike Metzger
Data Engineer




all_success - ทุก task ก่อนหน้าสำเร็จ ไม่มีการข้ามall_failed - ทุก task ก่อนหน้าล้มเหลวall_done - ทุก task ก่อนหน้าเสร็จสิ้น ไม่ว่าผลจะเป็นอย่างไรone_failed - มีอย่างน้อยหนึ่ง task ที่ล้มเหลวone_success - มีอย่างน้อยหนึ่ง task ที่สำเร็จnone_failed - ทุก task ก่อนหน้าสำเร็จหรือถูกข้าม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 - พารามิเตอร์ของ task ที่กำหนดจำนวนครั้งที่ Airflow จะลองรัน task ที่ล้มเหลวใหม่ก่อนจะมาร์กว่าล้มเหลวretry_delay - พารามิเตอร์ของ task ที่รับค่า 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 - บล็อก task ถัดไปจนกว่าจะเสร็จ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 # Run first task, then kick off child Dag
trigger_child >> cleanup # Run a cleanup task after child Dag completes
Apache Airflow เบื้องต้นด้วย Python