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

Форма DAG

TL;DR

DAG с LLM-оператором состоит из четырёх задач: выгрузка, подготовка входа, вызов модели с подтверждением, публикация. Модель занимает одну задачу, остальные три выполняет обычный код.

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

ЗадачаОператорЧто делаетЧто отдаёт дальше
ВыгрузкаОбычный оператор или Python-задачаЧитает поставки за интервал данныхФайл с полной таблицей
Подготовка входаPython-задачаСчитает агрегаты, отбирает просрочки, убирает лишние поляКороткий JSON
СводкаLLMOperatorФормулирует текст по готовому JSONОбъект с полями сводки
ПубликацияPython-задачаЗаписывает утверждённую сводку в общую папкуФайл за этот интервал

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

Результат оператора Airflow кладёт в XCom. Для небольших структур это удобно, для крупных файлов между задачами передают путь, а содержимое остаётся в общем хранилище. Так история запусков остаётся лёгкой, а чувствительные тексты остаются вне служебной базы планировщика.

Вход модели

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

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

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

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

Повторный запуск

Airflow позволяет перезапустить DAG за прошлый интервал данных, и сделать это придётся: поправили выгрузку, нашли ошибку, сменили инструкцию. Идемпотентность здесь означает, что повторный запуск за тот же интервал заменяет черновик, и второго экземпляра на выходе нет. Для этого имя файла черновика строят из интервала данных, а новый запуск перезаписывает черновик с этим именем. Защита внешнего вызова от дублей на уровне одной попытки разобрана в материале про Airflow retry, а здесь речь о перезапуске всей цепочки.

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

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

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

Какой отчёт в вашей компании собирают вручную каждую неделю?

Прийти на Discovery →

Проверка выхода

У LLMOperator есть режим подтверждения: параметр require_approval приостанавливает выполнение до решения человека, allow_modifications разрешает рецензенту править текст, approval_timeout задаёт срок ожидания, а on_approval_timeout принимает значения fail, approve или reject. По документации, механизм рассчитан на Airflow 3.1 и новее. Для истечения срока выбирайте fail или reject: автоматическое approve превращает проверку в формальность.

Параметр output_type со структурой на базе Pydantic делает ответ объектом с полями, и рецензент видит сводку, список упомянутых поставок и отдельное поле «сомнения модели». Что именно рецензент сверяет, лучше записать шагами и повесить рядом с задачей подтверждения. Список короткий намеренно: длинный чек-лист читают по диагонали, и подпись под ним превращается в привычку. Каждый пункт сформулирован так, чтобы на него можно было ответить «совпало» или «расхождение» за один взгляд на экран.

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

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

Ошибка и разбор

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

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

// с чего начать

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

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

Чем LLMOperator отличается от AgentOperator?
LLMOperator отправляет один запрос и возвращает текст или структуру. AgentOperator ведёт многошаговую работу с инструментами. Для сводки по готовым данным берут первый, агентный подключают отдельным решением.
Как сделать DAG с LLM идемпотентным?
Стройте имя черновика из интервала данных и перезаписывайте его при повторном запуске. Утверждённую версию фиксируйте отдельно, замену делает человек.
Как подтвердить ответ модели до публикации?
Включите режим подтверждения у оператора, а рецензенту дайте список сверки: числа, идентификаторы, причины, тон. При истечении срока ожидания выбирайте отказ.
Что делать, если оператор вернул ошибку?
DAG останавливается до публикации, потребитель получает прошлую утверждённую версию с пометкой о сбое, причину разбирает ответственный по журналу запуска.
Нужен ли отдельный сервер для DAG с моделью?
Нет, DAG запускается там же, где работает Airflow. Отдельный сервер понадобится лишь для открытых весов, если вы решили держать модель у себя.