diff --git a/Makefile b/Makefile index 6f15f1d..80c7670 100644 --- a/Makefile +++ b/Makefile @@ -9,6 +9,8 @@ contract-test config-test test lint COMPOSE ?= docker compose +STARTUP_HISTORY_EXPORT_ARTIFACT ?= $(or $(ARTIFACT),/tmp/clickstream-startup-history.json) +STARTUP_HISTORY_IMPORT_ARTIFACT ?= $(or $(ARTIFACT),data/startup_history/reference-world.json.xz) PROFILE ?= daily-wave PYTEST_ARGS ?= -q -p no:cacheprovider PYTHON_CHECK_PATHS := airflow/dags configs generator scripts superset tests @@ -63,11 +65,11 @@ generated-history-chain-check: # Сгенерировать стартовую историю и сохранить портативный артефакт startup-history-export: - COMPOSE_BIN="$(COMPOSE)" bash ./scripts/export_startup_history_artifact.sh + ARTIFACT="$(STARTUP_HISTORY_EXPORT_ARTIFACT)" COMPOSE_BIN="$(COMPOSE)" bash ./scripts/export_startup_history_artifact.sh # Воспроизвести портативный артефакт стартовой истории в Kafka startup-history-import: - COMPOSE_BIN="$(COMPOSE)" bash ./scripts/import_startup_history_artifact.sh + ARTIFACT="$(STARTUP_HISTORY_IMPORT_ARTIFACT)" COMPOSE_BIN="$(COMPOSE)" bash ./scripts/import_startup_history_artifact.sh # Сверить ClickHouse после импорта с manifest артефакта startup-history-check: diff --git a/README.md b/README.md index 875eedb..e88b0d4 100644 --- a/README.md +++ b/README.md @@ -59,7 +59,7 @@ Airflow-образ и подтягивает новые зависимости make generated-history-analytics ``` -По умолчанию используется учебный профиль `daily-wave`: 2 суток с суточной +По умолчанию используется учебный профиль `daily-wave`: 3 суток с суточной волной. В live-продолжении он идёт с ×60: модельные сутки проходят примерно за 24 настенные минуты. Плоский профиль `ci` на 6 часов остаётся служебным для автоматических тестов. @@ -162,7 +162,7 @@ flowchart LR - [Запуск и эксплуатация](./docs/OPERATIONS.md) — сценарий запуска, параметры DAG-ов, мониторинг, частые проблемы. - [Runbook стартовой истории](./docs/runbooks/startup-history.md) — экспорт, - импорт и live-продолжение из готового артефакта. + импорт эталонного мира из Git по умолчанию и live-продолжение. - [Карта репозитория](./docs/REPO_MAP.md) — где какие файлы и что менять. - [Курс «Кликстрим на ClickHouse»](./docs/course/README.md) — учебная программа на этом стенде. diff --git a/airflow/dags/generator_control_dag.py b/airflow/dags/generator_control_dag.py index be66ed1..db5847f 100644 --- a/airflow/dags/generator_control_dag.py +++ b/airflow/dags/generator_control_dag.py @@ -252,7 +252,7 @@ with DAG( title="Артефакт", description=( "Backfill: куда сохранить файл; пусто — не сохранять. " - "Import: что читать; пусто — путь по умолчанию в data." + "Import: что читать; пусто — эталонный мир из репозитория." ), ), "expected_t_end": Param( diff --git a/data/startup_history/README.md b/data/startup_history/README.md new file mode 100644 index 0000000..bfd96dd --- /dev/null +++ b/data/startup_history/README.md @@ -0,0 +1,6 @@ +# Эталонный мир + +Здесь хранится `reference-world.json.xz` — готовая история за три модельных дня. +Файл собирает сопровождающий проекта по инструкции в +[`docs/runbooks/startup-history.md`](../../docs/runbooks/startup-history.md). +Ожидаемый размер сжатого файла — около 32 МБ. diff --git a/data/startup_history/reference-world.json.xz b/data/startup_history/reference-world.json.xz new file mode 100644 index 0000000..8db5276 Binary files /dev/null and b/data/startup_history/reference-world.json.xz differ diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index 480fc12..d44d18b 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -69,7 +69,8 @@ volumes или live-генератором. Для стыка backfill/live от - `profile` — список берётся из `PROFILES` генератора; - `duration` — `6h`, `2d` и т.п.; пусто означает длительность профиля; - `seed`, `model_time_speed` — необязательные переопределения мира; - - `artifact_path` — для `backfill` путь сохранения, для `import` путь чтения; + - `artifact_path` — для `backfill` путь сохранения; для `import` путь чтения. + При пустом поле импортируется эталонный мир из репозитория; - `expected_t_end` — необязательная ожидаемая граница перед `next-day`. При расхождении запуск показывает ожидаемое и фактическое значения. @@ -252,7 +253,7 @@ CHECK_LIVE_SEAM=1 GEN_LIVE_CHECK_MINUTES=10 make generated-history-check ``` Для commit gate issue 17 есть короткий runtime-путь без полного `daily-wave` на -2 суток и без Superset UI: +3 суток и без Superset UI: ```bash make generated-history-runtime-check @@ -270,7 +271,7 @@ make generated-history-runtime-check LIVE_SECONDS=45 WAIT_STG_SECONDS=10 make generated-history-runtime-check ``` -По умолчанию команда использует учебный профиль `daily-wave`: 2 суток с +По умолчанию команда использует учебный профиль `daily-wave`: 3 суток с суточной волной. В live-продолжении он идёт с ×60 и тикает раз в секунду, поэтому модельные сутки проходят примерно за 24 настенные минуты. Плоский профиль `ci` на 6 часов остаётся служебным для автоматических тестов. diff --git a/docs/TEST_PLAN.md b/docs/TEST_PLAN.md index df9bc8c..57d1e59 100644 --- a/docs/TEST_PLAN.md +++ b/docs/TEST_PLAN.md @@ -124,7 +124,7 @@ curl -s -u admin:admin "http://localhost:3000/api/dashboards/uid/airflow-overvie ### B.1 Полная стартовая история и ETL ```bash -# История на 2 суток с видимой суточной волной +# История на 3 суток с видимой суточной волной PROFILE=daily-wave make generated-history-analytics make up ``` diff --git a/docs/runbooks/startup-history.md b/docs/runbooks/startup-history.md index 9f2eae4..758ce7c 100644 --- a/docs/runbooks/startup-history.md +++ b/docs/runbooks/startup-history.md @@ -29,11 +29,43 @@ | Профиль | Для чего | Длительность | Live-ход | |---------|----------|--------------|----------| | `ci` | Быстрая проверка и CI | `6h` | ×1, тик 60 с | -| `daily-wave` | История с видимой суточной волной | `2d` | ×60, тик 1 с | +| `daily-wave` | История с видимой суточной волной | `3d` | ×60, тик 1 с | Длительность можно переопределить через `GEN_HISTORY_DURATION`, например `2d`. Команда сама считает `GEN_MODEL_T_END` от `GEN_MODEL_T0`. +## Эталонный мир + +В репозитории хранится готовый мир за три модельных дня: +`data/startup_history/reference-world.json.xz`. Это обычный JSON-артефакт, +сжатый xz. Его добавляют прямо в Git, без Git LFS. + +Сопровождающий проекта собирает файл одним запуском из корня репозитория: + +```bash +ARTIFACT="$PWD/data/startup_history/reference-world.json.xz" \ +make startup-history-export +``` + +Команда выполняет backfill с профилем `daily-wave` и сразу пишет xz-файл. +Несжатый JSON занимает около 850 МБ, сжатый файл — около 32 МБ. После проверки +файл можно добавить обычной командой: + +```bash +git add data/startup_history/reference-world.json.xz +``` + +При импорте DAG дважды читает и распаковывает артефакт: во время предпроверки и +перед записью в Kafka. Это увеличивает время импорта, но не меняет результат. + +Менти в форме `generator_control` выбирает `import` и оставляет +`artifact_path` пустым. Тогда читается эталонный мир из репозитория. Из консоли +тот же импорт запускается без указания пути: + +```bash +make startup-history-import +``` + ## Пульт в Airflow Основной ручной путь — DAG `generator_control` в Airflow UI: @@ -58,7 +90,8 @@ - `duration` можно оставить пустым, тогда берётся длительность профиля. - `seed` и `model_time_speed` — необязательные переопределения мира. - `artifact_path`: для `backfill` — куда сохранить файл; пусто — не сохранять. - Для `import` — что читать; пусто — `/opt/airflow/data/startup-history-import.json`. + Для `import` — что читать; пусто — эталонный мир из репозитория: + `/opt/airflow/data/startup_history/reference-world.json.xz`. Backfill и import работают только на чистом стенде. Если Kafka-топики данных или STG уже непустые, DAG упадёт до записи и подскажет `make clean`. Консольные @@ -79,7 +112,7 @@ STG уже непустые, DAG упадёт до записи и подска ## Экспорт -По умолчанию создаётся двухсуточный артефакт `daily-wave` с суточной волной: +По умолчанию создаётся трёхсуточный артефакт `daily-wave` с суточной волной: ```bash make startup-history-export @@ -104,11 +137,11 @@ make clean docker compose up -d clickhouse kafka make ddl -ARTIFACT=/tmp/clickstream-startup-history.json make startup-history-import +make startup-history-import sleep 10 make transform -ARTIFACT=/tmp/clickstream-startup-history.json make startup-history-check +ARTIFACT=data/startup_history/reference-world.json.xz make startup-history-check CHECK_LIVE_SEAM=0 make generated-history-check ``` diff --git a/docs/specs/2026-06-14-generator-model-time-and-startup-history.md b/docs/specs/2026-06-14-generator-model-time-and-startup-history.md index 9954e44..9e9a93b 100644 --- a/docs/specs/2026-06-14-generator-model-time-and-startup-history.md +++ b/docs/specs/2026-06-14-generator-model-time-and-startup-history.md @@ -105,7 +105,7 @@ ADR-0005 решил отвязать время генератора от реа - `reset` — подставляет `GEN_RUN_MODE=live` и `GEN_STATE_RESET=true`. Профиль `ci` даёт быстрый 6-часовой прогон и остаётся на ×1 с тиком 60 с. -Профиль `daily-wave` даёт 2 суток, чтобы была видна суточная волна, а в live +Профиль `daily-wave` даёт 3 суток, чтобы была видна суточная волна, а в live идёт с `GEN_MODEL_TIME_SPEED=60` и `GEN_TICK_SECONDS=1`: модельные сутки проходят примерно за 24 настенные минуты. Пара «тик 1 с, скорость ×60» выбрана, чтобы модельный шаг тика остался 60 секунд. Так событийный бюджет и форма волны diff --git a/docs/specs/2026-07-19-mentee-path-redesign.md b/docs/specs/2026-07-19-mentee-path-redesign.md index 895e5ea..78ba25c 100644 --- a/docs/specs/2026-07-19-mentee-path-redesign.md +++ b/docs/specs/2026-07-19-mentee-path-redesign.md @@ -60,8 +60,8 @@ 1. **Эталонный мир**: 3 модельных дня на профиле `daily-wave` (суточная волна видна, есть «средний» день и сравнение день-к-дню). Собирает мейнтейнер один раз через `backfill`; хранится в git по фиксированному - пути `data/startup_history/` сжатым `xz` (~13 МБ; замер 2026-07-19: - `xz -9e` жмёт артефакт в 48 раз). Формат не меняется: `import` требует + пути `data/startup_history/` сжатым `xz` (факт 2026-07-22: ~850 МБ + сырой JSON, ~32 МБ после `xz -9e`). Формат не меняется: `import` требует `raw_topics`, сжатие снимает вопрос размера. Паттерн-образец — `airflow-greenplum-solution` (сид `demo.sql.xz` в git, стрим-распаковка при загрузке). diff --git a/generator/README.md b/generator/README.md index 1d7a145..8e4f305 100644 --- a/generator/README.md +++ b/generator/README.md @@ -126,7 +126,7 @@ make generated-history-analytics проверяет Superset metadata. Файлы `data/*.jsonl` при этом не грузятся в Kafka: они пока используются только как фактура для генератора. -По умолчанию используется учебный профиль `daily-wave`: 2 суток с суточной +По умолчанию используется учебный профиль `daily-wave`: 3 суток с суточной волной. В live он идёт с `GEN_MODEL_TIME_SPEED=60` и `GEN_TICK_SECONDS=1`: модельные сутки проходят примерно за 24 настенные минуты. Плоский профиль `ci` на 6 часов остаётся служебным для автоматических diff --git a/generator/src/clickstream_generator/airflow_control.py b/generator/src/clickstream_generator/airflow_control.py index eacc21a..75ca623 100644 --- a/generator/src/clickstream_generator/airflow_control.py +++ b/generator/src/clickstream_generator/airflow_control.py @@ -335,12 +335,13 @@ def target_dag_trigger_error(dag_id: str, dag_model) -> str | None: def default_artifact_path(operation: str) -> str: """Возвращает путь артефакта по умолчанию в общем томе data.""" - filename = ( - "startup-history-import.json" - if operation == "import" - else "startup-history.json" - ) - return str(Path(AIRFLOW_DATA_DIR) / filename) + if operation == "import": + return str( + Path(AIRFLOW_DATA_DIR) + / "startup_history" + / "reference-world.json.xz" + ) + return str(Path(AIRFLOW_DATA_DIR) / "startup-history.json") @contextmanager diff --git a/generator/src/clickstream_generator/launch.py b/generator/src/clickstream_generator/launch.py index fb1328f..4ba62db 100644 --- a/generator/src/clickstream_generator/launch.py +++ b/generator/src/clickstream_generator/launch.py @@ -47,7 +47,7 @@ PROFILES = { }, ), "daily-wave": LaunchProfile( - duration="2d", + duration="3d", env={ "GEN_SEED": "4242", "GEN_MODEL_T0": "2026-01-01T00:00:00+00:00", diff --git a/generator/src/clickstream_generator/startup_history_artifact.py b/generator/src/clickstream_generator/startup_history_artifact.py index e77f9e6..3ae2483 100644 --- a/generator/src/clickstream_generator/startup_history_artifact.py +++ b/generator/src/clickstream_generator/startup_history_artifact.py @@ -5,6 +5,7 @@ from __future__ import annotations import argparse import hashlib import json +import lzma import logging from copy import deepcopy from datetime import datetime, timezone @@ -353,6 +354,15 @@ def write_startup_history_artifact(path: str | Path, artifact: dict) -> None: """Пишет артефакт в JSON-файл.""" target = Path(path) target.parent.mkdir(parents=True, exist_ok=True) + if target.suffix == ".xz": + with lzma.open( + target, + "wt", + encoding="utf-8", + preset=9 | lzma.PRESET_EXTREME, + ) as output: + json.dump(artifact, output, ensure_ascii=True, indent=2) + return target.write_text( json.dumps(artifact, ensure_ascii=True, indent=2), encoding="utf-8", @@ -361,7 +371,11 @@ def write_startup_history_artifact(path: str | Path, artifact: dict) -> None: def load_startup_history_artifact(path: str | Path) -> dict: """Читает артефакт из JSON-файла.""" - return json.loads(Path(path).read_text(encoding="utf-8")) + source = Path(path) + if source.suffix == ".xz": + with lzma.open(source, "rt", encoding="utf-8") as input_file: + return json.load(input_file) + return json.loads(source.read_text(encoding="utf-8")) def validate_startup_history_artifact( diff --git a/generator/tests/test_airflow_control.py b/generator/tests/test_airflow_control.py index 931f926..a4806b9 100644 --- a/generator/tests/test_airflow_control.py +++ b/generator/tests/test_airflow_control.py @@ -10,6 +10,18 @@ from types import SimpleNamespace import pytest +def test_default_artifact_paths_keep_import_and_backfill_separate(): + """Import читает эталонный мир, а backfill пишет в прежний файл.""" + from clickstream_generator.airflow_control import default_artifact_path + + assert default_artifact_path("import") == ( + "/opt/airflow/data/startup_history/reference-world.json.xz" + ) + assert default_artifact_path("backfill") == ( + "/opt/airflow/data/startup-history.json" + ) + + def test_build_control_env_uses_profile_duration_and_airflow_data_dir(): """Пульт строит env для backfill из профиля и каталога Airflow.""" from clickstream_generator.airflow_control import build_control_env @@ -47,7 +59,7 @@ def test_import_env_requires_artifact_path_and_uses_backfill_contract(): assert env["GEN_RUN_MODE"] == "backfill" assert env["GEN_STATE_RESET"] == "true" - assert env["GEN_HISTORY_DURATION"] == "2d" + assert env["GEN_HISTORY_DURATION"] == "3d" def test_next_day_env_uses_world_settings_and_current_manifest_boundary(): diff --git a/generator/tests/test_generator_control_dag_contract.py b/generator/tests/test_generator_control_dag_contract.py index cd8517a..e7dbc0d 100644 --- a/generator/tests/test_generator_control_dag_contract.py +++ b/generator/tests/test_generator_control_dag_contract.py @@ -247,6 +247,22 @@ def test_console_generator_paths_precheck_clean_stand_before_writes(): ) +def test_reference_artifact_is_import_default_only(): + """Экспорт не может по умолчанию перезаписать эталонный мир.""" + makefile = (REPO_ROOT / "Makefile").read_text(encoding="utf-8") + export_script = ( + REPO_ROOT / "scripts" / "export_startup_history_artifact.sh" + ).read_text(encoding="utf-8") + import_script = ( + REPO_ROOT / "scripts" / "import_startup_history_artifact.sh" + ).read_text(encoding="utf-8") + + assert "STARTUP_HISTORY_EXPORT_ARTIFACT" in makefile + assert "STARTUP_HISTORY_IMPORT_ARTIFACT" in makefile + assert "/tmp/clickstream-startup-history.json" in export_script + assert "data/startup_history/reference-world.json.xz" in import_script + + def test_console_continue_checks_state_or_clean_stand_before_live_start(): """Continue не стартует новый live-мир на грязном стенде без state.""" run_generator = (REPO_ROOT / "scripts" / "run_generator.sh").read_text( diff --git a/generator/tests/test_launch.py b/generator/tests/test_launch.py index 6b96b71..10ca177 100644 --- a/generator/tests/test_launch.py +++ b/generator/tests/test_launch.py @@ -26,15 +26,15 @@ def test_duration_2d_sets_model_t_end_from_t0(): assert env["GEN_MODEL_T_END"] == "2026-01-03T00:00:00+00:00" -def test_daily_wave_profile_uses_two_days_by_default(): +def test_daily_wave_profile_uses_three_days_by_default(): """Профиль daily-wave даёт историю с суточной волной.""" from clickstream_generator.launch import build_launch_env env = build_launch_env("backfill", profile_name="daily-wave") assert env["GEN_LAUNCH_PROFILE"] == "daily-wave" - assert env["GEN_HISTORY_DURATION"] == "2d" - assert env["GEN_MODEL_T_END"] == "2026-01-03T00:00:00+00:00" + assert env["GEN_HISTORY_DURATION"] == "3d" + assert env["GEN_MODEL_T_END"] == "2026-01-04T00:00:00+00:00" assert env["GEN_MODEL_TIME_SPEED"] == "60" assert env["GEN_TICK_SECONDS"] == "1" diff --git a/generator/tests/test_startup_history_artifact.py b/generator/tests/test_startup_history_artifact.py index c90d511..d70c2d6 100644 --- a/generator/tests/test_startup_history_artifact.py +++ b/generator/tests/test_startup_history_artifact.py @@ -4,12 +4,55 @@ from dataclasses import replace from datetime import datetime, timezone +import json +import lzma import pytest from generator import GeneratorState +def test_plain_json_roundtrip_keeps_existing_serialization(tmp_path): + """Обычный JSON по-прежнему пишется без сжатия и читается обратно.""" + from clickstream_generator.startup_history_artifact import ( + load_startup_history_artifact, + write_startup_history_artifact, + ) + + artifact = {"text": "мир", "nested": {"value": 42}} + path = tmp_path / "artifact.json" + + write_startup_history_artifact(path, artifact) + + assert path.read_text(encoding="utf-8") == json.dumps( + artifact, + ensure_ascii=True, + indent=2, + ) + assert load_startup_history_artifact(path) == artifact + + +def test_xz_artifact_roundtrip(tmp_path): + """XZ-артефакт потоково пишется и читается без изменения данных.""" + from clickstream_generator.startup_history_artifact import ( + load_startup_history_artifact, + write_startup_history_artifact, + ) + + artifact = {"text": "мир", "nested": {"value": 42}} + path = tmp_path / "artifact.json.xz" + + write_startup_history_artifact(path, artifact) + + assert load_startup_history_artifact(path) == artifact + with lzma.open(path, "rt", encoding="utf-8") as input_file: + assert input_file.read() == json.dumps( + artifact, + ensure_ascii=True, + indent=2, + ) + + def _state() -> GeneratorState: import random diff --git a/scripts/assert_stand_clean.sh b/scripts/assert_stand_clean.sh index a3872ab..8ed54a3 100755 --- a/scripts/assert_stand_clean.sh +++ b/scripts/assert_stand_clean.sh @@ -26,7 +26,8 @@ fi if [[ "${CHECK_MODE}" == "clean" ]]; then echo "Предпроверка: live-генератор не должен быть запущен." PYTHONPATH="${REPO_ROOT}/generator/src" \ - uv run python -m clickstream_generator.stand_clean \ + uv run --with-requirements generator/requirements.txt \ + python -m clickstream_generator.stand_clean \ --mode live \ --live-metrics-url "${GENERATOR_METRICS_URL_HOST}" fi @@ -54,7 +55,8 @@ fi KAFKA_BOOTSTRAP_SERVERS="${KAFKA_BOOTSTRAP_SERVERS_HOST}" \ GEN_DATA_DIR="${REPO_ROOT}/data" \ PYTHONPATH="${REPO_ROOT}/generator/src" \ - uv run python -m clickstream_generator.stand_clean \ + uv run --with-requirements generator/requirements.txt \ + python -m clickstream_generator.stand_clean \ --mode "${CHECK_MODE}" \ --kafka-bootstrap-servers "${KAFKA_BOOTSTRAP_SERVERS_HOST}" \ --stg-counts "${stg_counts}" diff --git a/scripts/check_startup_history_manifest.sh b/scripts/check_startup_history_manifest.sh index 06b9ef0..9d0f444 100755 --- a/scripts/check_startup_history_manifest.sh +++ b/scripts/check_startup_history_manifest.sh @@ -27,7 +27,7 @@ clickhouse_datetime_literal() { [[ -s "${ARTIFACT}" ]] || fail "артефакт не найден или пуст: ${ARTIFACT}" -ARTIFACT_DIR="$(dirname "${ARTIFACT}")" +ARTIFACT_DIR="$(cd "$(dirname "${ARTIFACT}")" && pwd)" ARTIFACT_FILE="$(basename "${ARTIFACT}")" manifest_summary="$(${COMPOSE_BIN} run --rm --no-deps \ diff --git a/scripts/import_startup_history_artifact.sh b/scripts/import_startup_history_artifact.sh index 8020699..8e91f9e 100755 --- a/scripts/import_startup_history_artifact.sh +++ b/scripts/import_startup_history_artifact.sh @@ -8,7 +8,7 @@ SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" REPO_ROOT="$(cd "${SCRIPT_DIR}/.." && pwd)" COMPOSE_BIN="${COMPOSE_BIN:-docker compose}" -ARTIFACT="${ARTIFACT:-/tmp/clickstream-startup-history.json}" +ARTIFACT="${ARTIFACT:-data/startup_history/reference-world.json.xz}" PROFILE="${PROFILE:-${GEN_LAUNCH_PROFILE:-daily-wave}}" duration_args=() if [[ -n "${GEN_HISTORY_DURATION:-}" ]]; then @@ -31,7 +31,7 @@ if [[ ! -s "${ARTIFACT}" ]]; then exit 1 fi -ARTIFACT_DIR="$(dirname "${ARTIFACT}")" +ARTIFACT_DIR="$(cd "$(dirname "${ARTIFACT}")" && pwd)" ARTIFACT_FILE="$(basename "${ARTIFACT}")" echo "=== Импорт стартовой истории ==="