Airflow sensor ставят перед ИИ-шагом, когда модель нельзя запускать без свежего входа: сенсор ждёт файл или подтверждённое событие, и только после его успеха стартуют следующие задачи. Если данные приходят от подрядчика или другой системы без точного времени, такая задача заменяет ожидание «по часам» и защищает от разбора пустого или вчерашнего файла. Когда вход создаёт ваш же DAG, сенсор лишний.
Что ждёт сенсор
Сенсор превращает фразу «когда придут данные» в проверяемое условие с границей по времени. Модель получает вход только после успеха сенсора, а при тишине конвейер сообщает о проблеме вместо запуска на пустоте.
По документации Airflow, сенсор — это специальный оператор, который ждёт наступления события. Он либо завершается успехом и запускает зависимые задачи, либо падает по истечении времени ожидания и поднимает тревогу. Для ИИ-конвейера это точка контроля: до неё идёт приём данных, после неё работает модель.
Представьте подрядчика: его бухгалтерия выкладывает ежедневную выгрузку платежей в общую папку, а ваш DAG просит модель разметить подозрительные строки для проверяющего. Время выкладки гуляет: то утро, то обед. Фиксированный запуск в девять утра иногда читает вчерашний файл, иногда ничего. Сенсор снимает эту неопределённость и делает ожидание видимым в интерфейсе.
Нужно ли вообще ждать? Ответ зависит от цены ошибки: запуск модели на вчерашних данных даёт уверенный и неверный отчёт, а ожидание обходится в несколько минут. Заготовок несколько. FileSensor следит за файлом в файловой системе, BashSensor ждёт успешного выполнения команды, PythonSensor ждёт, пока ваша функция вернёт истину, ExternalTaskSensor ждёт завершения задачи в другом DAG. Для внешнего события из чужой системы чаще всего пригоден PythonSensor с проверкой через API или по записи в вашей базе.
Poke или reschedule
Режимов два. В poke, который включён изначально, сенсор занимает слот воркера на всё время ожидания и проверяет условие с заданным интервалом. В reschedule он занимает слот только в момент проверки, а между проверками освобождает его. Документация подсказывает ориентир: poke подходит для частых проверок, раз в секунды, а reschedule лучше при интервале от минуты и выше.
| Ситуация | Режим | Почему |
|---|---|---|
| Файл ожидается в течение ближайших минут, проверка каждые несколько секунд | poke | Частая проверка быстро замечает файл, слот занят ненадолго |
| Выгрузка подрядчика приходит в неопределённый час, проверка раз в несколько минут | reschedule | Слоты воркеров остаются свободными для других задач |
| Десятки DAG ждут разные источники одновременно | reschedule | Иначе ожидание займёт весь пул воркеров |
| Одна проверка обходится дорого, например платный вызов API | reschedule с увеличением интервала | Меньше обращений между попытками |
Для растущих пауз есть параметр exponential_backoff: интервалы между проверками увеличиваются до предела max_wait. Это уместно для событий, которые чем дольше задерживаются, тем реже их имеет смысл проверять. Базовый интервал задаёт poke_interval. Подбирайте его под реальное поведение источника: проверка каждые три секунды файла, который приходит раз в сутки, бессмысленно нагружает систему. Запишите выбранный интервал и причину рядом с кодом DAG, иначе через полгода его поменяют наугад.
Учитывайте и пул слотов. Если все воркеры заняты ожидающими сенсорами в режиме poke, полезные задачи встают в очередь, а команда ищет причину в самой модели. Перед запуском десятка однотипных DAG прикиньте, сколько сенсоров будут ждать одновременно, и выберите режим по этому числу и оставьте привычку в стороне.
Timeout и пропуск
Без ограничения сенсор может висеть сутками, и о проблеме узнают, когда руководитель спросит про отчёт. Параметр timeout задаёт максимальное время ожидания: по документации, оно отсчитывается от первой попытки, и по истечении сенсор падает. Выбирайте значение из реальной жизни процесса: если выгрузка обычно приходит до обеда, а ждать вы готовы до вечера, ставьте timeout на вечер: всё позже считается сбоем.
Что делать при тишине, решает бизнес-правило, разработчик его лишь исполняет. Параметр soft_fail помечает сенсор пропущенным вместо падения, и зависимые задачи тоже пропускаются. Подходит для необязательных отчётов, отсутствие которых допустимо. Для обязательных входов пропуск опасен, потому что зелёно-серый запуск легко принять за норму. Там падение с оповещением владельца источника уместнее.
Что у вас происходит, если данные от подрядчика опаздывают?
Опишите в регламенте, кто получает оповещение о таймауте и в какие часы. Различайте повторы упавшего шага и ожидание входа. Если вызов модели упал по временной причине, его повторяют, и это забота самого шага. Если сенсор дождался таймаута, повторять модель бессмысленно: ей нечего читать. Поэтому в сообщении об ошибке указывайте, чего именно ждали и у кого это спросить.
Проверка готовности
Появление файла ещё мало говорит о его готовности. Подрядчик может загружать его постепенно, и в момент проверки в папке лежит половина данных. Надёжных соглашений два, и оба простые: источник сначала пишет под временным именем и переименовывает после окончания, либо рядом кладёт маркерный файл о завершении. Сенсор ждёт именно маркер, а сам файл вторичен.
- Договоритесь с источником о признаке готовности: итоговое имя файла или отдельный маркер.
- Напишите проверку в PythonSensor: файл существует, имя соответствует дате, размер больше нуля.
- Добавьте проверку структуры: ожидаемые столбцы и число строк в допустимых рамках.
- Передайте дальше только путь к проверенному файлу, содержимое остаётся в хранилище.
Проверку структуры выполняйте кодом до обращения к модели. Скрипт убеждается, что столбцов нужное число, даты разумны, а пустых строк мало. Модель получает уже очищенный фрагмент. Тестовую выгрузку для проверки возьмите обезличенную: реальные платёжные данные в песочнице хранить незачем. Так полные данные обрабатывает программа, а языковая модель объясняет и предлагает, и экономия токенов выходит побочным эффектом. Подробнее о качестве входа при разборе данных рассказывает статья про анализ больших данных с ИИ.
Гейт перед моделью
Соберём цепочку целиком. Сенсор готовности входа, затем задача проверки структуры, затем задача вызова модели, затем публикация черновика. Каждое звено имеет своего владельца и свой сигнал тревоги, причём сигналы отправляются в канал, который читают в рабочие часы: опоздание источника уходит владельцу источника, сбой структуры — разработчику конвейера, спорный вывод модели — аналитику. Так сообщения остаются раздельными и читаются по адресу.
Решение по выводу модели остаётся за человеком, и в журнале запуска фиксируется, кто именно его принял. Если аналитик отклонил черновик, причина записывается: накопленные отказы показывают, где инструкцию пора править. Результат публикуется как черновик с пометкой источника и версии инструкции, аналитик сверяет ключевые цифры с исходными данными и подтверждает отправку. Права на запись в бизнес-системы проверяет сервер, а DAG служит лишь порядком шагов. Подробный пример аналитического контура с ответственностью сторон вы найдёте на странице про ИИ-аналитика данных под ключ.
Для проверки готовности устройте репетицию трёх сценариев на тестовой среде, прежде чем подключать боевой источник: данные пришли вовремя, данные опоздали, данные пришли неполными. Каждый сценарий должен заканчиваться понятным состоянием запуска и сообщением, которое поймёт дежурный без чтения кода. Пока хотя бы один исход приводит к тишине, гейт остаётся недоделанным.
Выберите один DAG, который сейчас стартует по часам, и замените запуск на сенсор с явным timeout. Решите заранее, падает он или пропускается при тишине, и пропишите получателя оповещения. Возможно, потребуется обсудить сроки с источником данных, зато конвейер перестанет работать вслепую. Если хотите собрать такой контур вместе с проверкой и согласованием, напишите нам на странице про внедрение ИИ в рабочие процессы.