Ответ
В Airflow сенсоры (Sensors) я использовал для ожидания выполнения внешних условий перед запуском DAG или задачи. Это позволяет строить отказоустойчивые и зависимые пайплайны.
Основные сценарии применения:
- Ожидание появления файла: Запуск обработки только после того, как в S3-бакете или на SFTP-сервере появится файл с данными за определённую дату.
- Зависимость от завершения внешнего процесса: Ожидание, пока другая система (например, Spark-джоба или сторонний API) завершит свою работу и выставит флаг (например, запишет метку в БД).
- Проверка доступности данных в БД: Ожидание, когда в целевой таблице появятся свежие записи, необходимые для следующего этапа трансформации.
Пример сенсора для ожидания файла в S3:
from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor
wait_for_file = S3KeySensor(
task_id='wait_for_daily_file',
bucket_name='my-data-bucket',
bucket_key='data/landing/{{ ds }}/input.csv',
aws_conn_id='aws_default',
mode='poke',
poke_interval=60, # Проверять каждые 60 секунд
timeout=60 * 30 # Таймаут 30 минут
)
Сенсор будет «пинговать» S3 раз в минуту, и задача process_data запустится только после того, как файл появится.