Airflow provider для ИИ — это пакет apache-airflow-providers-common-ai, который добавляет в DAG операторы и декораторы для вызова модели как обычной задачи. Подключение сводится к четырём решениям: версия пакета, соединение с ключом, формат ответа и проверка на эталонном наборе. Подход подходит, когда ИИ-шаг в конвейере один, вход проверяем, а результат можно сравнить с ручной разметкой.
Пакет и версия
Официальный провайдер называется 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, поле пароля остаётся пустым. Актуальные идентификаторы моделей берите на странице вендора.
- Установите пакет с группой выбранного поставщика в то же окружение, где работают планировщик и исполнители задач.
- Создайте соединение типа
pydanticaiпод отдельным идентификатором для этого конвейера: общийpydanticai_defaultна все DAG сразу размывает ответственность. - Запишите модель в формате «поставщик:идентификатор-модели», ключ положите в поле пароля.
- Храните соединение в secrets backend или в переменной окружения, а в веб-интерфейсе оставьте только идентификатор. Просмотр соединений выдавайте узкому кругу администраторов.
- Для собственного сервера с открытыми весами укажите адрес в поле 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, а тестовая задача прогоняет шаг на всём наборе и сравнивает ответ модели с ручной разметкой. Допустимую долю совпадений владелец процесса задаёт до прогона: порог, подобранный после результатов, теряет смысл.
| Итог прогона | Что это значит | Действие |
|---|---|---|
| Категория совпала с эталоном | Шаг справился на этом примере | Оставить пример в наборе |
| Категория отличается, основание разумное | Спорная позиция или двусмысленный перечень | Бухгалтер решает, правим перечень или инструкцию |
| Ответ вне перечня | Нарушен формат | Поправить инструкцию и тип ответа, повторить прогон |
| В основании названы факты, отсутствующие во входе | Выдумка модели | Шаг в рабочий поток допускать нельзя |
Собирать набор начинают с реальных примеров прошлого квартала, взятых по разным поставщикам и по разным категориям, включая редкие. Перекос в сторону лёгких позиций даёт красивую цифру совпадений и пустое доверие. Размечает набор тот, кто отвечает за учёт, а второй сотрудник выборочно перепроверяет разметку.
Набор пополняют каждым спорным случаем из рабочей жизни. Через несколько месяцев он превращается в живое описание того, как компания понимает свои категории затрат, и это ценнее самого промпта. После обновления пакета или смены модели весь набор прогоняют заново, а расхождения с прошлым результатом читает человек.
Какие повторяющиеся решения в вашем отделе можно сверить с эталоном?
Повтор и отказ
Сбои сети и ограничения поставщика Airflow повторяет через обычные параметры задачи, retries и retry_delay. Классификация ошибок и защита от дублей разобраны в материале про Airflow retry для LLM-шага, а уведомление владельца о сбое — в материале про callback Airflow. Для подключения важнее правило отказа: когда попытки исчерпаны, категория остаётся пустой, а позиция уходит бухгалтеру на ручной разбор. Подставлять категорию по умолчанию нельзя, иначе ошибка превращается в запись в учёте.
Запись результата в учётную систему делает отдельная задача после проверки человеком или после выборочного контроля. Права на запись проверяет сервер учётной системы, а ключ модели прав на запись лишён. Если шаг нужно встроить в процесс компании целиком, посмотрите страницу про автоматизацию бизнес-процессов.
Тестовое и рабочее окружения получают разные соединения с разными ключами и разными лимитами. Ключ из теста нельзя использовать в рабочем потоке, даже временно: так проще отозвать доступ одной из сторон и проще понять по журналу поставщика, чьи это были вызовы. Раз в квартал ключи перевыпускают, а старые отзывают.
Первым делом закрепите версию пакета и создайте отдельное соединение под один конвейер. Вторым идёт разметка эталонного набора, и лишь третьим — подключение шага к рабочему потоку.