Airflow trigger для ИИ-DAG — это TriggerDagRunOperator в конце основного DAG: он создаёт запуск целевого DAG сразу после выгрузки данных и передаёт ему параметры через conf. Такая связка нужна, когда выгрузка и разбор моделью живут в разных DAG, у каждого свой владелец и своё расписание. Если ИИ-шаг всего один и идёт сразу за выгрузкой, проще держать его задачей внутри того же DAG.

Связка двух DAG

TL;DR

Основной DAG заканчивает выгрузку и последней задачей вызывает TriggerDagRunOperator с trigger_dag_id целевого DAG; всё, что нужно ИИ-шагу, передаётся в conf.

Условный пример: ночной DAG собирает заказы из учётной системы в хранилище, а второй DAG просит модель подготовить для руководителя пояснительную записку по отклонениям. Эти процессы удобно развести. Выгрузкой занимается группа данных, запиской отдел аналитики, и у модельного шага другие сроки, другие сбои и другая стоимость повтора. Запуск второго DAG от первого сохраняет порядок: записка стартует только после того, как данные легли в хранилище.

Документация описывает оператор лаконично: он запускает DAG из другого DAG, а в примере указаны task_id, trigger_dag_id и conf. Остальные параметры приведены в справочнике провайдера standard. Связку важно отличать от ожидания: триггер сам создаёт новый запуск, а Sensor ждёт события внутри уже идущего DAG, и это другой механизм с другими параметрами.

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

Параметры запуска

ПараметрЧто делает по документацииКак применить для ИИ-шага
trigger_dag_idУказывает DAG, который нужно запуститьИдентификатор DAG с вызовом модели
confПередаёт конфигурацию запускаНомер партии выгрузки и путь к данным, без самих данных
trigger_run_idЗадаёт идентификатор запуска, иначе он генерируетсяСтрока из номера партии, чтобы повтор узнавался
logical_dateЗадаёт логическую дату запускаДата отчётного периода
fail_when_dag_is_pausedВызывает ошибку, если целевой DAG на паузеГромкий отказ вместо тихого пропуска

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

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

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

Какой процесс у вас сейчас запускается по часам вместо события?

Прийти на Discovery →

Идентификатор и дубли

Основной DAG легко перезапустить: повторить упавшую задачу, прогнать за прошлую дату, запустить вручную. Каждый такой перезапуск дойдёт до триггера и создаст ещё один запуск ИИ-DAG, то есть ещё один платный вызов модели и, возможно, второе письмо руководителю. Защита строится на двух параметрах: trigger_run_id и skip_when_already_exists.

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

  • Одна партия выгрузки — один идентификатор запуска, собранный из её номера.
  • Повторный запуск за ту же дату пропускается; пересчёт возможен только по явной просьбе владельца процесса.
  • Пересчёт выполняется явно и записывается в журнал с указанием, кто его назначил.
  • Побочные действия ИИ-DAG, например отправка письма, защищены собственным ключом идемпотентности.

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

Ждать ли результата

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

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

  1. Определите, нужен ли основному DAG результат ИИ-DAG для следующих задач.
  2. Если нужен, включите ожидание завершения и перечислите допустимые итоговые состояния явно.
  3. Если ИИ-шаг может идти долго, выберите отложенное ожидание, чтобы освободить слот воркера.
  4. Проверьте на тестовой партии оба исхода: успешный разбор и провал вызова модели.

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

Журнал и остановка

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

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

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

// начните отсюда

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

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

Как запустить один DAG из другого в Airflow?
Используйте TriggerDagRunOperator: укажите trigger_dag_id целевого DAG и при необходимости передайте параметры через conf. Остальные настройки приведены в справочнике провайдера standard.
Как передать данные в запускаемый DAG?
Через параметр conf, но только небольшие значения вроде номера партии или адреса файла. Тексты, персональные данные и секреты в конфигурации хранить нельзя, целевой DAG читает их из хранилища сам.
Как защититься от дублей запуска?
Собирайте trigger_run_id из номера партии и включайте skip_when_already_exists, чтобы запуск для той же логической даты пропускался. Побочные действия вроде отправки письма защищайте отдельным ключом.
Нужно ли ждать завершения запущенного DAG?
Только если основному DAG нужен результат. Параметр wait_for_completion включает ожидание, а режим deferrable освобождает слот воркера на время ожидания.
Чем TriggerDagRunOperator отличается от Sensor?
Оператор создаёт новый запуск другого DAG. Sensor ждёт события внутри уже идущего DAG и выпускает следующие задачи только после наступления события.