feat(generator): сериализатор, приёмники, проигрыватель и запуск контейнером #61

Merged
ddmitry merged 2 commits from feat/41-serializer-sinks-cli into main 2026-08-07 14:37:25 +03:00
20 changed files with 1204 additions and 16 deletions
+29
View File
@@ -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
+6
View File
@@ -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
+1
View File
@@ -31,6 +31,7 @@ credentials/
# Логи и временные файлы
logs/
tmp/
*.log
*.tmp
*.bak
+20
View File
@@ -13,6 +13,26 @@
складывается фраза — это вопрос владельцу, а не строчка кода: проверки в этом
репозитории однажды уже переросли продукт.
**Усложнение не бесплатно, и платит за него учебная ценность.** Чем сложнее
стенд, тем хуже он как учебный материал: внимание менти конечно, и каждая
конструкция, которую он обязан расшифровать по дороге к уроку, списывается с
этого счёта. Поэтому вопрос выше имеет вторую половину: **что менти платит,
чтобы это прочесть, и что получает взамен?** Ответ «вопрос — и ничего»
означает, что писать не нужно.
Оборонительный код дороже прочего. Читатель не отличает «стоит здесь, потому
что есть настоящая опасность» от «стоит здесь, потому что кто-то был
остроумен», — и предполагает первое. Лишний сторож не безобиден: он
утверждает, будто опасность заслуживает внимания, и этим врёт. Особенно это
касается защиты от ошибки, которую человек делает сам себе и тут же видит:
разбирательство с ней и есть урок, отнимать его не надо.
Ревью само по себе умеет только прибавлять: каждая находка просится в правку,
и за прогон их набираются десятки. Поэтому перед приёмкой заметной работы
идёт отдельный **проход на вычитание** — с правом только резать и с явным
списком неприкосновенного. Сомнение «резать или нет» решается в пользу
реза: вещь, которая себя не защитила, себя не защитила.
## Кому что поручать
Читатель здесь не побочный потребитель, а тот, ради кого стенд существует:
+12 -1
View File
@@ -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
+69 -3
View File
@@ -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;
+19
View File
@@ -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}
+112 -8
View File
@@ -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 в 514 раз,
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;
+50
View File
@@ -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"]
+20 -4
View File
@@ -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` — тесты генератора;
+9
View File
@@ -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`,
@@ -0,0 +1,9 @@
"""Точка входа пакета: `python -m clickstream_generator`.
Код возврата уходит наружу как есть по нему судит о прогоне и контейнер, и
даг, который его запустит (спека генератора, раздел 9).
"""
from clickstream_generator.cli import main
raise SystemExit(main())
+207
View File
@@ -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)
@@ -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}"
@@ -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()
@@ -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}")
+169
View File
@@ -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()
+72
View File
@@ -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]
+22
View File
@@ -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"
+6
View File
@@ -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