Airflow provider для ИИ — это пакет apache-airflow-providers-common-ai, который добавляет в DAG операторы и декораторы для вызова модели как обычной задачи. Подключение сводится к четырём решениям: версия пакета, соединение с ключом, формат ответа и проверка на эталонном наборе. Подход подходит, когда ИИ-шаг в конвейере один, вход проверяем, а результат можно сравнить с ручной разметкой.

Пакет и версия

TL;DR

Официальный провайдер называется apache-airflow-providers-common-ai и рассчитан на Airflow версии 3.0.0 и новее. Он превращает вызов модели в обычную задачу DAG: оператор LLMOperator или декоратор @task.llm.

Установка идёт через менеджер пакетов вашего окружения. У пакета есть группы необязательных зависимостей (extras) под поставщиков моделей: openai, anthropic, google, bedrock, а также mcp, sql, pdf и другие. Ставьте группу только выбранного поставщика, остальные приносят в окружение Airflow лишние библиотеки. Версию пакета закрепите в файле зависимостей, а обновление делайте отдельным шагом с повторным прогоном эталонного набора: интерфейс операторов вправе меняться между выпусками.

  • Версия Airflow на вашем сервере: на установке старее 3.0.0 сначала обновляется сам Airflow.
  • Поставщик модели: сервис вендора напрямую или открытые веса на вашем сервере.
  • Список данных, которые разрешено отправлять за пределы компании: его согласует владелец данных, а автор DAG исполняет.
  • Шаг конвейера, который станет вызовом модели: один, с понятным входом и выходом.

Для первого проекта берут самый узкий вариант, оператор одного запроса LLMOperator. Остальные операторы пакета решают соседние задачи: LLMBatchOperator обрабатывает пакеты, LLMFileAnalysisOperator разбирает файлы и изображения, LLMBranchOperator выбирает ветку DAG, AgentOperator ведёт многошаговую работу с инструментами. Агентный оператор расширяет поверхность риска, поэтому его подключают отдельным решением после того, как простой шаг доказал свою пригодность.

Эта статья про подключение и первый проверенный шаг. Запуск DAG по событию описан в материале про Airflow trigger, ожидание входного файла — в материале про Airflow sensor.

Соединение и ключ

Провайдер берёт настройки из соединения (connection) Airflow, поэтому ключ остаётся за пределами кода DAG. Тип соединения называется pydanticai, идентификатор по умолчанию — pydanticai_default. В поле модели указывают строку вида «поставщик:идентификатор-модели», а ключ API кладут в поле пароля. Для поставщиков, которые авторизуются через окружение, например Bedrock или Vertex AI, поле пароля остаётся пустым. Актуальные идентификаторы моделей берите на странице вендора.

  1. Установите пакет с группой выбранного поставщика в то же окружение, где работают планировщик и исполнители задач.
  2. Создайте соединение типа pydanticai под отдельным идентификатором для этого конвейера: общий pydanticai_default на все DAG сразу размывает ответственность.
  3. Запишите модель в формате «поставщик:идентификатор-модели», ключ положите в поле пароля.
  4. Храните соединение в secrets backend или в переменной окружения, а в веб-интерфейсе оставьте только идентификатор. Просмотр соединений выдавайте узкому кругу администраторов.
  5. Для собственного сервера с открытыми весами укажите адрес в поле host: в документации приведён пример http://localhost:11434/v1 для Ollama.

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

Один шаг с моделью

Декоратор @task.llm превращает функцию в задачу: значение, которое она возвращает, становится промптом, а параметр llm_conn_id указывает соединение. Параметр output_type задаёт ожидаемый тип ответа. Если передать класс Pydantic BaseModel, ответ придёт структурированным, и следующая задача получит готовый объект вместо строки для разбора. Параметр system_prompt поддерживает шаблоны Jinja.

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

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

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

Эталонный набор

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

Итог прогонаЧто это значитДействие
Категория совпала с эталономШаг справился на этом примереОставить пример в наборе
Категория отличается, основание разумноеСпорная позиция или двусмысленный переченьБухгалтер решает, правим перечень или инструкцию
Ответ вне перечняНарушен форматПоправить инструкцию и тип ответа, повторить прогон
В основании названы факты, отсутствующие во входеВыдумка моделиШаг в рабочий поток допускать нельзя

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

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

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

Какие повторяющиеся решения в вашем отделе можно сверить с эталоном?

Прийти на Discovery →

Повтор и отказ

Сбои сети и ограничения поставщика Airflow повторяет через обычные параметры задачи, retries и retry_delay. Классификация ошибок и защита от дублей разобраны в материале про Airflow retry для LLM-шага, а уведомление владельца о сбое — в материале про callback Airflow. Для подключения важнее правило отказа: когда попытки исчерпаны, категория остаётся пустой, а позиция уходит бухгалтеру на ручной разбор. Подставлять категорию по умолчанию нельзя, иначе ошибка превращается в запись в учёте.

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

Тестовое и рабочее окружения получают разные соединения с разными ключами и разными лимитами. Ключ из теста нельзя использовать в рабочем потоке, даже временно: так проще отозвать доступ одной из сторон и проще понять по журналу поставщика, чьи это были вызовы. Раз в квартал ключи перевыпускают, а старые отзывают.

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

Первым делом закрепите версию пакета и создайте отдельное соединение под один конвейер. Вторым идёт разметка эталонного набора, и лишь третьим — подключение шага к рабочему потоку.

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

Как установить Airflow provider для ИИ?
Пакет называется apache-airflow-providers-common-ai, нужен Airflow версии 3.0.0 или новее. Ставьте его с группой зависимостей выбранного поставщика и закрепляйте версию в файле зависимостей.
Где хранить ключ API в Airflow?
В соединении типа pydanticai, в поле пароля. Само соединение держите в secrets backend или переменной окружения. В коде DAG и в журналах ключ исключён.
Какие операторы даёт провайдер?
Среди них LLMOperator для одного запроса, LLMBatchOperator для пакетной обработки, LLMBranchOperator для ветвления и AgentOperator для многошаговых задач. Для первого шага хватает LLMOperator.
Можно ли подключить модель на своём сервере?
Да. В поле host соединения указывают адрес совместимого endpoint, в документации приведён пример для Ollama. Модель на вашем сервере проходит тот же тест на эталонном наборе.
Что делать при ошибке вызова модели?
Настройте повторы через параметры задачи. Когда попытки исчерпаны, передайте запись человеку на ручной разбор и оставьте результат пустым.