Airflow Sensors

Apache Airflow เบื้องต้นด้วย Python

Mike Metzger

Data Engineer

Sensors

  • Operator ที่รอจนกว่าเงื่อนไขจะเป็นจริง
    • การสร้างไฟล์
    • การอัปโหลดระเบียนในฐานข้อมูล
    • การตอบสนองจาก Web Request
  • กำหนดความถี่ในการตรวจสอบเงื่อนไขได้
  • กำหนดให้กับ Task ได้

ภาพประกอบ Airflow Sensor ที่รอให้เงื่อนไขเป็นจริง

Apache Airflow เบื้องต้นด้วย Python

รายละเอียด Sensor

  • สืบทอดมาจาก airflow.sdk.BaseSensorOperator
  • อาร์กิวเมนต์ของ Sensor:
  • mode - วิธีตรวจสอบเงื่อนไข
    • mode='poke' - ค่าเริ่มต้น รันซ้ำต่อเนื่อง
    • mode='reschedule' - คืน Task Slot แล้วลองใหม่ภายหลัง
  • poke_interval - ระยะเวลารอระหว่างการตรวจสอบ
  • timeout - ระยะเวลารอก่อนที่ Task จะล้มเหลว
  • รองรับ Attribute ปกติของ Operator ด้วย
Apache Airflow เบื้องต้นด้วย Python

File Sensor

  • เป็นส่วนหนึ่งของไลบรารี airflow.providers.standard.sensors
  • ตรวจสอบการมีอยู่ของไฟล์ในตำแหน่งที่กำหนด
  • ตรวจสอบได้ด้วยว่ามีไฟล์ใดอยู่ในไดเรกทอรีหรือไม่
from airflow.providers.standard.sensors.filesystem import FileSensor

file_sensor_task = FileSensor(task_id='file_sense',
                              filepath='salesdata.csv',
                              poke_interval=300,
                              timeout=3000
                              )

init_sales_cleanup() >> file_sensor_task >> generate_report()
Apache Airflow เบื้องต้นด้วย Python

Sensor ประเภทอื่น

  • Sensor อื่น ๆ อีกมากใน airflow.providers.*.sensors
  • ExternalTaskSensor - รอให้ Task ใน DAG อื่นเสร็จสิ้น
  • HttpSensor - ร้องขอ URL และตรวจสอบเนื้อหา
  • SqlSensor - รัน SQL Query เพื่อตรวจสอบเนื้อหา
Apache Airflow เบื้องต้นด้วย Python

ควรใช้ Sensor เมื่อใด?

  • ไม่แน่ใจว่าเงื่อนไขจะเป็นจริงเมื่อใด
  • ต้องการหลีกเลี่ยงความล้มเหลวทันที
  • เพิ่มการรัน Task ซ้ำโดยไม่ใช้ Loop

 

 

ภาพประกอบ Airflow Sensor ที่รอให้เงื่อนไขเป็นจริง

Apache Airflow เบื้องต้นด้วย Python

มาฝึกกันเถอะ!

Apache Airflow เบื้องต้นด้วย Python

Preparing Video For Download...