diff --git a/Makefile b/Makefile index e3b17f1..d07808a 100644 --- a/Makefile +++ b/Makefile @@ -14,7 +14,7 @@ COMPOSE ?= docker compose # ============================================================================ up: - $(COMPOSE) up -d + $(COMPOSE) up -d --build clickhouse kafka kafka-ui postgres-metadata airflow-webserver airflow-scheduler prometheus grafana kafka-exporter statsd-exporter # Остановить и удалить контейнеры/сети текущего проекта down: diff --git a/README.md b/README.md index 16253e7..070534f 100644 --- a/README.md +++ b/README.md @@ -28,10 +28,14 @@ ## Быстрый старт +Перед первой командой нужны `Docker` с `docker compose`, `make`, `bash`, `curl`, +`git` и `uv`. `uv` нужен для локальных Python-проверок и команд разработки. + Для ручной работы поднимите стенд и создайте стартовую историю через Airflow: ```bash make up +make ddl docker compose ps ``` @@ -40,6 +44,11 @@ docker compose ps запустит ETL и выполнит `check`. `make up` не запускает live-генератор; live включается отдельно командой `make generator-continue`. +После обновления репозитория снова выполните `make up`: команда пересобирает +Airflow-образ и подтягивает новые зависимости и DAG-и. Superset-дэшборд +собирается позже, когда DM уже готов: через `make generated-history-analytics` +или `make superset-init`. + Для полностью автоматического чистого прогона из консоли остаётся команда: ```bash diff --git a/airflow/dags/generator_control_dag.py b/airflow/dags/generator_control_dag.py index c773ffc..d24ec60 100644 --- a/airflow/dags/generator_control_dag.py +++ b/airflow/dags/generator_control_dag.py @@ -10,10 +10,13 @@ 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 @@ -41,6 +44,8 @@ default_args = { "retry_delay": timedelta(minutes=1), } +ETL_DAG_ID = "etl_pipeline" + def _conf(context) -> dict: dag_run = context.get("dag_run") @@ -92,9 +97,9 @@ def choose_operation(**context) -> str: """Выбирает ветку пульта по параметру operation.""" operation = _operation(context) if operation == "backfill": - return "precheck_backfill" + return "check_etl_not_paused_before_backfill" if operation == "import": - return "precheck_import" + return "check_etl_not_paused_before_import" if operation == "check": return "check_only" raise ValueError(f"Неизвестная операция: {operation}") @@ -144,6 +149,27 @@ def check_manifest_task(**context) -> None: 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() + if dag_model is None: + raise AirflowException( + f"DAG {dag_id} ещё не найден Airflow. Подождите парсинга DAG-файлов " + "или перезапустите Airflow через make up." + ) + if dag_model.is_paused: + raise AirflowException( + f"DAG {dag_id} стоит на паузе: снимите паузу в Airflow UI, " + "затем повторите generator_control." + ) + + with DAG( dag_id="generator_control", description="Пульт стартовой истории генератора", @@ -224,9 +250,20 @@ with DAG( python_callable=run_import_task, ) + check_etl_not_paused_before_backfill = PythonOperator( + task_id="check_etl_not_paused_before_backfill", + python_callable=assert_target_dag_not_paused, + op_kwargs={"dag_id": ETL_DAG_ID}, + ) + check_etl_not_paused_before_import = PythonOperator( + task_id="check_etl_not_paused_before_import", + 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_pipeline", + trigger_dag_id=ETL_DAG_ID, conf={"full_refresh": True}, wait_for_completion=True, allowed_states=["success"], @@ -250,7 +287,13 @@ with DAG( trigger_rule=TriggerRule.NONE_FAILED_MIN_ONE_SUCCESS, ) - route >> [precheck_backfill_task, precheck_import_task, check_only] + route >> [ + check_etl_not_paused_before_backfill, + check_etl_not_paused_before_import, + check_only, + ] + check_etl_not_paused_before_backfill >> precheck_backfill_task + check_etl_not_paused_before_import >> precheck_import_task precheck_backfill_task >> backfill_task >> trigger_etl precheck_import_task >> import_task >> trigger_etl trigger_etl >> check_after_etl >> done diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index ffb0fbe..11ef263 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -6,7 +6,7 @@ Базовые команды: -- `make up` (или `docker compose up -d`) +- `make up` (поднимает базовый стенд без Superset init; пересобирает Airflow-образ) - `make down` (остановить и удалить контейнеры/сети проекта) - `make clean` (полная очистка: `down -v --remove-orphans`) - `make generated-history-analytics` (штатный чистый прогон: стартовая история @@ -63,6 +63,17 @@ Backfill/import требуют чистый стенд: пустые data-топ в DAG нет: live-генератор — долгоживущий сервис, его запускают с консоли через `make generator-continue`. +Перед запуском `etl_pipeline` пульт проверяет, что DAG не стоит на паузе. Если +стоит, задача падает сразу с подсказкой снять паузу в UI или командой: + +```bash +docker compose exec -T airflow-webserver airflow dags unpause etl_pipeline +``` + +После обновления репозитория выполните `make up`, чтобы Airflow получил новые +зависимости и DAG-и через пересборку образа. Superset metadata собирается после +готового DM: через `make generated-history-analytics` или `make superset-init`. + ### `ddl_init` - Запуск: ручной (`Trigger DAG`) diff --git a/docs/course/README.md b/docs/course/README.md index ff1b11c..2d1121c 100644 --- a/docs/course/README.md +++ b/docs/course/README.md @@ -25,6 +25,8 @@ слова с живым стендом. - **Железо:** стек тяжёлый — Kafka, ClickHouse, Airflow, Superset, Prometheus и Grafana поднимаются одновременно. Нужна машина, которая это потянет. +- **Инструменты:** `Docker` с `docker compose`, `make`, `bash`, `curl`, `git` и `uv`. + `uv` нужен для локальных Python-проверок и команд разработки. - **Подними стенд и создай стартовую историю** (из корня репозитория) — этого хватит, чтобы начать, и прогон быстрый: diff --git a/docs/runbooks/startup-history.md b/docs/runbooks/startup-history.md index bb041b4..8ac8d29 100644 --- a/docs/runbooks/startup-history.md +++ b/docs/runbooks/startup-history.md @@ -142,6 +142,11 @@ ClickHouse уже успел прочитать частичные сообще PROFILE=daily-wave make generator-continue ``` +У `daily-wave` скорость ×60. Долгий простой стенда создаёт большую дыру в +модельном времени: ночь простоя может стать десятками модельных суток без +событий. Для чистой демонстрации лучше запустите `make generator-reset` или +повторите импорт стартовой истории через `make startup-history-import`. + Если читаемый state есть, но настройки не совпадают, генератор падает с перечнем полей. Это защита от смешения разных миров. Для намеренного нового мира используйте `make generator-reset` или `make clean`. diff --git a/generator/tests/test_generator_control_dag_contract.py b/generator/tests/test_generator_control_dag_contract.py index 8cb350d..c05c3cb 100644 --- a/generator/tests/test_generator_control_dag_contract.py +++ b/generator/tests/test_generator_control_dag_contract.py @@ -42,12 +42,83 @@ def test_dag_branches_and_waits_for_etl_completion(): assert "BranchPythonOperator" in text assert "TriggerDagRunOperator" in text - assert 'trigger_dag_id="etl_pipeline"' in text + assert "trigger_dag_id=ETL_DAG_ID" in text assert "wait_for_completion=True" in text assert 'allowed_states=["success"]' in text assert 'failed_states=["failed"]' in text +def test_generator_control_prechecks_etl_dag_not_paused_before_waiting(): + """Пульт проверяет паузу etl_pipeline до долгого ожидания.""" + text = DAG_PATH.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 "Airflow 2.10.5" in text + assert "fail_when_dag_is_paused" in text + assert "снимите паузу" in text + assert 'return "check_etl_not_paused_before_backfill"' in text + assert 'return "check_etl_not_paused_before_import"' 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 text.index("check_etl_not_paused_before_backfill >> precheck_backfill_task") < text.index( + "precheck_backfill_task >> backfill_task" + ) + assert text.index("check_etl_not_paused_before_import >> precheck_import_task") < text.index( + "precheck_import_task >> import_task" + ) + + +def test_make_up_rebuilds_airflow_images_after_repo_update(): + """make up пересобирает Airflow, но не запускает Superset до DM.""" + text = (REPO_ROOT / "Makefile").read_text(encoding="utf-8") + + up_block = text.split("\nup:\n", maxsplit=1)[1].split( + "\n\n# Остановить", + maxsplit=1, + )[0] + assert "$(COMPOSE) up -d --build" in up_block + assert "airflow-webserver" in up_block + assert "airflow-scheduler" in up_block + assert "superset" not in up_block + assert "superset-init" not in up_block + + +def test_readme_quick_start_lists_prerequisites_and_runs_ddl_before_generator_control(): + """Быстрый старт называет инструменты и DDL до generator_control.""" + 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") + + +def test_course_readme_lists_uv_before_first_command(): + """Курс называет uv до первой команды.""" + text = (REPO_ROOT / "docs" / "course" / "README.md").read_text(encoding="utf-8") + + assert "uv" in text + assert text.index("uv") < text.index("make generated-history-analytics") + + +def test_startup_history_runbook_warns_about_daily_wave_idle_gap(): + """Runbook предупреждает о дыре модельного времени при простое daily-wave.""" + text = (REPO_ROOT / "docs" / "runbooks" / "startup-history.md").read_text( + encoding="utf-8" + ) + + assert "daily-wave" in text + assert "прост" in text + assert "дыр" in text + assert "make generator-reset" in text + assert "startup-history-import" in text + + def test_no_docker_or_continue_operation_in_dag(): """Пульт не управляет Docker и не содержит операцию continue.""" text = DAG_PATH.read_text(encoding="utf-8")