From ac7a97350430750b547fc54370463410183da8b5 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Wed, 22 Jul 2026 21:45:02 +0300 Subject: [PATCH] =?UTF-8?q?feat(airflow):=20=D0=BF=D1=83=D0=BB=D1=8C=D1=82?= =?UTF-8?q?=20=D1=81=D1=82=D0=B0=D0=BB=20world=5Finit,=20=D0=B4=D0=BE?= =?UTF-8?q?=D0=B1=D0=B0=D0=B2=D0=BB=D0=B5=D0=BD=20DAG=20world=5Fnext=5Fday?= =?UTF-8?q?=20(#4)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - список DAG'ов должен читаться лесенкой ddl_init → world_init → world_next_day, а путь менти — проходиться пустыми формами (issue #4, спека редизайна пути менти, решения 2–3). - Что: - generator_control переименован в world_init, дефолт операции — import; next-day ушёл из выпадашки в отдельный DAG; - новый беспараметрный world_next_day: расписание */30 * * * *, создаётся на паузе, catchup=False, max_active_runs=1; общие задачи вынесены в airflow/dags/utils/startup_history_tasks.py; - доки и контрактные тесты обновлены синхронно; быстрый старт README — без make ddl, схему создаёт DAG ddl_init. - Проверка: - make test (210 + 31) и make lint зелёные; - живая приёмка на чистом стенде: world_init пустой формой импортировал эталонный мир за 217 с (3 дня, 280 437 событий), world_next_day после снятия с паузы добавляет ровно один день за прогон, дашборд Superset собирается. Co-Authored-By: Claude Fable 5 --- CONTEXT.md | 2 +- README.md | 29 +++-- airflow/dags/utils/startup_history_tasks.py | 59 ++++++++++ ...rator_control_dag.py => world_init_dag.py} | 104 ++---------------- airflow/dags/world_next_day_dag.py | 69 ++++++++++++ docs/ARCHITECTURE.md | 45 +++++--- docs/OPERATIONS.md | 54 ++++----- docs/REPO_MAP.md | 4 +- docs/TEST_PLAN.md | 2 +- docs/runbooks/startup-history.md | 20 ++-- .../clickstream_generator/airflow_control.py | 2 +- ...ontract.py => test_world_dags_contract.py} | 70 ++++++------ 12 files changed, 268 insertions(+), 192 deletions(-) create mode 100644 airflow/dags/utils/startup_history_tasks.py rename airflow/dags/{generator_control_dag.py => world_init_dag.py} (69%) create mode 100644 airflow/dags/world_next_day_dag.py rename generator/tests/{test_generator_control_dag_contract.py => test_world_dags_contract.py} (84%) diff --git a/CONTEXT.md b/CONTEXT.md index b985e33..1c25e86 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -87,7 +87,7 @@ user_domain_id (пользователь, постоянный) это не новые значения слова, а та же сущность) — сгенерированное прошлое (заливка `K → ∞` + заморозка состояния), с которого живой стенд стартует непрерывно. По ADR-0006 она **несущая**: именно с неё свежий стенд получает историю с первой минуты. - Механизм реализован через Airflow DAG `generator_control` и служебный чистый + Механизм реализован через Airflow DAG `world_init` и служебный чистый путь `make generated-history-analytics`; решение про часы — [ADR-0005](./docs/adr/0005-generator-model-clock.md). ### Мир (стенда) diff --git a/README.md b/README.md index b32e929..86a6420 100644 --- a/README.md +++ b/README.md @@ -9,7 +9,7 @@ витринами. Поток данных коротко: -- **стартовая история**: `generator backfill → Kafka → ClickHouse (STG) → +- **стартовая история**: `world_init → Kafka → ClickHouse (STG) → batch STG → ODS → DDS → DM → Superset`. - **живое продолжение**: `generator live → Kafka → ClickHouse (STG) → batch ETL → Superset`. @@ -31,20 +31,29 @@ Перед первой командой нужны `Docker` с `docker compose`, `make`, `bash`, `curl`, `git` и `uv`. `uv` нужен для локальных Python-проверок и команд разработки. -Для ручной работы поднимите стенд и создайте стартовую историю через Airflow: +Для ручной работы поднимите стенд: ```bash make up -make ddl docker compose ps ``` -Откройте Airflow: `http://localhost:8080` (`admin/admin`). Сначала снимите -паузу с DAG `etl_pipeline` (переключатель слева от имени): на свежем стенде он -создаётся на паузе, и `backfill` откажется стартовать. Затем запустите -`generator_control` с операцией `backfill`: DAG создаст стартовую историю, -запустит ETL и выполнит `check`. `make up` не запускает live-генератор; live -включается отдельно командой `make generator-continue`. +Дальше всё делается в Airflow: `http://localhost:8080` (`admin/admin`). +Список DAG'ов читается лесенкой сверху вниз; на свежем стенде все DAG'и +создаются на паузе, поэтому перед запуском снимайте паузу переключателем +слева от имени. + +1. `ddl_init` — снимите паузу и запустите: DAG создаст схему ClickHouse + (отдельная команда в терминале не нужна). +2. `etl_pipeline` — только снимите паузу: его запустит следующий шаг. +3. `world_init` — снимите паузу и запустите с пустой формой: DAG импортирует + эталонный мир, запустит ETL и сверит витрины. +4. `world_next_day` — когда захотите добавить ровно один модельный день, + запустите его с пустой формой. Расписание задано каждые 30 минут, но по + умолчанию DAG стоит на паузе. + +`make up` не запускает live-генератор; live включается отдельно командой +`make generator-continue`. После обновления репозитория снова выполните `make up`: команда пересобирает Airflow-образ и подтягивает новые зависимости и DAG-и. Superset-дэшборд @@ -53,7 +62,7 @@ Airflow-образ и подтягивает новые зависимости Для полностью автоматического чистого прогона из консоли есть команда — это тот же путь, что выше через Airflow UI, но одной командой и без ручных шагов -(DDL она применяет сама, отдельный `make ddl` не нужен): +(схему ClickHouse она применяет сама): ```bash make generated-history-analytics diff --git a/airflow/dags/utils/startup_history_tasks.py b/airflow/dags/utils/startup_history_tasks.py new file mode 100644 index 0000000..5aa87d6 --- /dev/null +++ b/airflow/dags/utils/startup_history_tasks.py @@ -0,0 +1,59 @@ +"""Общие задачи Airflow для операций над миром стенда.""" + +from airflow.exceptions import AirflowException +from airflow.models.dag import DagModel +from airflow.utils.session import provide_session +from airflow_clickhouse_plugin.hooks.clickhouse import ClickHouseHook + +from clickstream_generator.airflow_control import ( + assert_clickhouse_matches_manifest, + assert_live_generator_not_running_for_next_day, + assert_next_day_snapshot, + build_next_day_env, + load_manifest_from_kafka, + load_state_from_kafka, + run_next_day, + target_dag_trigger_error, +) + + +def precheck_next_day(**context) -> None: + """Проверяет точку продолжения и готовит настройки мира.""" + assert_live_generator_not_running_for_next_day() + manifest = load_manifest_from_kafka() + state = load_state_from_kafka() + assert_next_day_snapshot(manifest, state) + context["ti"].xcom_push( + key="generator_env", + value=build_next_day_env(manifest), + ) + + +def run_next_day_task(**context) -> None: + """Генерирует следующий модельный день от проверенного state.""" + env = context["ti"].xcom_pull( + task_ids="precheck_next_day", + key="generator_env", + ) + run_next_day(env) + + +def check_manifest_task(**context) -> None: + """Сверяет витрину ClickHouse с manifest стартовой истории.""" + manifest = load_manifest_from_kafka() + hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default", database="default") + assert_clickhouse_matches_manifest(manifest, hook) + + +@provide_session +def assert_target_dag_not_paused(dag_id: str, session=None) -> None: + """ + Проверяет, что зависимый DAG можно запустить. + + Context7: для Airflow 2.10.5 у старого TriggerDagRunOperator нет надёжного + fail_when_dag_is_paused, поэтому паузу проверяем заранее через DagModel. + """ + dag_model = session.query(DagModel).filter(DagModel.dag_id == dag_id).one_or_none() + error = target_dag_trigger_error(dag_id, dag_model) + if error: + raise AirflowException(error) diff --git a/airflow/dags/generator_control_dag.py b/airflow/dags/world_init_dag.py similarity index 69% rename from airflow/dags/generator_control_dag.py rename to airflow/dags/world_init_dag.py index db5847f..3d93de9 100644 --- a/airflow/dags/generator_control_dag.py +++ b/airflow/dags/world_init_dag.py @@ -10,36 +10,28 @@ from __future__ import annotations from datetime import datetime, timedelta from airflow import DAG -from airflow.exceptions import AirflowException -from airflow.models.dag import DagModel from airflow.models.param import Param from airflow.operators.empty import EmptyOperator from airflow.operators.python import BranchPythonOperator, PythonOperator from airflow.operators.trigger_dagrun import TriggerDagRunOperator -from airflow.utils.session import provide_session from airflow.utils.trigger_rule import TriggerRule from airflow_clickhouse_plugin.hooks.clickhouse import ClickHouseHook from clickstream_generator.airflow_control import ( KAFKA_BOOTSTRAP_SERVERS, - assert_clickhouse_matches_manifest, - assert_expected_t_end, assert_live_generator_not_running, - assert_live_generator_not_running_for_next_day, - assert_next_day_snapshot, assert_stand_clean, build_control_env, - build_next_day_env, default_artifact_path, - load_manifest_from_kafka, - load_state_from_kafka, run_backfill, run_import, - run_next_day, - target_dag_trigger_error, validate_import_artifact, ) from clickstream_generator.launch import PROFILES +from utils.startup_history_tasks import ( + assert_target_dag_not_paused, + check_manifest_task, +) default_args = { @@ -107,8 +99,6 @@ def choose_operation(**context) -> str: return "check_etl_not_paused_before_backfill" if operation == "import": return "check_etl_not_paused_before_import" - if operation == "next-day": - return "check_etl_not_paused_before_next_day" if operation == "check": return "check_only" raise ValueError(f"Неизвестная операция: {operation}") @@ -151,56 +141,9 @@ def run_import_task(**context) -> None: run_import(env, artifact_path) -def precheck_next_day(**context) -> None: - """Проверяет точку продолжения и готовит неизменные настройки мира.""" - assert_live_generator_not_running_for_next_day() - manifest = load_manifest_from_kafka() - state = load_state_from_kafka() - assert_next_day_snapshot(manifest, state) - assert_expected_t_end( - str(_param(context, "expected_t_end") or "").strip(), - str(manifest["model_t_end"]), - ) - context["ti"].xcom_push( - key="generator_env", - value=build_next_day_env(manifest), - ) - - -def run_next_day_task(**context) -> None: - """Генерирует следующий модельный день от проверенного state.""" - env = context["ti"].xcom_pull( - task_ids="precheck_next_day", - key="generator_env", - ) - run_next_day(env) - - -def check_manifest_task(**context) -> None: - """Сверяет витрину ClickHouse с manifest стартовой истории.""" - manifest = load_manifest_from_kafka() - hook = ClickHouseHook(clickhouse_conn_id="clickhouse_default", database="default") - assert_clickhouse_matches_manifest(manifest, hook) - - -@provide_session -def assert_target_dag_not_paused(dag_id: str, session=None) -> None: - """ - Проверяет, что зависимый DAG можно запустить. - - Context7: для Airflow 2.10.5 у старого TriggerDagRunOperator нет надёжного - fail_when_dag_is_paused, поэтому паузу проверяем заранее через DagModel. - """ - dag_model = session.query(DagModel).filter(DagModel.dag_id == dag_id).one_or_none() - error = target_dag_trigger_error(dag_id, dag_model) - if error: - raise AirflowException(error) - - -# Context7, Airflow 2.10.5: Param поддерживает enum, а max_active_runs ограничивает -# число одновременных DAG run. Поэтому форму и блокировку next-day держим в DAG. +# Context7, Airflow 2.10.5: Param поддерживает enum. with DAG( - dag_id="generator_control", + dag_id="world_init", description="Пульт стартовой истории генератора", default_args=default_args, schedule=None, @@ -211,13 +154,13 @@ with DAG( tags=["generator", "startup-history"], params={ "operation": Param( - "backfill", + "import", type="string", - enum=["backfill", "import", "next-day", "check"], + enum=["backfill", "import", "check"], title="Операция", description=( - "Что сделать: создать историю, импортировать артефакт, " - "добавить следующий день или проверить витрины." + "Что сделать: импортировать артефакт, создать историю " + "или проверить витрины." ), ), "profile": Param( @@ -255,15 +198,6 @@ with DAG( "Import: что читать; пусто — эталонный мир из репозитория." ), ), - "expected_t_end": Param( - None, - type=["null", "string"], - title="Ожидаемая граница next-day", - description=( - "Необязательный model_t_end до запуска. Защищает от " - "повторной доливки того же дня." - ), - ), }, ) as dag: route = BranchPythonOperator( @@ -289,15 +223,6 @@ with DAG( python_callable=run_import_task, ) - precheck_next_day_task = PythonOperator( - task_id="precheck_next_day", - python_callable=precheck_next_day, - ) - next_day_task = PythonOperator( - task_id="run_next_day", - python_callable=run_next_day_task, - ) - check_etl_not_paused_before_backfill = PythonOperator( task_id="check_etl_not_paused_before_backfill", python_callable=assert_target_dag_not_paused, @@ -308,12 +233,6 @@ with DAG( python_callable=assert_target_dag_not_paused, op_kwargs={"dag_id": ETL_DAG_ID}, ) - check_etl_not_paused_before_next_day = PythonOperator( - task_id="check_etl_not_paused_before_next_day", - python_callable=assert_target_dag_not_paused, - op_kwargs={"dag_id": ETL_DAG_ID}, - ) - trigger_etl = TriggerDagRunOperator( task_id="trigger_etl", trigger_dag_id=ETL_DAG_ID, @@ -343,14 +262,11 @@ with DAG( route >> [ check_etl_not_paused_before_backfill, check_etl_not_paused_before_import, - check_etl_not_paused_before_next_day, check_only, ] check_etl_not_paused_before_backfill >> precheck_backfill_task check_etl_not_paused_before_import >> precheck_import_task - check_etl_not_paused_before_next_day >> precheck_next_day_task precheck_backfill_task >> backfill_task >> trigger_etl precheck_import_task >> import_task >> trigger_etl - precheck_next_day_task >> next_day_task >> trigger_etl trigger_etl >> check_after_etl >> done check_only >> done diff --git a/airflow/dags/world_next_day_dag.py b/airflow/dags/world_next_day_dag.py new file mode 100644 index 0000000..7c6beb6 --- /dev/null +++ b/airflow/dags/world_next_day_dag.py @@ -0,0 +1,69 @@ +"""DAG добавления одного модельного дня в мир стенда.""" + +from datetime import datetime, timedelta + +from airflow import DAG +from airflow.operators.python import PythonOperator +from airflow.operators.trigger_dagrun import TriggerDagRunOperator + +from utils.startup_history_tasks import ( + assert_target_dag_not_paused, + check_manifest_task, + precheck_next_day, + run_next_day_task, +) + + +ETL_DAG_ID = "etl_pipeline" + +default_args = { + "owner": "airflow", + "depends_on_past": False, + "email_on_failure": False, + "email_on_retry": False, + "retries": 0, + "retry_delay": timedelta(minutes=1), +} + + +# Context7, Airflow 2.10.5: schedule принимает cron-строку. Расписание задано +# заранее, но новый DAG остаётся на паузе до отдельного решения. +with DAG( + dag_id="world_next_day", + description="Добавление одного модельного дня в мир стенда", + default_args=default_args, + schedule="*/30 * * * *", + start_date=datetime(2024, 1, 1), + catchup=False, + max_active_runs=1, + is_paused_upon_creation=True, + tags=["generator", "startup-history"], +) as dag: + check_etl_not_paused = PythonOperator( + task_id="check_etl_not_paused", + python_callable=assert_target_dag_not_paused, + op_kwargs={"dag_id": ETL_DAG_ID}, + ) + precheck = PythonOperator( + task_id="precheck_next_day", + python_callable=precheck_next_day, + ) + generate = PythonOperator( + task_id="run_next_day", + python_callable=run_next_day_task, + ) + trigger_etl = TriggerDagRunOperator( + task_id="trigger_etl", + trigger_dag_id=ETL_DAG_ID, + conf={"full_refresh": True}, + wait_for_completion=True, + allowed_states=["success"], + failed_states=["failed"], + poke_interval=30, + ) + check = PythonOperator( + task_id="check_after_etl", + python_callable=check_manifest_task, + ) + + check_etl_not_paused >> precheck >> generate >> trigger_etl >> check diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index eb079c9..cbb0f12 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -27,8 +27,9 @@ flowchart LR subgraph AF["Airflow"] DAG1["ddl_init"] - DAG2["generator_control"] - DAG3["etl_pipeline"] + DAG2["world_init"] + DAG3["world_next_day"] + DAG4["etl_pipeline"] end subgraph GEN["Generator"] @@ -57,7 +58,8 @@ flowchart LR V[витрины VIEW] end - DAG2 -->|startup-history backfill/import| K + DAG2 -->|import/backfill| K + DAG3 -->|следующий день| K G -->|live после make generator-continue| K K -->|MV| S S -->|batch| O @@ -66,12 +68,13 @@ flowchart LR D1 & D2 -->|VIEW| V DAG1 -.->|DDL| STG & ODS & DDS & DM - DAG3 -.->|batch| ODS & DDS + DAG4 -.->|batch| ODS & DDS ``` В учебном стенде предусмотрены два пути загрузки: -- `startup-history`: DAG `generator_control` создаёт или импортирует историю, добавляет следующий модельный день, запускает ETL и проверяет витрины; +- `startup-history`: `world_init` импортирует или создаёт историю, а + `world_next_day` добавляет один модельный день; оба запускают ETL и проверяют витрины; - `live`: генератор запускается явно через `make generator-continue`, когда нужна непрерывная подача новых событий. ### Слои и их назначение @@ -80,8 +83,9 @@ flowchart LR flowchart LR subgraph AF["Airflow"] DAG1["ddl_init"] - DAG2["generator_control"] - DAG3["etl_pipeline"] + DAG2["world_init"] + DAG3["world_next_day"] + DAG4["etl_pipeline"] end subgraph GEN["Generator"] @@ -106,7 +110,8 @@ flowchart LR DM_T["VIEW"] end - DAG2 -->|startup-history| KAFKA + DAG2 -->|стартовый мир| KAFKA + DAG3 -->|следующий день| KAFKA G -->|live| KAFKA KAFKA -->|MV| STG_T STG_T -->|batch| ODS_T @@ -114,7 +119,7 @@ flowchart LR ODS_T -.->|ошибки| DQ DAG1 -.->|DDL| L1 & L2 & L3 & L4 - DAG3 -.->|batch| ODS_T & DDS_T + DAG4 -.->|batch| ODS_T & DDS_T ``` --- @@ -387,10 +392,12 @@ sequenceDiagram CH-->>User: ✅ Структура БД создана alt Startup-history режим - User->>Airflow: Trigger generator_control (backfill/import) + User->>Airflow: Trigger world_init с пустой формой Airflow->>K: события стартовой истории Airflow->>Airflow: trigger etl_pipeline + check K-->>User: ✅ История в Kafka и витринах + User->>Airflow: Trigger world_next_day с пустой формой + Airflow->>K: события следующего модельного дня else Live режим User->>Compose: make generator-continue loop каждые 1-10 секунд @@ -637,7 +644,8 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic; ```python # airflow/dags/ddl_init_dag.py — создание баз/таблиц -# airflow/dags/generator_control_dag.py — backfill/import/next-day/check стартовой истории +# airflow/dags/world_init_dag.py — import/backfill/check стартового мира +# airflow/dags/world_next_day_dag.py — добавление одного модельного дня # airflow/dags/etl_pipeline_dag.py — основной ETL (STG→ODS→DDS→DM) # airflow/dags/kafka_load_dag.py — архивный ручной путь из JSONL, не основной контур @@ -645,22 +653,25 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic; # - DDL и трансформации выполняются явными SQL-task через ClickHouseOperator; # - SQL-файлы вызываются по фиксированным путям; # - загрузка может идти двумя путями: -# 1) startup-history через DAG `generator_control`; +# 1) стартовый мир через `world_init` и рост через `world_next_day`; # 2) live-поток через явный `make generator-continue`. # # Базовый demo-сценарий: -# ddl_init -> generator_control(backfill/import) -> etl_pipeline -> check +# ddl_init -> world_init(import) -> etl_pipeline -> check # Расширенный учебный сценарий: # make generator-continue + периодический etl_pipeline ``` -**DAG `generator_control`**: -- `backfill`: создаёт стартовую историю через генератор -- `import`: импортирует портативный артефакт стартовой истории -- `next-day`: пакетно добавляет следующий модельный день от текущего слепка мира +**DAG `world_init`**: +- `import` по умолчанию импортирует портативный артефакт стартового мира +- `backfill` создаёт стартовую историю через генератор - `check`: сверяет ClickHouse с manifest стартовой истории - После `backfill` и `import` запускает `etl_pipeline` с `full_refresh` +**DAG `world_next_day`** без параметров пакетно добавляет следующий модельный +день, запускает `etl_pipeline` с `full_refresh` и сверяет manifest. У него задано +расписание каждые 30 минут, но DAG создаётся на паузе и не выполняет пропущенные интервалы. + **Подключение к ClickHouse:** - Connection: `clickhouse_default` - URL: `clickhouse://default:123456@clickhouse:9000/default` (native TCP для Airflow plugin) diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index d44d18b..8cfa870 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -48,45 +48,48 @@ volumes или live-генератором. Для стыка backfill/live от ## Airflow DAGs -Штатный ручной путь начинается с `generator_control`: чистый стенд получает -стартовую историю генератора, затем этот же DAG запускает ETL и проверку. +Штатный ручной путь начинается с `world_init`: пустая форма импортирует +эталонный мир, затем этот же DAG запускает ETL и проверку. `kafka_load` остаётся для экспериментов и совместимости учебного стенда. -### `generator_control` +### `world_init` - Запуск: ручной (`Trigger DAG`). -- Назначение: пульт стартовой истории генератора. +- Назначение: импорт или служебная сборка стартового мира. - Операции: + - `import` — операция по умолчанию: импортировать портативный артефакт, затем + запустить `etl_pipeline` и дождаться `success`; - `backfill` — создать стартовую историю, затем запустить `etl_pipeline` и дождаться `success`; - - `import` — импортировать портативный артефакт, затем запустить - `etl_pipeline` и дождаться `success`; - - `next-day` — восстановить мир из state, добавить 24 модельных часа, - затем запустить `etl_pipeline` и дождаться `success`; - `check` — сверить ClickHouse с manifest из Kafka. - Параметры: - - `operation` (`backfill` / `import` / `next-day` / `check`); + - `operation` (`import` / `backfill` / `check`); - `profile` — список берётся из `PROFILES` генератора; - `duration` — `6h`, `2d` и т.п.; пусто означает длительность профиля; - `seed`, `model_time_speed` — необязательные переопределения мира; - `artifact_path` — для `backfill` путь сохранения; для `import` путь чтения. - При пустом поле импортируется эталонный мир из репозитория; - - `expected_t_end` — необязательная ожидаемая граница перед `next-day`. - При расхождении запуск показывает ожидаемое и фактическое значения. + При пустом поле импортируется эталонный мир из репозитория. Backfill/import требуют чистый стенд: пустые data-топики Kafka и пустые `stg.*_raw`. При отказе очистите стенд через `make clean`. Операции `continue` в DAG нет: live-генератор — долгоживущий сервис, его запускают с консоли через `make generator-continue`. -`next-day` работает на непустом стенде и не использует проверку чистоты. -Перед записью пульт требует manifest, state ровно на его `T_end` и остановленный -live-генератор. Настройки мира берутся из manifest; поля `profile`, `duration`, -`seed` и `model_time_speed` формы для этой операции не применяются. Один запуск -добавляет полуоткрытый диапазон `[T_end, T_end + 24h)` в UTC. Новая граница -появляется в `boundaries`; старый manifest без поля читается как `[T0, T_end]`. -Расписание остаётся выключенным (`schedule=None`), а `max_active_runs=1` не даёт -двум доливкам выполняться параллельно. +### `world_next_day` + +- Запуск: вручную с пустой формой. +- Параметров нет. +- Один запуск восстанавливает мир из state, добавляет 24 модельных часа, + запускает `etl_pipeline` с полной пересборкой и сверяет витрины с manifest. +- Расписание задано каждые 30 минут, но DAG создаётся на паузе; `catchup=False`. + Не включайте расписание до внедрения накопительных счётчиков manifest. +- `max_active_runs=1` не даёт двум доливкам выполняться параллельно. + +`world_next_day` работает на непустом стенде и не использует проверку чистоты. +Перед записью он требует manifest, state ровно на его `T_end` и остановленный +live-генератор. Настройки мира берутся из manifest. Один запуск добавляет +полуоткрытый диапазон `[T_end, T_end + 24h)` в UTC. Новая граница появляется в +`boundaries`; старый manifest без поля читается как `[T0, T_end]`. После двух доливок проверьте завершённые стыки: @@ -98,10 +101,11 @@ make generated-history-chain-check смену browser/referer/utm внутри переходящих визитов. Она не меняет `make generated-history-runtime-check` для стыка backfill/live. -Точка фиксации `next-day` — новый manifest. Порядок записи: data-топики, state, -manifest. Автоматического отката нет. Если запуск упал до публикации manifest, +Результат запуска `world_next_day` фиксируется новым manifest. Порядок записи: +data-топики, state, manifest. Автоматического отката нет. Если запуск упал до +публикации manifest, не повторяйте доливку поверх возможного хвоста. Очистите стенд и переимпортируйте -последний исправный портативный артефакт, затем повторите `next-day`. +последний исправный портативный артефакт, затем повторите запуск `world_next_day`. Текущая версия пересчитывает накопительные счётчики и контрольные суммы по всей доступной истории data-топиков Kafka. Поэтому время выполнения и расход памяти @@ -689,8 +693,8 @@ make generated-history-runtime-check ## Быстрые проверки - Kafka ingest: наличие данных генератора в `stg.*` и типизированных строк в `ods.*`. -- Airflow UI: `http://localhost:8080` показывает DAG `ddl_init`, `generator_control`, - `kafka_load`, `etl_pipeline`; основной ручной пульт генератора — `generator_control`. +- Airflow UI: `http://localhost:8080` показывает лестницу `ddl_init` → + `world_init` → `world_next_day`, а также `kafka_load` и `etl_pipeline`. - BI: витрина `dm.v_events_enriched` отвечает за разумное время при фильтре по дате. --- diff --git a/docs/REPO_MAP.md b/docs/REPO_MAP.md index faa0254..b998da1 100644 --- a/docs/REPO_MAP.md +++ b/docs/REPO_MAP.md @@ -7,7 +7,9 @@ ### Airflow (ручной и учебный путь запуска) - `airflow/dags/ddl_init_dag.py` — инициализация схемы ClickHouse -- `airflow/dags/generator_control_dag.py` — Airflow-пульт стартовой истории: backfill/import/next-day/check +- `airflow/dags/world_init_dag.py` — импорт или служебная сборка стартового мира и проверка витрин +- `airflow/dags/world_next_day_dag.py` — беспараметрное добавление одного модельного дня +- `airflow/dags/utils/startup_history_tasks.py` — общие задачи DAG для роста и проверки мира - `airflow/dags/kafka_load_dag.py` — архивная загрузка в Kafka из JSONL; не основной источник аналитики - `airflow/dags/etl_pipeline_dag.py` — ETL процесс STG -> ODS -> DDS -> DM - `airflow/dags/utils/kafka_helpers.py` — helper-функции для Kafka diff --git a/docs/TEST_PLAN.md b/docs/TEST_PLAN.md index 57d1e59..aff7467 100644 --- a/docs/TEST_PLAN.md +++ b/docs/TEST_PLAN.md @@ -16,7 +16,7 @@ - Для smoke и CI явно задаём служебный профиль `ci`. - Полный прогон выполняем отдельно через `PROFILE=daily-wave`. -- Основной ручной путь запуска — через Airflow DAG `generator_control`. +- Основной ручной путь запуска — через Airflow DAG `world_init` с пустой формой. - Консольный чистый прогон `make generated-history-analytics` остаётся коротким повторяемым сценарием для smoke и CI. - Критерий успеха: не только `Success` DAG, но и проверки данных/ошибок/мониторинга. diff --git a/docs/runbooks/startup-history.md b/docs/runbooks/startup-history.md index 758ce7c..9cb20e9 100644 --- a/docs/runbooks/startup-history.md +++ b/docs/runbooks/startup-history.md @@ -58,9 +58,9 @@ git add data/startup_history/reference-world.json.xz При импорте DAG дважды читает и распаковывает артефакт: во время предпроверки и перед записью в Kafka. Это увеличивает время импорта, но не меняет результат. -Менти в форме `generator_control` выбирает `import` и оставляет -`artifact_path` пустым. Тогда читается эталонный мир из репозитория. Из консоли -тот же импорт запускается без указания пути: +Менти запускает `world_init` с пустой формой. Тогда читается эталонный мир из +репозитория. Следующий модельный день добавляет отдельный беспараметрный DAG +`world_next_day`. Из консоли тот же импорт запускается без указания пути: ```bash make startup-history-import @@ -68,20 +68,24 @@ make startup-history-import ## Пульт в Airflow -Основной ручной путь — DAG `generator_control` в Airflow UI: +Основной учебный путь в Airflow UI: 1. Поднимите стенд: `make up`. 2. Если DDL ещё не применён, запустите `ddl_init`. 3. Снимите паузу с `etl_pipeline`, если он ещё paused: `docker compose exec -T airflow-webserver airflow dags unpause etl_pipeline`. -4. Откройте `generator_control` и выберите `operation`. +4. Запустите `world_init` с пустой формой. По умолчанию он импортирует эталонный мир. +5. Когда нужен ещё один модельный день, запустите `world_next_day` с пустой формой. -Операции: +`world_next_day` имеет расписание каждые 30 минут, но по умолчанию стоит на паузе. +Не включайте расписание до внедрения накопительных счётчиков manifest. + +## Операции сопровождающего + +В форме `world_init` сопровождающему дополнительно доступны операции: - `backfill` — создать стартовую историю. После записи в Kafka DAG сам запускает `etl_pipeline`, ждёт завершения и выполняет `check`. -- `import` — прочитать артефакт из `artifact_path`. Несовместимый артефакт - отклоняется до записи в Kafka. - `check` — сверить ClickHouse с manifest из Kafka. Поля формы: diff --git a/generator/src/clickstream_generator/airflow_control.py b/generator/src/clickstream_generator/airflow_control.py index 75ca623..68d0171 100644 --- a/generator/src/clickstream_generator/airflow_control.py +++ b/generator/src/clickstream_generator/airflow_control.py @@ -328,7 +328,7 @@ def target_dag_trigger_error(dag_id: str, dag_model) -> str | None: if dag_model.is_paused: return ( f"DAG {dag_id} стоит на паузе: снимите паузу в Airflow UI, " - "затем повторите generator_control." + "затем повторите запуск DAG." ) return None diff --git a/generator/tests/test_generator_control_dag_contract.py b/generator/tests/test_world_dags_contract.py similarity index 84% rename from generator/tests/test_generator_control_dag_contract.py rename to generator/tests/test_world_dags_contract.py index e7dbc0d..748cd82 100644 --- a/generator/tests/test_generator_control_dag_contract.py +++ b/generator/tests/test_world_dags_contract.py @@ -1,6 +1,4 @@ -""" -Контракт DAG generator_control без запуска Airflow. -""" +"""Контракты DAG world_init и world_next_day без запуска Airflow.""" import ast from pathlib import Path @@ -8,7 +6,8 @@ from pathlib import Path from clickstream_generator.launch import PROFILES -DAG_PATH = Path(__file__).parents[2] / "airflow" / "dags" / "generator_control_dag.py" +DAG_PATH = Path(__file__).parents[2] / "airflow" / "dags" / "world_init_dag.py" +NEXT_DAY_DAG_PATH = Path(__file__).parents[2] / "airflow" / "dags" / "world_next_day_dag.py" REPO_ROOT = Path(__file__).parents[2] @@ -65,7 +64,7 @@ def test_dag_file_exists_and_uses_dynamic_profiles(): """DAG берёт варианты профилей из PROFILES, а не из ручного списка.""" text = DAG_PATH.read_text(encoding="utf-8") - assert "dag_id=\"generator_control\"" in text + assert "dag_id=\"world_init\"" in text assert "sorted(PROFILES)" in text for profile in PROFILES: assert profile not in {"hardcoded-profile"} @@ -76,12 +75,13 @@ def test_trigger_form_has_expected_param_enums(): text = DAG_PATH.read_text(encoding="utf-8") param_defaults = _declared_param_defaults() - assert 'enum=["backfill", "import", "next-day", "check"]' in text + assert 'enum=["backfill", "import", "check"]' in text + assert param_defaults["operation"] == "import" assert "enum=sorted(PROFILES)" in text assert param_defaults["profile"] == "daily-wave" assert '"duration": Param(' in text assert '"artifact_path": Param(' in text - assert '"expected_t_end": Param(' in text + assert '"expected_t_end": Param(' not in text def test_trigger_form_marks_only_optional_params_as_nullable(): @@ -93,7 +93,6 @@ def test_trigger_form_marks_only_optional_params_as_nullable(): "seed", "model_time_speed", "artifact_path", - "expected_t_end", }: assert "null" in param_types[name] for name in {"operation", "profile"}: @@ -101,7 +100,7 @@ def test_trigger_form_marks_only_optional_params_as_nullable(): def test_dag_branches_and_waits_for_etl_completion(): - """Backfill/import/next-day запускают ETL и ждут завершения перед check.""" + """Backfill/import запускают ETL и ждут завершения перед check.""" text = DAG_PATH.read_text(encoding="utf-8") assert "BranchPythonOperator" in text @@ -112,41 +111,43 @@ def test_dag_branches_and_waits_for_etl_completion(): assert 'failed_states=["failed"]' in text -def test_next_day_has_own_boundary_precheck_and_serial_execution(): - """Next-day не использует clean-guard и не допускает параллельных запусков.""" - text = DAG_PATH.read_text(encoding="utf-8") +def test_world_next_day_is_parameterless_paused_half_hour_dag(): + """DAG следующего дня запускается пустой формой каждые полчаса.""" + text = NEXT_DAY_DAG_PATH.read_text(encoding="utf-8") - assert "schedule=None" in text + assert 'dag_id="world_next_day"' in text + assert 'schedule="*/30 * * * *"' in text + assert "catchup=False" in text + assert "is_paused_upon_creation=True" in text assert "max_active_runs=1" in text - assert 'return "check_etl_not_paused_before_next_day"' in text - assert "assert_next_day_snapshot" in text - assert "assert_expected_t_end" in text - next_day_precheck = text.split("def precheck_next_day", maxsplit=1)[1].split( - "\ndef ", maxsplit=1 - )[0] - assert "assert_stand_clean" not in next_day_precheck - assert "assert_live_generator_not_running" in next_day_precheck - assert "run_next_day" in text + assert "params=" not in text + assert "precheck_next_day" in text + assert "run_next_day_task" in text + assert "TriggerDagRunOperator" in text + assert "wait_for_completion=True" in text + assert "check_manifest_task" in text + assert "check_etl_not_paused >> precheck" in text -def test_generator_control_prechecks_etl_dag_not_paused_before_waiting(): +def test_world_dags_precheck_etl_dag_not_paused_before_waiting(): """Пульт проверяет паузу etl_pipeline до долгого ожидания.""" text = DAG_PATH.read_text(encoding="utf-8") + shared_tasks = ( + REPO_ROOT / "airflow" / "dags" / "utils" / "startup_history_tasks.py" + ).read_text(encoding="utf-8") assert "assert_target_dag_not_paused" in text assert 'ETL_DAG_ID = "etl_pipeline"' in text assert 'op_kwargs={"dag_id": ETL_DAG_ID}' in text - assert "session.query(DagModel)" in text - assert "DagModel.dag_id == dag_id" in text - assert "target_dag_trigger_error" in text - assert "Airflow 2.10.5" in text - assert "fail_when_dag_is_paused" in text + assert "session.query(DagModel)" in shared_tasks + assert "DagModel.dag_id == dag_id" in shared_tasks + assert "target_dag_trigger_error" in shared_tasks + assert "Airflow 2.10.5" in shared_tasks + assert "fail_when_dag_is_paused" in shared_tasks assert 'return "check_etl_not_paused_before_backfill"' in text assert 'return "check_etl_not_paused_before_import"' in text - assert 'return "check_etl_not_paused_before_next_day"' in text assert "check_etl_not_paused_before_backfill >> precheck_backfill_task" in text assert "check_etl_not_paused_before_import >> precheck_import_task" in text - assert "check_etl_not_paused_before_next_day >> precheck_next_day_task" in text assert text.index("check_etl_not_paused_before_backfill >> precheck_backfill_task") < text.index( "precheck_backfill_task >> backfill_task" ) @@ -170,15 +171,16 @@ def test_make_up_rebuilds_airflow_images_after_repo_update(): assert "superset-init" not in up_block -def test_readme_quick_start_lists_prerequisites_and_runs_ddl_before_generator_control(): - """Быстрый старт называет инструменты и DDL до generator_control.""" +def test_readme_quick_start_creates_schema_via_ddl_init_before_world_init(): + """Быстрый старт: терминал — только make up, схему создаёт ddl_init до world_init.""" text = (REPO_ROOT / "README.md").read_text(encoding="utf-8") assert "uv" in text assert "Docker" in text assert "docker compose" in text - assert "make up\nmake ddl" in text - assert text.index("make ddl") < text.index("generator_control") + assert "make ddl" not in text + quick_start = text.split("## Быстрый старт", maxsplit=1)[1] + assert quick_start.index("ddl_init") < quick_start.index("world_init") def test_course_readme_lists_uv_before_first_command():