diff --git a/.env.example b/.env.example index 55b193b..bd4dd69 100644 --- a/.env.example +++ b/.env.example @@ -15,6 +15,11 @@ GRAFANA_PORT=23000 AIRFLOW_PORT=28080 SUPERSET_PORT=28088 +# GID группы `docker` на этой машине: без него планировщик Airflow не достучится +# до сокета докера, и пульт мира не поднимет контейнер генератора. +# Подсмотреть свой — `getent group docker`. +DOCKER_GID=127 + # Учётные данные и ключи учебного стенда. Для VPS замените значения образца. GRAFANA_ADMIN_USER=admin GRAFANA_ADMIN_PASSWORD=admin diff --git a/CONTEXT.md b/CONTEXT.md index 4288cf2..454aeee 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -114,7 +114,8 @@ _Избегать_: манифест, мини-манифест **Живой день**: Проигрывание текущего модельного дня в реальном времени с ускорением; -включается по требованию, не постоянный фон. +включается по требованию — постоянным фоном идёт, только пока включён +выключатель пульта. **Пакетный режим**: Проигрывание готового дня пачкой, без темпа: заливка снимка при старте diff --git a/README.md b/README.md index fe91b6b..ad7ce26 100644 --- a/README.md +++ b/README.md @@ -50,7 +50,10 @@ make check-services Дальше, в рабочей петле, обычно хватает `make smoke`. `.env.example` — образец обязательной локальной настройки. В нём живут имя -экземпляра, внешние порты, учебные учётные данные и ключи. `compose.yaml` не +экземпляра, внешние порты, учебные учётные данные и ключи. Одно значение образца +верно не везде — `DOCKER_GID`: у группы `docker` на каждой машине свой номер, и +свой подскажет `getent group docker`. Пока он не сойдётся, стенд работает весь, +кроме пульта мира — о нём ниже. `compose.yaml` не дублирует их умолчаниями: если значения нет в `.env`, Compose сразу предложит скопировать образец. Внутренние адреса, имена топиков и версии образов описывают сам стенд, поэтому записаны литералами в `compose.yaml` и Dockerfile. @@ -219,6 +222,33 @@ uv run --project generator python -m clickstream_generator batch \ код возврата, ещё до первого события. Поэтому умолчание дня и живёт в `Makefile`: день называет тот, кто запускает, а не тот, кого запускают. +### Пульт мира + +Считать дни самому необязательно. В Airflow живёт **пульт мира** — три дага, +которые зовут тот же образ генератора и ведут позицию на оси за вас. + +| Даг | Что делает | +| --- | --- | +| `world_next_day` | играет следующие дни пачкой; сколько — параметр запуска, по умолчанию один | +| `world_live_day` | играет следующий день в темпе модельного времени, около двадцати четырёх минут | +| `world_live` | выключатель: пока снят с паузы, дёргает `world_live_day` день за днём | + +Позицию хранит переменная Airflow `world_position` — номер первого несыгранного +дня. Ставит её работник, сыгравший день, и только по успеху: оборванный прогон +позицию не двигает, и следующий запуск играет тот же день заново. Нет +переменной — мир в стартовом состоянии. + +Выключатель создаётся на паузе. Снимите — мир поедет сам; поставите обратно — +встанет на границе модельных суток, доиграв начатый день. Форма пульта и доводы +целиком — [ADR 0009](docs/adr/0009-world-control.md). + +**Плата названа вслух: планировщику Airflow отдан сокет докера** — иначе +контейнер генератора ему не поднять. Доступ к сокету равен праву root на +машине; для локального учебного стенда размен принят, но знать о нём надо. +Открывает дверь не монтирование, а членство в группе: GID группы `docker` у +каждой машины свой, живёт в `.env` и подсматривается командой +`getent group docker`. + ## Состав и доступ - `clickhouse-01` — инициатор DDL и точка подключения Airflow; diff --git a/compose.yaml b/compose.yaml index 10e79f2..916aee4 100644 --- a/compose.yaml +++ b/compose.yaml @@ -29,6 +29,17 @@ x-clickhouse-common: &clickhouse-common retries: 30 start_period: 10s +# Факты стенда, у которых стало по два потребителя: разовая служба генератора и +# пульт мира. Названы по одному разу — иначе однажды разойдутся, и заметит это +# не проверка, а менти с пустым топиком. +x-generator-image: &generator-image clickstream-generator:local + +x-kafka-target: &kafka-target + KAFKA_BOOTSTRAP_SERVERS: kafka:9092 + KAFKA_TOPIC: hits + +x-world-starting-days: &world-starting-days "8" + x-airflow-common: &airflow-common image: clickstream-airflow:local build: @@ -50,6 +61,13 @@ x-airflow-common: &airflow-common AIRFLOW_ADMIN_USER: ${AIRFLOW_ADMIN_USER:?Скопируйте .env.example в .env} AIRFLOW_ADMIN_PASSWORD: ${AIRFLOW_ADMIN_PASSWORD:?Скопируйте .env.example в .env} CLICKHOUSE_ETL_PASSWORD: ${CLICKHOUSE_ETL_PASSWORD:?Скопируйте .env.example в .env} + # Пульт мира поднимает генератор сам, отдельным контейнером, — значит те же + # факты стенда, что compose даёт разовой службе генератора, нужны и дагам. + # Имя сети собирается из имени проекта: у второй копии стенда оно другое. + <<: *kafka-target + GENERATOR_IMAGE: *generator-image + STAND_NETWORK: ${COMPOSE_PROJECT_NAME}_default + WORLD_STARTING_DAYS: *world-starting-days volumes: - ./dags:/opt/airflow/dags:ro - ./infra/airflow/init.sh:/opt/airflow/init.sh:ro @@ -57,7 +75,7 @@ x-airflow-common: &airflow-common - airflow_auth:/opt/airflow/auth x-generator-common: &generator-common - image: clickstream-generator:local + image: *generator-image build: context: . dockerfile: generator/Dockerfile @@ -75,8 +93,7 @@ x-generator-common: &generator-common # Адрес брокера и имя топика — факты стенда, и называет их стенд. # Остальное (день, зерно, число дней, предел пачки) приходит аргументами # от того, кто запускает: у службы нет позиции на оси мира. - KAFKA_BOOTSTRAP_SERVERS: kafka:9092 - KAFKA_TOPIC: hits + <<: *kafka-target x-superset-common: &superset-common image: clickstream-superset:local @@ -261,14 +278,15 @@ services: # Число дней стоит здесь числом: YAML не читает Python, и одно из двух мест # (второе — `STARTING_DAYS` в inventory.py) лишнее по построению. Правя одно, # правьте второе — на страже тут никто не стоит: залей эта служба лишний - # день, он лёг бы за рамкой дат описи и остался бы незамеченным. + # день, он лёг бы за рамкой дат описи и остался бы незамеченным. Внутри YAML + # число одно на всех: то же говорит дагам пульта, где кончается стартовый мир. # # Повторный `make up` заливает мир заново, и это не оплошность: `WatchID` у # событий те же, ReplacingMergeTree схлопнет повтор в ODS. Сырьё в STG при # этом честно удвоится — свойство слоя, описанное в storage.md. world-init: <<: *generator-common - command: ["batch", "--day", "0", "--days", "8"] + command: ["batch", "--day", "0", "--days", *world-starting-days] # Тот же образ для ручных прогонов: `make generate-batch`, `make generate-live`. # Под профилем — чтобы обычный подъём стенда её не трогал. @@ -350,6 +368,21 @@ services: airflow-init: condition: service_completed_successfully command: scheduler + # Пульт мира зовёт генератор отдельным контейнером, а при LocalExecutor + # задачи исполняет сам планировщик — значит сокет докера нужен ему одному. + # Плата названа вслух в README и ADR 0009: доступ к сокету равен праву root + # на машине. Открывает дверь не монтирование, а группа: у сокета права 660 + # и группа `docker`, чей GID на каждой машине свой и живёт в `.env`. + group_add: + - ${DOCKER_GID:?Скопируйте .env.example в .env} + # Тома перечислены заново: список службы общий не дополняет, а заменяет. + # Появится новый том у остальных служб Airflow — добавьте и сюда. + volumes: + - /var/run/docker.sock:/var/run/docker.sock + - ./dags:/opt/airflow/dags:ro + - ./infra/airflow/init.sh:/opt/airflow/init.sh:ro + - airflow_logs:/opt/airflow/logs + - airflow_auth:/opt/airflow/auth mem_limit: 640m healthcheck: test: ["CMD-SHELL", "curl -sf http://127.0.0.1:8974/health >/dev/null"] diff --git a/dags/world_control.py b/dags/world_control.py new file mode 100644 index 0000000..fcad3a4 --- /dev/null +++ b/dags/world_control.py @@ -0,0 +1,183 @@ +"""Пульт мира: даги, которыми двигают ось модельного времени. + +Их три, и они двух родов. **Работники** играют день и запускаются руками: +`world_next_day` — пачкой, без пауз, сколько дней попросили; `world_live_day` — +один день в темпе модельного времени. **Выключатель** `world_live` своей работы +не делает: он тикает по расписанию и дёргает работника живого дня, дожидаясь +конца. Снят с паузы — мир едет день за днём; поставлен на паузу — встал на +границе модельных суток. + +Разделение не косметическое. Расписание на самом работнике заставило бы кнопку +паузы значить две вещи разом — «мир не едет сам» и «даг выключен», — а работник +при этом выглядел бы в списке выключенным, хотя нажимают его каждый день. + +Календарь Airflow к оси мира отношения не имеет: какой день играть, работник +спрашивает у переменной, а не у логической даты прогона. + +Решения и доводы целиком — ADR 0009. +""" + +from __future__ import annotations + +import datetime +import os + +from airflow.providers.docker.operators.docker import DockerOperator +from airflow.providers.standard.operators.trigger_dagrun import TriggerDagRunOperator +from airflow.sdk import Param, Variable, dag, get_current_context, task + +# Факты стенда — образ генератора, сеть, адрес брокера, топик и размер +# стартового мира — приходят окружением, и называет их compose: тот же, что +# называет их разовой службе генератора. +GENERATOR_IMAGE = os.environ["GENERATOR_IMAGE"] +STAND_NETWORK = os.environ["STAND_NETWORK"] +GENERATOR_ENVIRONMENT = { + "KAFKA_BOOTSTRAP_SERVERS": os.environ["KAFKA_BOOTSTRAP_SERVERS"], + "KAFKA_TOPIC": os.environ["KAFKA_TOPIC"], +} +STARTING_DAYS = int(os.environ["WORLD_STARTING_DAYS"]) + +# Позиция на оси: номер первого несыгранного дня. Переменной нет — мир в +# стартовом состоянии, и играть надо сразу за ним. +# +# Позиция ставится, а не увеличивается. Наложись один прогон на другой, худшее +# при таком правиле — сыгранный дважды день, а повтор схлопнет +# ReplacingMergeTree. Увеличение молча съело бы день, и в мире осталась бы +# дыра, которой никто не заметит. +WORLD_POSITION = "world_position" + +# Тик выключателя. Каденцию задаёт не он, а сама длина живого дня — около +# двадцати четырёх минут: тик только спрашивает «не пора ли снова». +LIVE_TICK = datetime.timedelta(minutes=25) + +START_DATE = datetime.datetime(2026, 1, 1, tzinfo=datetime.UTC) +TAGS = ["пульт мира"] + + +@task +def first_unplayed_day() -> int: + """Номер дня, с которого играть.""" + return int(Variable.get(WORLD_POSITION, default=STARTING_DAYS)) + + +def _play(task_id: str, command: list[str]) -> DockerOperator: + """Задача, играющая дни в каноническом контейнере генератора. + + Внутрь образа Airflow генератор не поставить: он требует Python 3.14, а + образ несёт 3.13. Да и обещание побайтовой воспроизводимости дано для + зафиксированного образа генератора — держится оно только там. + """ + return DockerOperator( + task_id=task_id, + image=GENERATOR_IMAGE, + command=command, + network_mode=STAND_NETWORK, + environment=GENERATOR_ENVIRONMENT, + # Контейнер убирается в любом исходе: вывод генератора оператор уже + # перелил в журнал задачи, а мёртвые контейнеры копить незачем. + auto_remove="force", + # По умолчанию оператор монтирует контейнеру временный каталог. Здесь + # это ловушка: путь он заводит внутри Airflow, а монтирует демон с + # хоста, где такого пути нет. Генератору временный каталог не нужен. + mount_tmp_dir=False, + ) + + +@dag( + dag_id="world_next_day", + schedule=None, + start_date=START_DATE, + is_paused_upon_creation=False, + max_active_runs=1, + tags=TAGS, + params={"days": Param(1, type="integer", minimum=1, title="Сколько дней прожить")}, +) +def world_next_day(): + """Прожить следующие дни пачкой, без пауз. + + Запускается руками. День по умолчанию один, но разгон вперёд идёт одним + нажимом, а не десятью: сколько дней играть — параметр запуска. + """ + + @task + def remember_played(first_day: int) -> None: + """Позиция ставится по сыгранным дням и только по успеху.""" + days = get_current_context()["params"]["days"] + Variable.set(WORLD_POSITION, str(first_day + days)) + + first_day = first_unplayed_day() + played = _play( + "play_days", + [ + "batch", + "--day", + "{{ ti.xcom_pull(task_ids='first_unplayed_day') }}", + "--days", + "{{ params.days }}", + ], + ) + + first_day >> played >> remember_played(first_day) + + +@dag( + dag_id="world_live_day", + schedule=None, + start_date=START_DATE, + is_paused_upon_creation=False, + max_active_runs=1, + tags=TAGS, +) +def world_live_day(): + """Прожить следующий день в темпе модельного времени. + + Ускорение ×60: модельные сутки укладываются примерно в двадцать четыре + реальные минуты, и суточная волна разворачивается на глазах. Запускается + руками; чтобы мир жил так день за днём сам, есть выключатель `world_live`. + """ + + @task + def remember_played(first_day: int) -> None: + """Позиция ставится по сыгранному дню и только по успеху.""" + Variable.set(WORLD_POSITION, str(first_day + 1)) + + first_day = first_unplayed_day() + played = _play( + "play_day", + ["live", "--day", "{{ ti.xcom_pull(task_ids='first_unplayed_day') }}"], + ) + + first_day >> played >> remember_played(first_day) + + +@dag( + dag_id="world_live", + schedule=LIVE_TICK, + start_date=START_DATE, + catchup=False, + is_paused_upon_creation=True, + max_active_runs=1, + tags=TAGS, +) +def world_live(): + """Выключатель: пока включён, мир живёт день за днём. + + Своей работы у выключателя нет — он дёргает `world_live_day` и ждёт конца. + Ожидание тут несущая конструкция, а не вежливость: без него тик шёл бы + независимо от хода дня, лишние прогоны скопились бы очередью, и мир потом + промчался бы по ней без всякого темпа. + + Ждём триггером, а не сенсором: оператор опрашивает тот прогон, который сам + и создал, и ссылка на дочерний прогон видна прямо отсюда. Упал день — + краснеет и выключатель. + """ + TriggerDagRunOperator( + task_id="trigger_live_day", + trigger_dag_id="world_live_day", + wait_for_completion=True, + ) + + +world_next_day() +world_live_day() +world_live() diff --git a/infra/airflow/Dockerfile b/infra/airflow/Dockerfile index a82e296..b6f1f7b 100644 --- a/infra/airflow/Dockerfile +++ b/infra/airflow/Dockerfile @@ -4,6 +4,9 @@ USER root # Официальный провайдер Kafka использует тот же confluent-kafka. Пробникам # не нужны его подключения и обёртки, поэтому в образ добавлен сам клиент. +# +# Провайдера docker, которым пульт мира зовёт генератор, ставить не нужно: +# база несёт его сама — 4.5.7 к Airflow 3.3.0, замер 13 августа 2026 года. RUN uv pip install --python /home/airflow/.local/bin/python --no-cache \ "clickhouse-connect==1.6.0" \ "confluent-kafka==2.15.0" \