diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..1536767 --- /dev/null +++ b/.dockerignore @@ -0,0 +1,29 @@ +.git/ +.agents/ +.claude/ +.vscode/ +.idea/ + +.env +.env.* +*.pem +*.key +*.crt +secrets/ +credentials/ + +**/__pycache__/ +**/*.py[cod] +**/.venv/ +**/.pytest_cache/ +**/.ruff_cache/ +**/.mypy_cache/ +**/.ty_cache/ +**/*.egg-info/ +**/build/ +**/dist/ + +tmp/ +logs/ +*.log +*.tmp diff --git a/.env.example b/.env.example index 46d87c0..d944b4d 100644 --- a/.env.example +++ b/.env.example @@ -6,6 +6,12 @@ CLICKHOUSE_02_HTTP_PORT=28124 CLICKHOUSE_02_TCP_PORT=29001 KAFKA_IMAGE=apache/kafka:4.3.1 KAFKA_EXTERNAL_PORT=29092 + +# Разовая служба генератора. Адрес брокера здесь внутренний: команда идёт +# из сети Compose. +KAFKA_BOOTSTRAP_SERVERS=kafka:9092 +KAFKA_TOPIC=hits + PROMETHEUS_IMAGE=prom/prometheus:v3.13.2 PROMETHEUS_PORT=29090 GRAFANA_IMAGE=grafana/grafana:13.1.1 diff --git a/.gitignore b/.gitignore index 1968751..9cda0e1 100644 --- a/.gitignore +++ b/.gitignore @@ -31,6 +31,7 @@ credentials/ # Логи и временные файлы logs/ +tmp/ *.log *.tmp *.bak diff --git a/Makefile b/Makefile index 0d9d51e..c23fa54 100644 --- a/Makefile +++ b/Makefile @@ -1,6 +1,9 @@ COMPOSE ?= docker compose +GENERATOR_DAY ?= 0 +GENERATOR_LIMIT ?= +GENERATOR_SPEED ?= -.PHONY: up down clean ps logs config-test lint typecheck test docs smoke check-clickhouse check-services +.PHONY: up down clean ps logs generate-batch generate-live config-test lint typecheck test docs smoke check-clickhouse check-services up: $(COMPOSE) up --detach --build --wait --wait-timeout 600 @@ -17,6 +20,14 @@ ps: logs: $(COMPOSE) logs --follow +generate-batch: + $(COMPOSE) --profile generator run --rm generator batch --day "$(GENERATOR_DAY)" \ + $(if $(GENERATOR_LIMIT),--limit "$(GENERATOR_LIMIT)") + +generate-live: + $(COMPOSE) --profile generator run --rm generator live --day "$(GENERATOR_DAY)" \ + $(if $(GENERATOR_SPEED),--speed "$(GENERATOR_SPEED)") + config-test: COMPOSE_BIN="$(COMPOSE)" ./scripts/config-test.sh diff --git a/README.md b/README.md index e975c63..5ef3fce 100644 --- a/README.md +++ b/README.md @@ -9,9 +9,11 @@ [«Боевой реализм стенда (v2)»](docs/specs/2026-07-30-stand-v2-realism.md). Сейчас работают кластер ClickHouse из двух шардов, отдельный clickhouse-keeper, односерверная Kafka в режиме KRaft, Airflow 3.3, Superset, -Prometheus, Grafana и общая база Postgres для метаданных. Начат генератор -кликстрима: в [`generator/`](generator/) заведён контракт схемы события, из -которого собрано [описание выгрузки](docs/formats/clickstream-event.md). +Prometheus, Grafana и общая база Postgres для метаданных. Генератор +кликстрима в [`generator/`](generator/) собран целиком со стороны клиента: у +него есть контракт схемы события, из которого собрано +[описание выгрузки](docs/formats/clickstream-event.md), модельный мир и +проигрыватель, отправляющий дни в Kafka или в файл. ## Быстрый старт @@ -117,6 +119,70 @@ Kafka по той же причине спрашивают снаружи. Её полного сброса с удалением всех именованных томов используйте `make clean`. Повторный `make up` безопасен: одноразовая подготовка приложений идемпотентна. +## Как позвать генератор + +События производит генератор из [`generator/`](generator/). Модельный день — +функция зерна и номера дня, поэтому один и тот же день всегда даёт те же +события: обрыв лечится повтором. Позицию на оси генератор не помнит — какой +день играть, решает зовущий. + +Режима два, и различаются они темпом. **Пакетный** гонит день подряд, без пауз: +так заливается мир и так переигрывается день после обрыва. **Живой** держит темп +модельного времени — по умолчанию ×60: модельные сутки за 24 реальные минуты, +суточная волна разворачивается на глазах, отставание видно в логе. + +На стенде генератор ходит своим образом — разовой службой, которую поднимают +и убирают на один прогон. На каждый режим по цели: + +```bash +# день D0 целиком в топик hits +make generate-batch +# первая сотня событий дня D3 +make generate-batch GENERATOR_DAY=3 GENERATOR_LIMIT=100 +# день D3 живьём, ускорение ×1000 +make generate-live GENERATOR_DAY=3 GENERATOR_SPEED=1000 +``` + +| Переменная | Что задаёт | Умолчание | +| --- | --- | --- | +| `GENERATOR_DAY` | номер дня на оси мира (D0 — первый) | `0`, задано в `Makefile` | +| `GENERATOR_LIMIT` | потолок событий на прогон; только `generate-batch` | нет: день целиком | +| `GENERATOR_SPEED` | ускорение модельного времени; только `generate-live` | ×60, задано в генераторе | + +Куда отправлять, обе цели берут из `.env`: `KAFKA_BOOTSTRAP_SERVERS` (внутри +сети Compose это `kafka:9092`) и `KAFKA_TOPIC` (`hits`). + +Образ цели не пересобирают: нет образа — Compose соберёт его сам, есть — +возьмёт как есть. Пересобрать намеренно, после правки `generator/Dockerfile` +или зависимостей, — отдельной командой: + +```bash +docker compose --profile generator build generator +``` + +Разведено это нарочно: собрать образ и запустить контейнер — разные действия, и +цель запуска, молча пересобирающая образ, стирает между ними границу. Что образ +устарел, видно по собственному прогону — это обратная связь, а не ловушка. + +Посмотреть на события, не поднимая стенд, помогает файловый приёмник: одно +событие — одна строка. Флаг `--limit` берёт начало дня вместо целого дня — +тому, кто смотрит на конвейер, ждать полсотни тысяч событий незачем. + +```bash +uv run --project generator python -m clickstream_generator batch \ + --day 0 --limit 100 --file tmp/day0.jsonl +``` + +Несколько дней подряд играет один запуск: `--days 8`. Всё, что принимается +флагами, принимается и переменными окружения — так генератор позовёт даг +этапа 5; полный список у `--help`. + +**Умолчания есть не у всего, и это нарочно.** Зерно, число дней и темп живого +дня описывают модель мира — их генератор знает сам. День на оси, приёмник и имя +топика описывают стенд, на котором его запустили: не назвали — отказ и ненулевой +код возврата, ещё до первого события. Поэтому умолчание дня и живёт в +`Makefile`: день называет тот, кто запускает, а не тот, кого запускают. + ## Состав и доступ - `clickhouse-01` — инициатор DDL и точка подключения Airflow; diff --git a/compose.yaml b/compose.yaml index 74edbfc..c361ff8 100644 --- a/compose.yaml +++ b/compose.yaml @@ -194,6 +194,25 @@ services: exit 1 fi + generator: + profiles: ["generator"] + build: + context: . + dockerfile: generator/Dockerfile + restart: "no" + depends_on: + kafka-init: + condition: service_completed_successfully + environment: + # Адрес брокера и имя топика — факты стенда, и называет их стенд. + # Остальное (день, зерно, число дней, предел пачки) приходит аргументами + # от того, кто запускает: у службы нет позиции на оси мира. + KAFKA_BOOTSTRAP_SERVERS: ${KAFKA_BOOTSTRAP_SERVERS:-kafka:9092} + KAFKA_TOPIC: ${KAFKA_TOPIC:-hits} + command: ["batch"] + # До #42 от этой разовой службы нет зависимого: включать её в обычный + # `up --wait` нельзя, иначе успешное завершение будет считаться сбоем. + # DDL применяется с ноды 1 по порядку имён файлов после готовности всего кластера. clickhouse-init: image: ${CLICKHOUSE_IMAGE:-clickhouse/clickhouse-server:26.3.17.56} diff --git a/docs/specs/2026-08-01-generator.md b/docs/specs/2026-08-01-generator.md index 1268582..c33cb21 100644 --- a/docs/specs/2026-08-01-generator.md +++ b/docs/specs/2026-08-01-generator.md @@ -259,8 +259,7 @@ отдельным сообщением. Контракт транспорта, не деталь реализации — сторона хранилища читает топик байтами и кладёт сообщение строкой ([ADR 0005](../adr/0005-event-ingestion.md)), поэтому склейка нескольких - событий в одно сообщение сломала бы разбор целиком. Сторожится тестом - приёмника. + событий в одно сообщение сломала бы разбор целиком. - **Даты и время на проводе — ISO-8601.** `EventDate` уезжает как `2026-06-01`, `UTCEventTime` — как `2026-06-01T12:34:56Z`. Довод — читаемость сырья: весь смысл слоя STG в том, что менти открывает колонку `raw` в обычном клиенте и @@ -269,9 +268,10 @@ принимает ISO без плясок. Колонка `ecommerce` — строка, внутри которой лежит экранированный JSON, как отдаёт Метрика. - Форму пинит хранилище (#43) как первый потребитель, сериализатор (#41) - её соблюдает. Порядок тикетов обратный порядку зависимости, поэтому здесь она - и записана — иначе каждый выберет своё, и разойдётся это уже после приёмки. + Форму реализует сериализатор (#41), хранилище (#43) читает то, что он + положил: порядок тикетов развёрнут 6 августа 2026 года, и отправитель идёт + первым. Записана форма всё равно здесь — иначе каждый выберет своё, и + разойдётся это уже после приёмки. - **Рабочий выбор сериализатора — orjson**: быстрее stdlib json в 5–14 раз, numpy-массивы и datetime сериализует нативно (заметка исследования #31). Смена библиотеки меняет канонические байты, поэтому проходит как @@ -428,6 +428,16 @@ pytest-тест с маркером `perf` и таймаутом-обрубан теперь есть. - **Этап 5 (Airflow)**: живой день со стороны хранилища — обычный ETL-даг по расписанию (~раз в 24 минуты); генератор не дорабатывается. +- **Этап 5, позиция на оси времени.** Проигрыватель состояния не хранит + (раздел 9), поэтому вести позицию — работа того, кто его зовёт. Живёт она + переменной Airflow: отсутствие переменной означает мир в зерновом + состоянии, заводит и двигает её только даг `next_day` и только по успеху. + Довод — генератор отдельная и заменяемая сущность, привязывать его к + хранилищу незачем, а `make clean` сносит том метаданных Airflow вместе с + данными ClickHouse, так что позиция и мир чистятся одной командой. + Расхождение переменной с данными (менти почистил партицию руками) лечится + документацией, а не сторожем: постоянная сверка была бы проверкой, красной + в норме ([ADR 0002](../adr/0002-monitoring-scope.md)). - **Этап 7 (эталонный мир)**: пересборка — это манифест, не артефакт; CI-генерация на amd64 и arm64. - **Будущие лабы**: перезаливка дня X пакетным режимом проигрывателя — @@ -865,11 +875,105 @@ pytest-тест с маркером `perf` и таймаутом-обрубан порог скорости дня. Упрётся будущая правка — двигать надо бюджет и его причины, а не долю. +Решено при исполнении #41 (2026-08-07). Первые два пункта — решения владельца +грилингом 7 августа 2026 года, внесённые как есть; остальные приняты при +исполнении. + +- **Генератор живёт в своём контейнере**, и на этапе 5 даг будет запускать + контейнер, а не импортировать пакет в процессе воркера. Довод учебный: так + граница между оркестратором и задачей видна глазами — у задачи своё + окружение, свой образ, свой процесс, а оркестратор передаёт параметры и + читает результат. Это та же ступенька, с которой в бою вырастает + `KubernetesPodOperator`. Побочно это единственный вариант, при котором + «генератор — отдельная и заменяемая сущность» перестаёт быть словами: его + зависимости не смешиваются с окружением Airflow. Хостовый `uv` остаётся для + `make test`, `make lint` и `make typecheck`. Цена названа и принимается: + чтобы даг запускал контейнер, в Airflow пробрасывается сокет докера, а это + доступ, равный root на хосте; на учебном стенде размен допустим, но + записывается уроком — в бою так не делают. Отклонено: *генератор в образе + Airflow, даг импортирует пакет* (так у стенда v1) — механика проще, но менти + видит лишь «даг позвал функцию», а платить пришлось бы сменой базового + образа всего стенда, переездом контекста сборки и смешением зависимостей; + *`ExternalPythonOperator` — задача в отдельном venv внутри образа Airflow*: + изоляцию даёт и сокета не требует, но границу видно в конфигурации, а не в + `docker ps` — остаётся запасным путём, если проброс сокета однажды окажется + неприемлемым. +- **Клиент Kafka — `confluent-kafka`.** Он уже стоит в образе Airflow, на нём + написан пробник `dags/test_kafka.py`, и его же использует официальный + провайдер Kafka для Airflow: один клиент на стенде вместо двух. Объявлен + необязательной группой зависимостей генератора (`[kafka]`), а образ ставит + её из лока: приёмник импортирует клиента внутри функции, поэтому пакет без + этой группы остаётся живым и переносимым — с файловым приёмником. +- **Интерфейс запуска — `python -m clickstream_generator batch|live`.** + Зовущий — контейнер: параметры приходят аргументами или переменными + окружения (`GENERATOR_SEED`, `GENERATOR_DAY`, `GENERATOR_DAYS`, + `GENERATOR_LIMIT`, `GENERATOR_SPEED`, `GENERATOR_FILE`, + `KAFKA_BOOTSTRAP_SERVERS`, `KAFKA_TOPIC`), логи идут в стандартный вывод, + итог виден кодом возврата. Общее у режимов: `--seed`, `--day`, `--days` и + приёмник — `--file` либо `--brokers` с `--topic`. Состояния нет: позицию на + оси ведёт зовущий (хвост этапу 5 — раздел 8). `--days` играет несколько + дней подряд одним запуском — этим зальётся зерновой мир (#42), восемь + запусков службы внутри `make up` были бы плохим ответом. Отклонено: + *режимы флагом `--speed 0`* — различие тогда прячется за числом, тогда как + у режимов разное устройство: у пакетного есть пачка и нет ожидания, у + живого наоборот; *приёмник по умолчанию* — молчание однажды напишет в файл + то, чего ждали в Kafka, поэтому не назван приёмник — ошибка, названы оба — + тоже. +- **Умолчание есть только у того, что описывает мир** (решение владельца при + разборе ревью #41, 2026-08-07). Зерно (каноническое), число дней (1) и темп + живого дня (×60) — свойства модели: они заданы спекой и этим тикетом, и + умолчание здесь — тот же самый ответ. Позиция на оси (`--day`), адрес + брокера, имя топика и путь файла — факты стенда, на котором генератор + запущен; своего умолчания у них нет, и без них прогон падает с ненулевым + кодом до первой сгенерированной строки. Довод: генератор — отдельная и + переносимая сущность, факты про *этот* стенд обязан приносить зовущий + (служба compose, цель `make`, даг этапа 5). Умолчание подменило бы забытый + параметр чужим фактом и выдало бы неверный ответ за успех — даг `next_day`, + не передавший день, переиграл бы нулевой, вышел бы с нулём, был бы записан + удачей и подвинул позицию: мир и позиция разъехались бы молча. Отличие от + `--days 1` названо прямо: единица измерения — минимальный *верный* ответ, а + `0` в позиции — конкретный *неверный*, поданный тихо. Переменные окружения + (`GENERATOR_DAY`, `KAFKA_TOPIC`) умолчаниями не считаются: это тот же + зовущий, только другим каналом. Следствие для обвязки: имя топика и день + служба compose и цели `make` передают явно. +- **Ограниченная пачка — флаг `--limit N` пакетного режима**: потолок событий + на весь прогон, а не на каждый день. Довод от менти: тот, кто хочет + посмотреть на конвейер, не должен ждать полсотни тысяч событий. Сочетается + с файловым приёмником, даёт взять срез, ещё не приезжавший в ODS, а число + отправленного печатается всегда — счёт нужен обеим сторонам (на нём стоят + проверки #43). На живой день пачка не распространяется: ожидание там + снимает ускорение. +- **Ключа у сообщения Kafka нет.** `WatchID` ключом был бы ключом только на + вид: обещание Kafka про ключ — «сообщения одного ключа лежат в одном + разделе и сохраняют порядок», а у ряда, где ключи не повторяются, обещать + нечего; читатель же прочёл бы в нём смысл, которого нет. Осмысленный ключ + здесь — `ClientID`, но раскладку по разделам стенд намеренно оставляет + транспорту (раздел 2), и лаба «какая нода читала топик» живёт как раз тем, + что раскладка не предрешена. Проверено 2026-08-07 (Context7 по + `CONFIGURATION.md` librdkafka и снятие умолчаний с самой библиотеки, + поставленной `confluent-kafka` 2.15.0, через `rd_kafka_conf_dump`): + умолчание топикового `partitioner` — `consistent_random`, то есть у + сообщения без ключа раздел случайный; `sticky.partitioning.linger.ms` по + умолчанию 10 — окно, в котором подряд идущие сообщения без ключа липнут к + одному разделу. Документация умолчаний не называет, поэтому оба числа взяты + у библиотеки. Следствие: день целиком идёт секунды и ложится в оба раздела, + а короткая пачка укладывается в одно окно и уезжает в один — поэтому «топик + прочитан обеими нодами» проверяется днём и разово. +- **Приёмников два, а режимов три.** «Kafka пачкой» и «Kafka с темпом ×60» — + один и тот же приёмник у разных зовущих: темп держит проигрыватель, знающий + время события. Приёмник, умеющий ускорение, снова знал бы о содержимом — + ровно то, чего правило «глупых приёмников» (раздел 4) не хочет. +- **Раскладка образа повторяет раскладку репозитория**: `/app` — корень, + пакет лежит исходником в `/app/generator/src` и виден через `PYTHONPATH`, + каталог товаров — в `/app/data`. Причина не в красоте: путь к каталогу + считается от модуля вверх по дереву (`catalog.py`, `parents[3]`), а + поставленный не editable пакет переехал бы в `site-packages` — и те же + `parents[3]` указали бы внутрь `.venv`. Тогда каталог в образе есть, а + модуль его не находит. Проверено запуском в контейнере; заодно день, + сыгранный в образе, совпал побайтово с днём, сыгранным на машине. + Остаётся открытым, за тикетами: -- интерфейс запуска генератора (CLI / цели make) и как он делит режимы - проигрывателя; кто его зовёт в стенде — даги `world_init`/`next_day` из - оценки мастер-спеки (раздел 9) — и в каком контейнере он живёт (#41); - как фиксируется «зерновой» мир конца этапа 2 (раздел 9 мастер-спеки): с манифестным решением напрашивается мини-манифест зернового мира — форма за #42; diff --git a/generator/Dockerfile b/generator/Dockerfile new file mode 100644 index 0000000..cfc4f54 --- /dev/null +++ b/generator/Dockerfile @@ -0,0 +1,50 @@ +# Образ генератора: тот самый «канонический контейнер», внутри которого спека +# обещает побайтовую воспроизводимость (раздел 2). Отсюда правило: пакеты +# ставятся только из `uv.lock` — второй список зависимостей писал бы на стенде +# не те байты, что сторожит `make test`. +# +# Контекст сборки — корень репозитория, а не `generator/`: в образ едет ещё и +# каталог товаров. +# База записана здесь и прибита к patch-версии: это канонический контейнер из +# раздела 2 спеки, и его база входит в тот же фиксированный набор, что `uv.lock`. +# Из `.env` её не берём — иначе описание воспроизводимой сборки распалось бы +# между Dockerfile и обвязкой запуска. +FROM python:3.14.7-slim + +COPY --from=ghcr.io/astral-sh/uv:0.11.21 /uv /bin/uv + +# Раскладка внутри образа повторяет раскладку репозитория: `/app` — его корень. +# Так надо не для красоты. Каталог товаров генератор ищет от своего модуля +# вверх по дереву (`catalog.py`, `parents[3]`): каталог — часть мира, а не +# запуска, и снаружи не настраивается (решение #39). Значит, взаимное +# расположение пакета и `data/` обязано сохраниться, иначе первый же импорт +# каталога упрётся в отсутствующий файл. +WORKDIR /app/generator + +ENV UV_LINK_MODE=copy \ + UV_PYTHON_DOWNLOADS=never + +# Зависимости — своим слоем: он пересобирается, только когда меняется лок. +# `--extra kafka` ставит клиента, которым живёт приёмник Kafka; `--no-dev` +# оставляет за бортом pytest, ruff и ty — их дом на машине разработчика. +COPY generator/pyproject.toml generator/uv.lock ./ +RUN uv sync --frozen --no-dev --no-install-project --extra kafka + +# Пакет в `.venv` не ставится: он лежит исходником и виден через `PYTHONPATH`. +# Поставленный не editable, он переехал бы в `site-packages`, и те же +# `parents[3]` указали бы внутрь `.venv` — каталог в образе есть, а модуль его +# не находит. Заодно правка кода не трогает слой зависимостей. +COPY generator/src ./src +COPY data /app/data + +ENV PATH="/app/generator/.venv/bin:$PATH" \ + PYTHONPATH="/app/generator/src" \ + PYTHONUNBUFFERED=1 + +# Не root: генератору хватает права читать своё и писать в сеть. +RUN useradd --create-home --uid 1000 generator +USER generator + +# Параметры прогона передаёт зовущий — аргументами или окружением; итог он +# читает кодом возврата, а логи идут в стандартный вывод. +ENTRYPOINT ["python", "-m", "clickstream_generator"] diff --git a/generator/README.md b/generator/README.md index b2f73a8..a801a44 100644 --- a/generator/README.md +++ b/generator/README.md @@ -4,9 +4,11 @@ образцу облачной выгрузки Яндекс Метрики. Устройство и принятые решения — спека [«Генератор (этап 2)»](../docs/specs/2026-08-01-generator.md). -События уже есть: день-функция отдаёт по паре (зерно, D) упорядоченный поток -трёх видов — просмотр страницы, корзина, покупка. Клиентская сторона на этом -целая; заказы бэкенда и запуск снаружи — за следующими тикетами. +События уже есть и уже уезжают: день-функция отдаёт по паре (зерно, D) +упорядоченный поток трёх видов — просмотр страницы, корзина, покупка, — +канонический сериализатор превращает его в JSON, а проигрыватель гонит в файл +или в Kafka. Клиентская сторона на этом целая; заказы бэкенда — за следующими +этапами. ## Как это работает @@ -60,6 +62,16 @@ D0 живёт предыстория, поэтому любой день соб Своя случайность, поэтому правка торговли трафик не двигает. - `src/clickstream_generator/ids.py` — номера событий: неповторяющиеся и ниже 2^53. Обещание одно на обе половины дня, поэтому и живёт отдельно. +- `src/clickstream_generator/serialize.py` — канонический сериализатор: + единственное место, где событие целиком превращается в JSON. Порядок ключей, + все 47 колонок всегда и форма на проводе — ISO-8601. +- `src/clickstream_generator/sinks.py` — приёмники: файл (одно событие — одна + строка) и Kafka (одно событие — одно сообщение). Про содержимое они не + знают; там же довод, почему у сообщения нет ключа. +- `src/clickstream_generator/player.py` — проигрыватель: гонит дни в приёмник + пачкой или с темпом живого дня. Состояния не хранит. +- `src/clickstream_generator/cli.py` — интерфейс запуска: параметры + аргументами или окружением, логи в стандартный вывод, итог кодом возврата. - `src/clickstream_generator/schema.py` — контракт схемы: чистые данные о колонках выгрузки. Собственность генератора; из него выводятся сам генератор, его валидация и описание выгрузки в доках. @@ -68,10 +80,14 @@ D0 живёт предыстория, поэтому любой день соб из контракта. Документ руками не правят — пересобирают. - `tests/` — инварианты контракта, свежесть описания и обещания мира: чистота от зерна, приток, гарантия двухкуковых пар, форма суточной волны - и сборка визитов по задокументированным правилам. + и сборка визитов по задокументированным правилам. Там же побайтовое + обещание, доведённое до диска: два прогона дня в файл дают тот же файл, а + строк в нём ровно столько, сколько событий. ## Команды +Проиграть день — [быстрый старт](../README.md#как-позвать-генератор) и `--help`. + Из корня репозитория: - `make test` — тесты генератора; diff --git a/generator/pyproject.toml b/generator/pyproject.toml index f551fa2..1942065 100644 --- a/generator/pyproject.toml +++ b/generator/pyproject.toml @@ -11,9 +11,18 @@ dependencies = [ "orjson>=3.11.9", ] +[project.optional-dependencies] +# Клиент Kafka — необязательная часть пакета: без него генератор пишет в файл +# и остаётся переносимым (спека генератора, раздел 9). Приёмник импортирует +# клиента внутри функции, поэтому пакет без этой группы живой. +kafka = ["confluent-kafka>=2.6"] + [dependency-groups] dev = [ "pytest>=8", + # Тот же клиент нужен и на машине разработчика: иначе ty не разберёт + # импорт приёмника, а ruff и тесты видели бы половину кода. + "confluent-kafka>=2.6", "ruff>=0.15", # Типчекер; молод (0.0.x), но версия закреплена в uv.lock. Настройка # сверена с доками Astral через Context7: запуск `uv run ty check`, diff --git a/generator/src/clickstream_generator/__main__.py b/generator/src/clickstream_generator/__main__.py new file mode 100644 index 0000000..451ea86 --- /dev/null +++ b/generator/src/clickstream_generator/__main__.py @@ -0,0 +1,9 @@ +"""Точка входа пакета: `python -m clickstream_generator`. + +Код возврата уходит наружу как есть — по нему судит о прогоне и контейнер, и +даг, который его запустит (спека генератора, раздел 9). +""" + +from clickstream_generator.cli import main + +raise SystemExit(main()) diff --git a/generator/src/clickstream_generator/cli.py b/generator/src/clickstream_generator/cli.py new file mode 100644 index 0000000..45fcb1c --- /dev/null +++ b/generator/src/clickstream_generator/cli.py @@ -0,0 +1,207 @@ +"""Интерфейс запуска: `python -m clickstream_generator batch|live`. + +Зовущий — контейнер (спека генератора, раздел 9), и интерфейс сделан под него: +параметры приходят аргументами или переменными окружения, логи идут в +стандартный вывод, итог виден кодом возврата. Ни файла настроек, ни состояния +между запусками нет — позиция на оси принадлежит тому, кто зовёт. + +Зовущих трое, и все трое видны в форме команд: + +- даги `world_init` и `next_day` этапа 5 — по дню за запуск, приёмник Kafka; +- заливка зернового мира (#42) — восемь дней подряд одним запуском: `--days`; +- проверки хранилища (#43) — ограниченная пачка в файл: `--limit` и `--file`. + +**Режимы разведены командами, а не флагом**, потому что различаются не темпом +в числе, а тем, что у них разное: у пакетного есть `--limit` и нет ожидания, у +живого есть `--speed` и нет пачки. Один флаг `--speed 0` прятал бы это +различие за числом, а команда называет его словом. Ограниченная пачка живому +дню не полагается: ждать там нечего — ожидание снимает ускорение. + +**Приёмник выбирается тем, что для него назвали**: `--file` или `--brokers`. +Оба сразу — ошибка, ни одного — тоже: молча выбранный по умолчанию приёмник +однажды напишет в файл то, чего ждали в Kafka. + +**Умолчания есть только у того, что описывает мир, а не окружение.** Зерно, +число дней и темп живого дня — свойства модели, они заданы спекой и тикетом, +и умолчание здесь — тот же самый ответ. Позиция на оси (`--day`), адрес +брокера, имя топика и путь файла — факты стенда, на котором нас запустили: +генератор их не знает и знать не должен, он отдельная и переносимая сущность. +Подставь он своё умолчание — забытый параметр превратился бы в неверный ответ, +поданный как успех. +""" + +import argparse +import logging +import os +import sys +from collections.abc import Callable, Sequence +from contextlib import closing +from pathlib import Path + +from clickstream_generator import player +from clickstream_generator.seeds import CANONICAL_SEED +from clickstream_generator.sinks import FileSink, KafkaSink, Sink + +DEFAULT_SPEED = 60.0 + + +def main(argv: Sequence[str] | None = None) -> int: + """Разобрать аргументы, проиграть дни, вернуть код возврата.""" + parser = _parser() + options = parser.parse_args(argv) + logging.basicConfig( + level=logging.INFO, format="%(message)s", stream=sys.stdout, force=True + ) + + _check(parser, options) + try: + with closing(_sink(parser, options)) as sink: + player.play( + sink, + seed=options.seed, + first_day=options.day, + days=options.days, + limit=options.limit, + speed=options.speed, + ) + except (OSError, RuntimeError) as failure: + logging.error("прогон не удался: %s", failure) + return 1 + return 0 + + +def _parser() -> argparse.ArgumentParser: + """Разбор командной строки; умолчания приходят из окружения.""" + parser = argparse.ArgumentParser( + prog="python -m clickstream_generator", + description="Проигрыватель модельных дней кликстрима в файл или Kafka.", + ) + modes = parser.add_subparsers(dest="mode", required=True) + + batch = modes.add_parser( + "batch", help="пачкой, без пауз: заливка снимка и переигровка дня" + ) + _common(batch) + batch.set_defaults(speed=None) + batch.add_argument( + "--limit", + type=int, + default=_env_int("GENERATOR_LIMIT", None), + metavar="N", + help="взять не больше N событий на весь прогон, а не день целиком", + ) + + live = modes.add_parser("live", help="живой день: темп модельного времени") + _common(live) + live.set_defaults(limit=None) + live.add_argument( + "--speed", + type=float, + default=_env_float("GENERATOR_SPEED", DEFAULT_SPEED), + metavar="X", + help=f"ускорение модельного времени (по умолчанию ×{DEFAULT_SPEED:.0f})", + ) + return parser + + +def _common(parser: argparse.ArgumentParser) -> None: + """Что спрашивают у обоих режимов: какой мир, какие дни и куда.""" + parser.add_argument( + "--seed", + type=int, + default=_env_int("GENERATOR_SEED", CANONICAL_SEED), + metavar="N", + help="зерно мира (по умолчанию каноническое)", + ) + parser.add_argument( + "--day", + type=int, + default=_env_int("GENERATOR_DAY", None), + metavar="D", + help="номер дня на оси мира (D0 — первый); умолчания нет —" + " позицию ведёт зовущий", + ) + parser.add_argument( + "--days", + type=int, + default=_env_int("GENERATOR_DAYS", 1), + metavar="N", + help="сколько дней подряд проиграть одним запуском", + ) + parser.add_argument( + "--file", + type=Path, + default=_env("GENERATOR_FILE", None, Path), + metavar="ПУТЬ", + help="приёмник — файл: одно событие в строке", + ) + parser.add_argument( + "--brokers", + default=os.environ.get("KAFKA_BOOTSTRAP_SERVERS"), + metavar="АДРЕС", + help="приёмник — Kafka: адреса брокеров через запятую", + ) + parser.add_argument( + "--topic", + default=os.environ.get("KAFKA_TOPIC"), + metavar="ИМЯ", + help="топик Kafka; умолчания нет — имя топика принадлежит стенду, а не пакету", + ) + + +def _check(parser: argparse.ArgumentParser, options: argparse.Namespace) -> None: + """День на оси обязан быть назван — иначе прогон не начинается.""" + if options.day is None: + parser.error( + "день на оси не назван: передайте --day или GENERATOR_DAY." + " Позицию ведёт зовущий — проигрыватель её не помнит" + ) + + +def _sink(parser: argparse.ArgumentParser, options: argparse.Namespace) -> Sink: + """Приёмник по названному: файл либо Kafka, но не оба и не ничего.""" + if options.file and options.brokers: + parser.error("названы оба приёмника: оставьте --file либо --brokers") + if options.file: + return FileSink(options.file) + if options.brokers: + if not options.topic: + parser.error( + "топик не назван: передайте --topic или KAFKA_TOPIC." + " Имя топика — факт стенда, генератор его не знает" + ) + try: + return KafkaSink(options.brokers, options.topic) + except ImportError: + # Ловим там, где возникает: обёрнутый вокруг всего прогона, + # этот перехват однажды объявил бы «нет клиента Kafka» о чужой + # сорванной загрузке модуля. + parser.error( + "приёмник Kafka требует клиента: поставьте пакет с группой" + " зависимостей kafka (`uv sync --extra kafka`)" + ) + parser.error("приёмник не назван: нужен --file либо --brokers") + + +def _env[T](name: str, fallback: T, kind: Callable[[str], T]) -> T: + """Значение переменной окружения нужного типа; нет переменной — умолчание. + + Пустая строка считается отсутствием: в compose так выглядит переменная, + которую не задали, — `${GENERATOR_DAYS:-}`. Мусор в переменной называется + вместе с её именем: из контейнера иначе не видно, чьё это значение. + """ + value = os.environ.get(name) + if not value: + return fallback + try: + return kind(value) + except ValueError: + raise SystemExit(f"переменная {name} не разбирается: {value!r}") from None + + +def _env_int(name: str, fallback: int | None) -> int | None: + return _env(name, fallback, int) + + +def _env_float(name: str, fallback: float) -> float: + return _env(name, fallback, float) diff --git a/generator/src/clickstream_generator/player.py b/generator/src/clickstream_generator/player.py new file mode 100644 index 0000000..95aeaea --- /dev/null +++ b/generator/src/clickstream_generator/player.py @@ -0,0 +1,157 @@ +"""Проигрыватель: гонит дни мира в приёмник — пачкой или с темпом живого дня. + +Состояния у него нет (спека генератора, раздел 9). Зерно и номер дня приходят +параметрами, позицию на оси он не хранит и из данных не выводит: её ведёт тот, +кто зовёт, — на этапе 5 это переменная Airflow у дага `next_day`. Отсюда и +переигровка обрыва: позвали тот же день заново — получили те же `WatchID`, и +дедупликация склеила повтор. + +Два режима отличаются только темпом. Пакетный шлёт события подряд, без пауз, — +это заливка снимка и переигровка дня. Живой держит модельное время: событие +уезжает тогда, когда до него дошли модельные часы, поделённые на ускорение. +По умолчанию ускорение ×60 — модельные сутки за 24 реальные минуты (спека, +раздел 5): суточная волна разворачивается на глазах. + +Тайминги печатаются раздельно — генерация и доставка, как требует спека +(раздел 5): это разные машины разной природы, и сложенные в одно число они +перестают что-либо говорить. Сериализация считается частью генерации: она +рождает те самые байты, которые сторожит манифест. +""" + +import logging +import time +from dataclasses import dataclass + +import numpy as np +from numpy.typing import NDArray + +from clickstream_generator import day as day_module +from clickstream_generator import serialize +from clickstream_generator.sinks import Sink + +log = logging.getLogger(__name__) + +# Как часто живой режим отчитывается о ходе дня. Минута реального времени — это +# час модельного при ×60: отчёт на каждый модельный час. +REPORT_SECONDS = 60.0 + + +@dataclass(frozen=True, slots=True) +class Played: + """Итог прогона: сколько уехало и за сколько.""" + + events: int + generated_seconds: float + delivered_seconds: float + + +def play( + sink: Sink, + seed: int, + first_day: int, + days: int = 1, + limit: int | None = None, + speed: float | None = None, +) -> Played: + """Проиграть `days` дней подряд начиная с `first_day` в приёмник `sink`. + + `limit` — потолок событий на весь прогон: срез для того, кто смотрит на + конвейер и не хочет ждать целый день. `speed` — ускорение живого режима; + `None` означает пакетный, то есть без пауз вовсе. + """ + events = 0 + generated = 0.0 + delivered = 0.0 + + for number in range(first_day, first_day + days): + left = None if limit is None else limit - events + if left is not None and left <= 0: + break + + clock = time.monotonic() + today = day_module.stream(seed, number) + payloads = serialize.events(today, limit=left) + seconds = _event_seconds(today, len(payloads)) + spent = time.monotonic() - clock + generated += spent + log.info("день %d: событий %d, генерация %.1f с", number, len(payloads), spent) + + clock = time.monotonic() + if speed is None: + _send(sink, payloads) + else: + _send_paced(sink, payloads, seconds, speed) + # Рубеж дня: доставка асинхронна, и без него напечатанное время + # означало бы только «события легли в очередь отправителя». + sink.flush() + spent = time.monotonic() - clock + delivered += spent + + events += len(payloads) + log.info( + "день %d: отправлено %d, доставка %.1f с", number, len(payloads), spent + ) + + log.info( + "итого отправлено %d событий: генерация %.1f с, доставка %.1f с", + events, + generated, + delivered, + ) + return Played( + events=events, generated_seconds=generated, delivered_seconds=delivered + ) + + +def _send(sink: Sink, payloads: list[bytes]) -> None: + """Пакетно: подряд и без пауз.""" + for payload in payloads: + sink.send(payload) + + +def _send_paced( + sink: Sink, payloads: list[bytes], seconds: NDArray[np.int64], speed: float +) -> None: + """С темпом: событие уезжает, когда до него дошло модельное время. + + Отставание не догоняется рывком и не прячется: спешить некуда — событие + всё равно уедет, — а вот увидеть отставание в логе нужно, иначе живой + режим врёт про темп. Обгонять модельное время нельзя, отставать можно, и + именно это печатает отчёт. + """ + started = time.monotonic() + origin = int(seconds[0]) if seconds.size else 0 + reported = started + lag = 0.0 + + for number, payload in enumerate(payloads): + due = started + (int(seconds[number]) - origin) / speed + now = time.monotonic() + if now < due: + time.sleep(due - now) + else: + lag = max(lag, now - due) + sink.send(payload) + + now = time.monotonic() + if now - reported >= REPORT_SECONDS: + log.info( + "проиграно %d из %d, модельное время %s, лаг %.1f с", + number + 1, + len(payloads), + _model_time(int(seconds[number]) - origin), + lag, + ) + reported = now + lag = 0.0 + + +def _event_seconds(today: day_module.Day, count: int) -> NDArray[np.int64]: + """Секунды событий абсолютной меткой — по ним живой режим держит темп.""" + times: NDArray[np.datetime64] = today.columns["UTCEventTime"][:count] + return times.astype("datetime64[s]").astype(np.int64) + + +def _model_time(elapsed: int) -> str: + """Прожитое модельное время дня в виде `ЧЧ:ММ` — от первого события.""" + return f"{elapsed // 3600:02d}:{elapsed % 3600 // 60:02d}" diff --git a/generator/src/clickstream_generator/serialize.py b/generator/src/clickstream_generator/serialize.py new file mode 100644 index 0000000..c2c893d --- /dev/null +++ b/generator/src/clickstream_generator/serialize.py @@ -0,0 +1,78 @@ +"""Канонический сериализатор: единственное место, где событие целиком → JSON. + +Правило «сериализатор один» (спека генератора, разделы 4 и 6) — не про +экономию строк, а про канон: два прогона одного дня обязаны дать те же байты, +а байты рождаются здесь. Второе место, собирающее событие руками, разошлось бы +с этим по экранированию, порядку ключей или записи чисел — и разошлось бы +молча. Граница правила проходит по событию, а не по всякому JSON: вложенный +блок `ecommerce` собирает `commerce`, и это часть содержимого колонки, а не +второй сериализатор. + +Что делает канон: + +- **Порядок ключей — порядок контракта схемы.** Он берётся из `schema.COLUMNS` + и нигде не повторяется: два источника порядка разъехались бы при первой же + вставке колонки. +- **Все 47 ключей всегда.** Пусто по контракту — пустое значение: пустой + массив, пустая строка, ноль. Пропавший ключ увёл бы событие в брак целиком: + строгий приём хранилища сверяет набор ключей (ADR 0005). +- **Даты и время — ISO-8601** (спека, раздел 4): `EventDate` уезжает как + `2026-06-01`, `UTCEventTime` — как `2026-06-01T12:34:56Z`. Довод — читаемость + сырья: менти открывает колонку `raw` обычным клиентом и разбирает событие + глазами, а число эпохи этот урок убивает. +- **Одно событие — один документ JSON**, без перевода строки внутри: приёмник + сам решает, чем их разделить. + +Колонки переводятся в питоновские значения целиком, а не по строкам: numpy +делает это одним вызовом на колонку, и на дне в полсотни тысяч событий разница +заметна. Обратная сторона — день лежит в памяти дважды; проигрыватель поэтому +и берёт его днями, а не горизонтом целиком. +""" + +from typing import Any + +import numpy as np +import orjson +from numpy.typing import NDArray + +from clickstream_generator import schema +from clickstream_generator.day import Day + +_ARRAY_PREFIX = "Array(" + + +def events(day: Day, limit: int | None = None) -> list[bytes]: + """Канонические байты событий дня: по документу JSON на событие. + + `limit` берёт первые события дня и на этом останавливается — срез для + того, кто смотрит на конвейер и не хочет ждать целый день (спека, + раздел 9). Ограничение считается до сериализации: платить за то, что не + поедет, незачем. + """ + count = len(day) if limit is None else min(limit, len(day)) + names = tuple(column.name for column in schema.COLUMNS) + values = [ + _values(column, day.columns[column.name][:count]) for column in schema.COLUMNS + ] + return [ + orjson.dumps(dict(zip(names, row, strict=True))) + for row in zip(*values, strict=True) + ] + + +def _values(column: schema.Column, values: NDArray[Any]) -> list[Any]: + """Колонка питоновскими значениями — в той записи, в какой уедет на провод. + + Массив узнаётся по типу ClickHouse, а не по `numpy_dtype`: у колонки-массива + там записан тип элемента (`uint32`), и от скалярной колонки её этим не + отличить. + """ + if column.clickhouse_type.startswith(_ARRAY_PREFIX): + # Колонка-массив: в ячейке лежит свой массив, пустой у события, + # которому эта колонка не по смыслу. + return [cell.tolist() for cell in values] + if column.numpy_dtype == "datetime64[D]": + return np.datetime_as_string(values, unit="D").tolist() + if column.numpy_dtype == "datetime64[s]": + return np.datetime_as_string(values, unit="s", timezone="UTC").tolist() + return values.tolist() diff --git a/generator/src/clickstream_generator/sinks.py b/generator/src/clickstream_generator/sinks.py new file mode 100644 index 0000000..5a32cef --- /dev/null +++ b/generator/src/clickstream_generator/sinks.py @@ -0,0 +1,137 @@ +"""Приёмники: куда уезжают канонические байты. Про содержимое они не знают. + +Приёмник глуп по замыслу (спека генератора, раздел 4): он берёт готовый байт +события и доставляет его. Ни формы, ни темпа он не решает — форму задал +сериализатор, темп задаёт проигрыватель. Отсюда и весь их интерфейс: принять +событие, дождаться принятого, закрыться. + +Приёмников два, а режимов три: «Kafka пачкой» и «Kafka с темпом ×60» +различаются не приёмником, а тем, кто его зовёт. Ускорять доставку, зная о +времени события, значило бы вернуть приёмнику знание о содержимом — ровно то, +чего правило не хочет. + +**Ключа у сообщения Kafka нет** — решение тикета #41, и вот довод. `WatchID` +уникален у каждого события, поэтому ключом он не был бы ключом: обещание Kafka +про ключ — «сообщения одного ключа лежат в одном разделе и сохраняют порядок», а +у ряда, где ключи не повторяются, обещать нечего. Читатель же прочёл бы такой +ключ как смысл, которого в нём нет. Настоящий ключ здесь — `ClientID` +(события одной куки по порядку), и он тикетом не назначен: раскладку по +разделам стенд намеренно оставляет транспорту (спека, раздел 2), а лаба +«какая нода читала топик» живёт как раз тем, что раскладка не предрешена. + +Раздел для сообщения без ключа librdkafka выбирает случайно, но подряд идущие +сообщения на десяток миллисекунд липнут к одному — измеренные умолчания и их +следствия записаны в спеке (раздел 9). +""" + +from pathlib import Path +from typing import Any, Protocol + +# Сколько ждать разбора очереди отправителя, когда она заполнилась, и сколько — +# доставки остатка при закрытии. Оба числа — потолок ожидания, а не пауза: +# обычно ждать не приходится вовсе. +QUEUE_WAIT_SECONDS = 1.0 +FLUSH_WAIT_SECONDS = 60.0 + + +class Sink(Protocol): + """Приёмник: принимает байты события и доводит их до места.""" + + def send(self, payload: bytes) -> None: + """Принять одно событие.""" + + def flush(self) -> None: + """Дождаться, пока принятое дойдёт до места. + + Нужно не только при закрытии: доставка асинхронна, и без этого рубежа + «доставка 0,1 с» в логе означала бы лишь то, что события успели лечь в + очередь отправителя. Проигрыватель ставит рубеж в конце каждого дня. + """ + + def close(self) -> None: + """Закрыть приёмник, дождавшись всего принятого.""" + + +class FileSink: + """Файл: одно событие — одна строка. + + Построчность — контракт файла, а не удобство: файл читают построчно, и + склейка двух событий в строку сломала бы разбор целиком. Файл — кэш чистой + функции (спека, раздел 4): потерял — пересчитал, поэтому места в + репозитории ему не отведено. + """ + + def __init__(self, path: Path) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + self._file = path.open("wb") + + def send(self, payload: bytes) -> None: + self._file.write(payload + b"\n") + + def flush(self) -> None: + self._file.flush() + + def close(self) -> None: + self._file.close() + + +class KafkaSink: + """Kafka: одно событие — одно сообщение. + + «Пачкой» относится к темпу отправки, а не к упаковке: события идут подряд + без пауз, но каждое отдельным сообщением (спека, раздел 4). Отправка + асинхронная — отправитель копит сообщения в своей очереди и шлёт их + пачками сам; отсюда `poll`, который отдаёт нам отчёты о доставке, и + `flush`, без которого хвост очереди уехал бы в никуда вместе с процессом. + """ + + def __init__(self, brokers: str, topic: str) -> None: + # Импорт внутри: клиент — необязательная часть пакета, и без него + # генератор пишет в файл (спека, раздел 9). + from confluent_kafka import Producer + + self._topic = topic + self._failure: str | None = None + self._producer = Producer({"bootstrap.servers": brokers}) + + def send(self, payload: bytes) -> None: + while True: + try: + self._producer.produce( + self._topic, value=payload, on_delivery=self._report + ) + break + except BufferError: + # Очередь отправителя полна: ждём, пока брокер её разберёт. + # `poll` здесь и работа, и пауза — он же отдаёт отчёты. + self._producer.poll(QUEUE_WAIT_SECONDS) + self._producer.poll(0) + self._raise_failure() + + def flush(self) -> None: + remaining = self._producer.flush(FLUSH_WAIT_SECONDS) + if remaining: + raise RuntimeError( + f"Kafka не приняла {remaining} сообщений за" + f" {FLUSH_WAIT_SECONDS:.0f} с: брокер недоступен или не успевает" + ) + self._raise_failure() + + def close(self) -> None: + # Закрывать у отправителя нечего — важно лишь не бросить хвост + # очереди: он уехал бы в никуда вместе с процессом. + self.flush() + + def _report(self, error: Any, message: Any) -> None: + """Отчёт о доставке: первую неудачу запоминаем, остальные не важны.""" + if error is not None and self._failure is None: + self._failure = str(error) + + def _raise_failure(self) -> None: + """Неудачную доставку превращаем в остановку прогона. + + Молча потерянное сообщение — худший исход: счёт в хранилище разойдётся + с числом отправленного, а причина будет забыта. + """ + if self._failure is not None: + raise RuntimeError(f"Kafka не приняла сообщение: {self._failure}") diff --git a/generator/tests/test_player.py b/generator/tests/test_player.py new file mode 100644 index 0000000..0022699 --- /dev/null +++ b/generator/tests/test_player.py @@ -0,0 +1,169 @@ +"""Проигрыватель и его интерфейс: побайтовый повтор, построчность, приёмник. + +Главное здесь — обещание раздела 2 спеки, доведённое до байтов на диске: два +прогона одного дня дают тот же файл. Тем же тестом сторожится построчность: +одно событие — одна строка. Дешёвая половина контракта транспорта — склейка +двух событий сломала бы разбор целиком, потому что хранилище читает топик +байтами. +""" + +import hashlib +import subprocess +import sys +from contextlib import closing +from dataclasses import replace + +import pytest + +from clickstream_generator import cli, player +from clickstream_generator import day as day_module +from clickstream_generator.seeds import CANONICAL_SEED +from clickstream_generator.sinks import FileSink + +DAY = 2 + +# Переменные, которыми зовущий задаёт прогон. Тест, читающий их из окружения +# машины, зелен у одного и красен у другого — а на этой машине они как раз и +# живут: стенд их экспортирует. +LAUNCH_VARIABLES = ( + "GENERATOR_SEED", + "GENERATOR_DAY", + "GENERATOR_DAYS", + "GENERATOR_LIMIT", + "GENERATOR_SPEED", + "GENERATOR_FILE", + "KAFKA_BOOTSTRAP_SERVERS", + "KAFKA_TOPIC", +) + + +@pytest.fixture(autouse=True) +def bare_environment(monkeypatch): + """Прогон тестов не зависит от того, что задано в окружении машины.""" + for name in LAUNCH_VARIABLES: + monkeypatch.delenv(name, raising=False) + + +def _play(path, **options): + """Проиграть в файл и вернуть итог прогона.""" + with closing(FileSink(path)) as sink: + return player.play(sink, seed=CANONICAL_SEED, first_day=DAY, **options) + + +def test_two_runs_give_the_same_file(tmp_path): + """Два прогона дня — одинаковые байты и строка на событие. + + Хеши сравниваются целиком, а не построчно: обещание побайтовое, и + расхождение в одной запятой обязано покраснеть так же, как расхождение в + наборе событий. + + Прогоны идут **разными процессами**, а не двумя вызовами в одном. Внутри + одного интерпретатора зерно хеширования общее на оба прогона, поэтому + зависимость канона от порядка обхода множества такой тест не увидел бы + никогда — а это ровно тот класс расхождений, ради которого обещание и + дано. Заодно день играется тем же путём, каким генератор зовут на самом + деле: через интерфейс запуска, а не через функцию. + """ + first, second = tmp_path / "first.jsonl", tmp_path / "second.jsonl" + _run_apart(first) + _run_apart(second) + + assert _digest(first) == _digest(second) + + events = len(day_module.stream(CANONICAL_SEED, DAY)) + assert len(first.read_bytes().splitlines()) == events + + +def _run_apart(path) -> None: + """Проиграть день отдельным процессом; окружение он берёт от нас.""" + finished = subprocess.run( + [ + sys.executable, + "-m", + "clickstream_generator", + "batch", + "--day", + str(DAY), + "--file", + str(path), + ], + capture_output=True, + text=True, + timeout=300, + check=False, + ) + assert finished.returncode == 0, finished.stderr + + +def test_days_play_in_a_row(tmp_path, monkeypatch): + """Дни идут подряд от названного, а пачка считается на весь прогон. + + Восемь дней одним запуском — то, чем зальётся зерновой мир (#42), поэтому + порядок дней проверяется, а не предполагается. День здесь подменён коротким: + проверяется ход проигрывателя, а не содержимое дня, и платить за полсотни + тысяч событий трижды незачем. + """ + short = _shorten(day_module.stream(CANONICAL_SEED, DAY), rows=2) + asked: list[int] = [] + + def stream(seed: int, number: int): + asked.append(number) + return short + + monkeypatch.setattr(player.day_module, "stream", stream) + path = tmp_path / "three-days.jsonl" + played = _play(path, days=3, limit=5) + + assert asked == [DAY, DAY + 1, DAY + 2] + # Двум дням хватило по два события, третьему досталось последнее: потолок + # считается на весь прогон, а не на каждый день заново. + assert played.events == 5 + assert len(path.read_bytes().splitlines()) == 5 + + +FILE = ["--file", "/dev/null"] +KAFKA = ["--brokers", "kafka:29092", "--topic", "hits"] + + +@pytest.mark.parametrize( + "argv", + [ + pytest.param(["batch", "--day", "0"], id="приёмник не назван"), + pytest.param(["batch", "--day", "0", *FILE, *KAFKA], id="названы оба"), + pytest.param( + ["batch", "--day", "0", "--brokers", "kafka:29092"], id="топик не назван" + ), + pytest.param(["batch", *FILE], id="день не назван"), + ], +) +def test_run_is_refused_loudly(argv): + """Всё, чего проигрыватель не знает, — отказ до первого события. + + Правило одно: параметр, описывающий окружение стенда или позицию на оси + мира, своего умолчания не имеет. Тихо подставленное умолчание — чужой факт, + выданный за наш, и прогон отчитается о нём успехом. + """ + with pytest.raises(SystemExit) as refusal: + cli.main(argv) + assert refusal.value.code == 2 + + +def test_cli_plays_a_limited_batch_into_a_file(tmp_path): + """Тот самый вызов, которым проверки #43 берут срез: код возврата — 0.""" + path = tmp_path / "cli.jsonl" + assert cli.main(["batch", "--day", "0", "--limit", "5", "--file", str(path)]) == 0 + assert len(path.read_bytes().splitlines()) == 5 + + +def _shorten(today: day_module.Day, rows: int) -> day_module.Day: + """Тот же день, но короткий: первые строки всех его рядов.""" + return replace( + today, + columns={name: value[:rows] for name, value in today.columns.items()}, + page=today.page[:rows], + product=today.product[:rows], + ) + + +def _digest(path) -> str: + return hashlib.sha256(path.read_bytes()).hexdigest() diff --git a/generator/tests/test_serialize.py b/generator/tests/test_serialize.py new file mode 100644 index 0000000..9f2927c --- /dev/null +++ b/generator/tests/test_serialize.py @@ -0,0 +1,72 @@ +"""Канон сериализатора: набор ключей, их порядок и форма на проводе. + +Сторожится здесь то, на чём стоит сторона хранилища: строгий приём сверяет +набор ключей и разбирает даты как ISO. Разойдись сериализатор с этим — событие +уйдёт в брак целиком, а поймается это уже на стенде. +""" + +import json + +import pytest + +from clickstream_generator import day as day_module +from clickstream_generator import schema, serialize +from clickstream_generator.seeds import CANONICAL_SEED + +DAY = 2 + + +@pytest.fixture(scope="module") +def events() -> list[dict[str, object]]: + """События дня, разобранные обратно из канонических байтов.""" + today = day_module.stream(CANONICAL_SEED, DAY) + return [json.loads(payload) for payload in serialize.events(today)] + + +def test_every_event_carries_every_column(events): + """Все 47 ключей всегда и в порядке контракта — у любого события. + + «Пусто» по контракту — пустое значение, а не отсутствие ключа: пропавший + ключ уводит событие в брак целиком (ADR 0005). Порядок ключей — часть + канона: от него зависят байты, а значит и хеши манифеста. + """ + names = [column.name for column in schema.COLUMNS] + for event in events: + assert list(event) == names + + +def test_empty_is_a_value_not_a_hole(events): + """У просмотра страницы торговые колонки пусты, но они есть.""" + pageview = next(event for event in events if event["EventType"] == "pageview") + assert pageview["purchaseID"] == [] + assert pageview["productPrice"] == [] + assert pageview["GoalsReached"] == [] + assert pageview["ecommerce"] == "" + + +def test_dates_go_as_iso(events): + """Даты читаются глазами: `2026-06-03` и `2026-06-03T12:34:56Z`. + + Весь смысл слоя STG в том, что менти открывает колонку `raw` обычным + клиентом и разбирает событие сам; число эпохи этот урок убивает. + """ + for event in events[:100]: + assert event["EventDate"] == "2026-06-03" + assert event["UTCEventTime"].endswith("Z") + assert len(event["UTCEventTime"]) == len("2026-06-03T12:34:56Z") + + +def test_ecommerce_is_a_string_with_json_inside(events): + """`ecommerce` уезжает строкой, как отдаёт Метрика, — материал лабы.""" + purchase = next(event for event in events if event["EventType"] == "purchase") + assert isinstance(purchase["ecommerce"], str) + inside = json.loads(purchase["ecommerce"]) + assert inside["purchase"]["actionField"]["id"] == purchase["purchaseID"][0] + + +def test_limit_takes_the_beginning_of_the_day(events): + """Ограниченная пачка — начало дня, а не его пересборка другими байтами.""" + today = day_module.stream(CANONICAL_SEED, DAY) + short = serialize.events(today, limit=10) + assert len(short) == 10 + assert [json.loads(payload) for payload in short] == events[:10] diff --git a/generator/uv.lock b/generator/uv.lock index 63c6dd5..065632a 100644 --- a/generator/uv.lock +++ b/generator/uv.lock @@ -11,8 +11,14 @@ dependencies = [ { name = "orjson" }, ] +[package.optional-dependencies] +kafka = [ + { name = "confluent-kafka" }, +] + [package.dev-dependencies] dev = [ + { name = "confluent-kafka" }, { name = "pytest" }, { name = "ruff" }, { name = "ty" }, @@ -20,12 +26,15 @@ dev = [ [package.metadata] requires-dist = [ + { name = "confluent-kafka", marker = "extra == 'kafka'", specifier = ">=2.6" }, { name = "numpy", specifier = ">=2" }, { name = "orjson", specifier = ">=3.11.9" }, ] +provides-extras = ["kafka"] [package.metadata.requires-dev] dev = [ + { name = "confluent-kafka", specifier = ">=2.6" }, { name = "pytest", specifier = ">=8" }, { name = "ruff", specifier = ">=0.15" }, { name = "ty", specifier = ">=0.0.65" }, @@ -40,6 +49,19 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/d1/d6/3965ed04c63042e047cb6a3e6ed1a63a35087b6a609aa3a15ed8ac56c221/colorama-0.4.6-py2.py3-none-any.whl", hash = "sha256:4f1d9991f5acc0ca119f9d443620b77f9d6b33703e51011c16baf57afb285fc6", size = 25335, upload-time = "2022-10-25T02:36:20.889Z" }, ] +[[package]] +name = "confluent-kafka" +version = "2.15.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/51/90/eb998fedefb63b42910b54b76b7d300ccddc56430a5175122eb60cedc4f3/confluent_kafka-2.15.0.tar.gz", hash = "sha256:7ad9bad1cbabf6713ec039b8204b48d322024fd11397eec88d912e048c732ba7", size = 323365, upload-time = "2026-06-30T19:48:43.395Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/18/20/5c078b9560b8005404e09a5889b9771e8deb5c9dddeb9003f58c7d028258/confluent_kafka-2.15.0-cp314-cp314-macosx_13_0_arm64.whl", hash = "sha256:e781d39ff31d5d10d0502471f92f95aff4f5be1eaae9a1f57aeac164bc9f9029", size = 4343440, upload-time = "2026-06-30T19:48:17.324Z" }, + { url = "https://files.pythonhosted.org/packages/aa/27/fd83281231b6179e3923822e20861dfc1d6c200edb910c3d8e9a15b9e95a/confluent_kafka-2.15.0-cp314-cp314-macosx_13_0_x86_64.whl", hash = "sha256:000463c822a7adc3293369b63688b68d17c03a1a4a64f86182efc08f8b150676", size = 4371051, upload-time = "2026-06-30T19:48:19.29Z" }, + { url = "https://files.pythonhosted.org/packages/5e/3e/fa82e6699144e707972bbf57489598defdf67e844ae2a7d29e0ea4e3a187/confluent_kafka-2.15.0-cp314-cp314-manylinux_2_28_aarch64.whl", hash = "sha256:9ddf4cf4647e5d633ef64e3f5a349d5288edfedf974203993e9d86548c2695be", size = 4979624, upload-time = "2026-06-30T19:48:21.454Z" }, + { url = "https://files.pythonhosted.org/packages/13/ee/b4b6a0da17584b432c83a0500ac79c04864ef92dafdb405614831b995232/confluent_kafka-2.15.0-cp314-cp314-manylinux_2_28_x86_64.whl", hash = "sha256:d8ed33f623ff2a104fb76f99a67e9917f0170fddb4e28380fcdc83347b1646b2", size = 4785678, upload-time = "2026-06-30T19:48:23.068Z" }, + { url = "https://files.pythonhosted.org/packages/a8/f4/19a852d16e8e4f8ac930037d8ecda21a220f6d5d050a59bbce10189ac9ec/confluent_kafka-2.15.0-cp314-cp314-win_amd64.whl", hash = "sha256:2e80bd96f61aae2ffba951754a769e6d3c5ebb5a5e778ed0ae8ad899ea91556b", size = 4744444, upload-time = "2026-06-30T19:48:24.698Z" }, +] + [[package]] name = "iniconfig" version = "2.3.0" diff --git a/scripts/config-test.sh b/scripts/config-test.sh index 4e8fcd4..86af4ae 100755 --- a/scripts/config-test.sh +++ b/scripts/config-test.sh @@ -13,6 +13,12 @@ trap cleanup EXIT "${COMPOSE_CMD[@]}" --project-directory "$ROOT_DIR" config --quiet config_json="$("${COMPOSE_CMD[@]}" --project-directory "$ROOT_DIR" config --format json)" +# Контекст сборки образов — корень репозитория, и локальный `.env` с настоящими +# паролями уехал бы в слой образа молча. Это единственная здешняя ошибка, о +# которой никто не узнает, пока образ не окажется у чужого. +for private_path in '.env' '.env.*' '*.pem' '*.key' '*.crt' 'secrets/' 'credentials/'; do + grep -qxF "$private_path" "$ROOT_DIR/.dockerignore" +done grep -qx 'ARG SUPERSET_BASE_IMAGE' "$ROOT_DIR/infra/superset/Dockerfile" jq -e ' .services.superset.image == "clickstream-superset:local" and