Airflow sensor ставят перед ИИ-шагом, когда модель нельзя запускать без свежего входа: сенсор ждёт файл или подтверждённое событие, и только после его успеха стартуют следующие задачи. Если данные приходят от подрядчика или другой системы без точного времени, такая задача заменяет ожидание «по часам» и защищает от разбора пустого или вчерашнего файла. Когда вход создаёт ваш же DAG, сенсор лишний.

Что ждёт сенсор

TL;DR

Сенсор превращает фразу «когда придут данные» в проверяемое условие с границей по времени. Модель получает вход только после успеха сенсора, а при тишине конвейер сообщает о проблеме вместо запуска на пустоте.

По документации Airflow, сенсор — это специальный оператор, который ждёт наступления события. Он либо завершается успехом и запускает зависимые задачи, либо падает по истечении времени ожидания и поднимает тревогу. Для ИИ-конвейера это точка контроля: до неё идёт приём данных, после неё работает модель.

Представьте подрядчика: его бухгалтерия выкладывает ежедневную выгрузку платежей в общую папку, а ваш DAG просит модель разметить подозрительные строки для проверяющего. Время выкладки гуляет: то утро, то обед. Фиксированный запуск в девять утра иногда читает вчерашний файл, иногда ничего. Сенсор снимает эту неопределённость и делает ожидание видимым в интерфейсе.

Нужно ли вообще ждать? Ответ зависит от цены ошибки: запуск модели на вчерашних данных даёт уверенный и неверный отчёт, а ожидание обходится в несколько минут. Заготовок несколько. FileSensor следит за файлом в файловой системе, BashSensor ждёт успешного выполнения команды, PythonSensor ждёт, пока ваша функция вернёт истину, ExternalTaskSensor ждёт завершения задачи в другом DAG. Для внешнего события из чужой системы чаще всего пригоден PythonSensor с проверкой через API или по записи в вашей базе.

Poke или reschedule

Режимов два. В poke, который включён изначально, сенсор занимает слот воркера на всё время ожидания и проверяет условие с заданным интервалом. В reschedule он занимает слот только в момент проверки, а между проверками освобождает его. Документация подсказывает ориентир: poke подходит для частых проверок, раз в секунды, а reschedule лучше при интервале от минуты и выше.

СитуацияРежимПочему
Файл ожидается в течение ближайших минут, проверка каждые несколько секундpokeЧастая проверка быстро замечает файл, слот занят ненадолго
Выгрузка подрядчика приходит в неопределённый час, проверка раз в несколько минутrescheduleСлоты воркеров остаются свободными для других задач
Десятки DAG ждут разные источники одновременноrescheduleИначе ожидание займёт весь пул воркеров
Одна проверка обходится дорого, например платный вызов APIreschedule с увеличением интервалаМеньше обращений между попытками

Для растущих пауз есть параметр exponential_backoff: интервалы между проверками увеличиваются до предела max_wait. Это уместно для событий, которые чем дольше задерживаются, тем реже их имеет смысл проверять. Базовый интервал задаёт poke_interval. Подбирайте его под реальное поведение источника: проверка каждые три секунды файла, который приходит раз в сутки, бессмысленно нагружает систему. Запишите выбранный интервал и причину рядом с кодом DAG, иначе через полгода его поменяют наугад.

Учитывайте и пул слотов. Если все воркеры заняты ожидающими сенсорами в режиме poke, полезные задачи встают в очередь, а команда ищет причину в самой модели. Перед запуском десятка однотипных DAG прикиньте, сколько сенсоров будут ждать одновременно, и выберите режим по этому числу и оставьте привычку в стороне.

Timeout и пропуск

Без ограничения сенсор может висеть сутками, и о проблеме узнают, когда руководитель спросит про отчёт. Параметр timeout задаёт максимальное время ожидания: по документации, оно отсчитывается от первой попытки, и по истечении сенсор падает. Выбирайте значение из реальной жизни процесса: если выгрузка обычно приходит до обеда, а ждать вы готовы до вечера, ставьте timeout на вечер: всё позже считается сбоем.

Что делать при тишине, решает бизнес-правило, разработчик его лишь исполняет. Параметр soft_fail помечает сенсор пропущенным вместо падения, и зависимые задачи тоже пропускаются. Подходит для необязательных отчётов, отсутствие которых допустимо. Для обязательных входов пропуск опасен, потому что зелёно-серый запуск легко принять за норму. Там падение с оповещением владельца источника уместнее.

● Discovery · 1 час · бесплатно

Что у вас происходит, если данные от подрядчика опаздывают?

Прийти на Discovery →

Опишите в регламенте, кто получает оповещение о таймауте и в какие часы. Различайте повторы упавшего шага и ожидание входа. Если вызов модели упал по временной причине, его повторяют, и это забота самого шага. Если сенсор дождался таймаута, повторять модель бессмысленно: ей нечего читать. Поэтому в сообщении об ошибке указывайте, чего именно ждали и у кого это спросить.

Проверка готовности

Появление файла ещё мало говорит о его готовности. Подрядчик может загружать его постепенно, и в момент проверки в папке лежит половина данных. Надёжных соглашений два, и оба простые: источник сначала пишет под временным именем и переименовывает после окончания, либо рядом кладёт маркерный файл о завершении. Сенсор ждёт именно маркер, а сам файл вторичен.

  1. Договоритесь с источником о признаке готовности: итоговое имя файла или отдельный маркер.
  2. Напишите проверку в PythonSensor: файл существует, имя соответствует дате, размер больше нуля.
  3. Добавьте проверку структуры: ожидаемые столбцы и число строк в допустимых рамках.
  4. Передайте дальше только путь к проверенному файлу, содержимое остаётся в хранилище.

Проверку структуры выполняйте кодом до обращения к модели. Скрипт убеждается, что столбцов нужное число, даты разумны, а пустых строк мало. Модель получает уже очищенный фрагмент. Тестовую выгрузку для проверки возьмите обезличенную: реальные платёжные данные в песочнице хранить незачем. Так полные данные обрабатывает программа, а языковая модель объясняет и предлагает, и экономия токенов выходит побочным эффектом. Подробнее о качестве входа при разборе данных рассказывает статья про анализ больших данных с ИИ.

Гейт перед моделью

Соберём цепочку целиком. Сенсор готовности входа, затем задача проверки структуры, затем задача вызова модели, затем публикация черновика. Каждое звено имеет своего владельца и свой сигнал тревоги, причём сигналы отправляются в канал, который читают в рабочие часы: опоздание источника уходит владельцу источника, сбой структуры — разработчику конвейера, спорный вывод модели — аналитику. Так сообщения остаются раздельными и читаются по адресу.

Решение по выводу модели остаётся за человеком, и в журнале запуска фиксируется, кто именно его принял. Если аналитик отклонил черновик, причина записывается: накопленные отказы показывают, где инструкцию пора править. Результат публикуется как черновик с пометкой источника и версии инструкции, аналитик сверяет ключевые цифры с исходными данными и подтверждает отправку. Права на запись в бизнес-системы проверяет сервер, а DAG служит лишь порядком шагов. Подробный пример аналитического контура с ответственностью сторон вы найдёте на странице про ИИ-аналитика данных под ключ.

Для проверки готовности устройте репетицию трёх сценариев на тестовой среде, прежде чем подключать боевой источник: данные пришли вовремя, данные опоздали, данные пришли неполными. Каждый сценарий должен заканчиваться понятным состоянием запуска и сообщением, которое поймёт дежурный без чтения кода. Пока хотя бы один исход приводит к тишине, гейт остаётся недоделанным.

// в первую очередь

Выберите один DAG, который сейчас стартует по часам, и замените запуск на сенсор с явным timeout. Решите заранее, падает он или пропускается при тишине, и пропишите получателя оповещения. Возможно, потребуется обсудить сроки с источником данных, зато конвейер перестанет работать вслепую. Если хотите собрать такой контур вместе с проверкой и согласованием, напишите нам на странице про внедрение ИИ в рабочие процессы.

Частые вопросы

Что такое Sensor в Airflow?
Это специальный оператор, который ждёт события: появления файла, завершения задачи в другом DAG или выполнения вашего условия. Он завершается успехом, и тогда следующие задачи стартуют, либо падает по истечении timeout.
Чем отличаются режимы poke и reschedule?
В poke сенсор занимает слот воркера всё время ожидания, в reschedule слот занят только во время проверки. Poke подходит для частых проверок, reschedule — для интервалов от минуты и дольше.
Как задать время ожидания сенсора?
Параметром timeout, который отсчитывается от первой попытки. Интервал между проверками задаётся параметром poke_interval. Значения подбирайте под реальное поведение источника данных.
Что делает soft_fail?
При истечении ожидания сенсор получает статус SKIPPED вместо FAILED, и зависимые задачи тоже пропускаются. Параметр подходит для необязательных входов, но для обязательных лучше падение с оповещением.
Как дождаться полностью загруженного файла?
Договоритесь с источником о признаке готовности: итоговое имя после переименования или маркерный файл. Затем сенсор проверяет именно этот признак, а скрипт отдельно проверяет структуру содержимого.