- Зачем: - нужно зафиксировать состояние трекера после рабочего коммита issue 12. - Что: - задача 12 переведена в done. - журнал coordinator-loop дополнен ревью-гейтом и коммит-гейтом. - Проверка: - git diff --cached --check.
15 KiB
15 KiB
Status: done
Airflow-DAG — пульт управления генератором
Parent
.scratch/generator-model-time-startup-history/PRD.md
Why
Даже с runbook и глаголами управление генератором остаётся консольным. Идея (пользователь, 2026-06-14; приоритет поднят 2026-07-04): параметризованный DAG в Airflow как основной человеческий интерфейс стенда — форма в веб-UI, где выбираются операция, профиль, длительность. Плюсы: валидация настроек против манифеста ещё до запуска, наглядный статус, понятные ошибки.
Учебный бонус: такой DAG сам по себе учебный материал — живой пример параметризованной оркестрации (тема урока 4), что ближе к цели курса, чем устройство генератора.
Решения (дооформление 2026-07-04)
- Без Docker-доступа из Airflow. Docker socket в контейнеры Airflow не
прокидываем. Конечные операции — это обычный Python-код генератора, которому
нужны только настройки и Kafka; таски DAG выполняют его сами, в worker'е
Airflow. Урок из проекта
airflow-greenplum-solution(пользователь): DAG говорит с управляемой системой по сетевому протоколу, жизненный цикл контейнеров — только хост/Makefile. - Границы пульта.
make up,make clean(полный сброс) и старт/стоп live-сервиса — консоль; операцииcontinueв пульте нет намеренно. Backfill и import работают только на чистом стенде (пустые топики данных); грязный стенд лечитсяmake cleanс консоли — пульт при отказе прямо подсказывает это. - Генератор не автостартует. Сейчас
make upподнимает live-генератор вместе со стендом (docker-compose.yml,restart: unless-stopped) — с пультом это неверно: любой backfill был бы сразу отбит предпроверкой, а выключить live из UI нельзя. Live становится явным действием (make generator-continue); механизм — например, compose-профиль для сервисаgenerator, выбрать при реализации и поправить make-цели. - Три операции, а не четыре.
backfill|import|check. Отдельногоexportнет: артефакт в коде пишется только по ходу backfill (StartupHistoryArtifactBuilderвнутри_run_backfill), «выгрузки задним числом» не существует — поэтому у backfill есть полеartifact_path(«сохранить артефакт», пусто — не сохранять). - Один DAG с параметром «операция», а не несколько DAG по операциям.
- После backfill/import DAG сам запускает ETL (
TriggerDagRunOperatorна существующий ETL-DAG) и дожидается его завершения перед check:wait_for_completionв Airflow 2.10.5 по умолчаниюFalse(сверено по Context7, ревью Codex 2026-07-04) — простой trigger запустил бы check раньше конца ETL. Менти видит цепочку «мир → пайплайн → витрины».
What to build
- Код генератора доступен таскам Airflow. Варианты: смонтировать
generator/srcтомом (как уже смонтирован./sql) и добавить вPYTHONPATH, либо ставить пакет вDockerfile.airflow. Монтировать нужно во все Airflow-сервисы (scheduler, webserver, worker): список профилей в форме читается изPROFILESпри разборе DAG. Критерий выбора — правка генератора не должна требовать лишних пересборок. - Выравнивание зависимостей: версии сейчас расходятся — kafka-python
2.0.5 у генератора против 2.0.6 в
airflow/requirements.txt; привести к одной.prometheus-clientв образе Airflow нет — либо добавить, либо (лучше) точка входа backfill не поднимает HTTP-сервер метрик (service.py:65-66вызывается безусловно — понадобится небольшая склейка; это единственное место, где честно появляется новый код, зафиксировать его в PR). - DAG
generator_controlвairflow/dags/:schedule=None,max_active_runs=1, форма запуска наParam:operation:backfill|import|check(выпадающий список);profile: изPROFILES(launch.py), список брать динамически;duration: строка вида6h/2d, пусто — из профиля;- переопределения мира для backfill (минимум скорость и seed) — через
механизм
overridesвbuild_launch_env; artifact_path: для backfill — куда сохранить артефакт (пусто — не сохранять), для import — что импортировать; дефолт в общем томе./data.
- Ветвление по операции —
BranchPythonOperator, в стиле существующих DAG стенда. - Каждая операция начинается с валидации, падение — до любых записей:
- backfill и import: топики данных пусты (переиспользовать
KafkaTopicInspector.assert_data_topics_empty) и STG-таблицы ClickHouse пусты — батч-ETL пересобирает ODS из всего STG (sql/ods/20_stg_to_ods.sql), поэтому пустых топиков мало: старый мир в STG смешался бы с новым. При отказе в логе — подсказка проmake clean; - import дополнительно: артефакт читается, манифест совместим с выбранными
настройками (готовые проверки
startup_history_artifact.py), в лог — перечень разошедшихся полей; - вспомогательно: проверка «live не работает» по метрикам
generator:9109внутри сети compose (см. Notes об ограничениях).
- backfill и import: топики данных пусты (переиспользовать
- Тела операций — вызовы существующего кода:
build_launch_env+ прогон генерации доT_end(backfill), функцииstartup_history_artifact(import), сверка контрольных чисел с манифестом (check). Источник манифеста для check — compact-топик (он есть и после backfill, и после import); артефакт — запасной вариант. retries=0у содержательных тасков: оператор должен сразу видеть ошибку (проверенный приём из greenplum-проекта).- Документация: в runbook
docs/runbooks/startup-history.md— раздел «Пульт в Airflow» как основной путь и явные границы пульта (что остаётся консолью и почему, включая «continue в пульте нет намеренно»); тонкость нестандартного мира:make generator-continueдля мира с переопределёнными настройками требует тех же настроек, громкий отказ подскажет разошедшиеся поля. Ссылки изdocs/OPERATIONS.mdиREADME.md; правки compose и Makefile описать в том же PR (правило репозитория).
Acceptance criteria
- После
make up(генератор не автостартует) стартовую историю выбранного профиля/длительности можно создать из веб-UI Airflow, не открывая консоль и не выставляя переменных окружения; следом ETL запускается из той же цепочки иcheckзелёный. - Импорт артефакта из веб-UI: несовместимый артефакт отклоняется до записи в топики, в логе таска — перечень разошедшихся полей.
- Backfill/import на непустом стенде отклоняются предпроверкой с
подсказкой про
make clean; миры не смешиваются. - Артефакт, сохранённый при backfill из веб-UI, пригоден для консольного импорта (формат один и тот же), и наоборот.
- Операция check сверяет контрольные числа ClickHouse с манифестом и падает при расхождении.
- Docker недоступен из контейнеров Airflow (socket не монтируется) — это граница решения, а не упущение.
- Консольные глаголы работают как раньше;
scripts/*и DAG сходятся в одномlaunch.py, дублирования логики запуска нет. make upбольше не запускает live;make generator-continueзапускает его явно; существующие сценарии (generated-history-analytics, CI) не сломаны.- Runbook и
docs/OPERATIONS.mdописывают пульт как основной путь и его границы. - DAG остаётся читаемым менти: витрина параметризованной оркестрации, сложность живёт в коде генератора.
Notes
- Перед реализацией сверить API формы (
Param, enum, описания полей) иTriggerDagRunOperatorдля Airflow 2.10 через MCP Context7 — правило репозитория. Configгенератора читает env процесса при создании: в таске собирать окружение черезbuild_launch_envи применять к процессу таска (executor — LocalExecutor, таск живёт в своём процессе). Не забытьGEN_DATA_DIR:build_launch_envего не задаёт, дефолт/data, а в контейнерах Airflow данные смонтированы в/opt/airflow/data— без явной установкиConfig()упадёт, словари*.jsonlне найдутся.- Права на файлы: артефакт из worker'а пишется uid'ом Airflow (50000); в
контейнере генератора
./dataсмонтирован read-only, консольные скрипты монтируют каталог артефактов отдельно. Чтение работает везде, перезапись чужого файла — не всегда; одну строку об этом — в runbook. - Проверка
generator:9109— вспомогательная эвристика, не защита: ошибка DNS означает «сервис не поднят» (не ошибку таска); при крашлупе после громкого отказа порт мигает (HTTP-сервер стартует до валидации state); эфемерные контейнерыdocker compose runпод этим именем не видны. Основная защита от смешивания — пустота топиков данных. - Backfill — блокирующий таск на минуты; при
max_active_runs=1для ручного стенда это нормально. - Грабли greenplum-проекта, применимые здесь: DAG создаются на паузе — в инструкции приёмки не забыть unpause перед trigger; автоматические прогоны гонять через REST API, а не CLI внутри контейнера.
- Возможное развитие (не в скоупе): операция «очистить данные» из пульта
(прецедент удаления/пересоздания топиков из DAG уже есть в
airflow/dags/utils/kafka_helpers.py) — сняла бы требованиеmake cleanперед новым миром. - Артефакты
daily-wave, сделанные до задачи 14, протухнут после неё (антисмешивание отработает громко) — при пересечении работ это ожидаемо. - Рекомендуемый режим ревью по coordinator-loop: обычный — state и сериализацию задача не трогает, код генератора переиспользуется как есть.
Зависимости (проверенные предпосылки)
07-startup-history-portable-artifact-and-usage-docs.md— сделана.11-generator-launch-verbs-and-profiles.md— сделана.- Мягкая связь с
14-fast-teaching-profile.md: быстрый профиль появится в выпадашке сам, если список профилей брать изPROFILESдинамически; блокером не является.