Dùng Sensor — operator chỉ làm một việc: kiểm tra định kỳ một điều kiện, đúng thì chuyển success để task sau chạy, quá timeout thì fail.
Sensor hay gặp:
- S3KeySensor, GCSObjectExistenceSensor — chờ file đối tác đẩy lên bucket.
- ExternalTaskSensor — chờ một task hoặc DAG khác chạy xong (DAG báo cáo chờ DAG ingest).
- SqlSensor — chờ một query trả về kết quả, ví dụ partition hôm nay đã có dữ liệu.
- DateTimeSensor, HttpSensor, FileSensor.
from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor
wait_file = S3KeySensor(
task_id="wait_partner_file",
bucket_key="s3://partner/sales/{{ ds }}/_SUCCESS",
poke_interval=300, # check every 5 minutes
timeout=6 * 60 * 60, # give up after 6 hours
mode="reschedule",
)
wait_file >> load_salesChế độ chạy — phần hay bị hỏi tiếp:
- poke (mặc định): chiếm một worker slot suốt thời gian chờ. Hợp khi kiểm tra dày, cỡ vài giây, và chờ ngắn.
- reschedule: nhả slot giữa các lần kiểm tra. Hợp khi chờ lâu, kiểm tra thưa.
- Deferrable (deferrable=True với sensor có hỗ trợ): giao việc chờ cho triggerer chạy bất đồng bộ, gần như không tốn worker. Nên ưu tiên khi phải chờ hàng giờ.
Nếu cả hai DAG đều nằm trong Airflow, lập lịch theo Asset (tên cũ Dataset ở Airflow 2) — DAG sau tự chạy khi DAG trước cập nhật dữ liệu — thường gọn hơn dùng ExternalTaskSensor.
Lưu ý: hàng chục sensor ở mode poke chờ file cả buổi sáng có thể chiếm hết slot, khiến task thật không có chỗ chạy dù cluster còn rảnh CPU. Luôn đặt timeout rõ ràng — mặc định của sensor là 7 ngày.