feat(generator): сериализатор, приёмники, проигрыватель и запуск контейнером #61
@@ -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,6 +6,12 @@ CLICKHOUSE_02_HTTP_PORT=28124
|
|||||||
CLICKHOUSE_02_TCP_PORT=29001
|
CLICKHOUSE_02_TCP_PORT=29001
|
||||||
KAFKA_IMAGE=apache/kafka:4.3.1
|
KAFKA_IMAGE=apache/kafka:4.3.1
|
||||||
KAFKA_EXTERNAL_PORT=29092
|
KAFKA_EXTERNAL_PORT=29092
|
||||||
|
|
||||||
|
# Разовая служба генератора. Адрес брокера здесь внутренний: команда идёт
|
||||||
|
# из сети Compose.
|
||||||
|
KAFKA_BOOTSTRAP_SERVERS=kafka:9092
|
||||||
|
KAFKA_TOPIC=hits
|
||||||
|
|
||||||
PROMETHEUS_IMAGE=prom/prometheus:v3.13.2
|
PROMETHEUS_IMAGE=prom/prometheus:v3.13.2
|
||||||
PROMETHEUS_PORT=29090
|
PROMETHEUS_PORT=29090
|
||||||
GRAFANA_IMAGE=grafana/grafana:13.1.1
|
GRAFANA_IMAGE=grafana/grafana:13.1.1
|
||||||
|
|||||||
@@ -31,6 +31,7 @@ credentials/
|
|||||||
|
|
||||||
# Логи и временные файлы
|
# Логи и временные файлы
|
||||||
logs/
|
logs/
|
||||||
|
tmp/
|
||||||
*.log
|
*.log
|
||||||
*.tmp
|
*.tmp
|
||||||
*.bak
|
*.bak
|
||||||
|
|||||||
@@ -13,6 +13,26 @@
|
|||||||
складывается фраза — это вопрос владельцу, а не строчка кода: проверки в этом
|
складывается фраза — это вопрос владельцу, а не строчка кода: проверки в этом
|
||||||
репозитории однажды уже переросли продукт.
|
репозитории однажды уже переросли продукт.
|
||||||
|
|
||||||
|
**Усложнение не бесплатно, и платит за него учебная ценность.** Чем сложнее
|
||||||
|
стенд, тем хуже он как учебный материал: внимание менти конечно, и каждая
|
||||||
|
конструкция, которую он обязан расшифровать по дороге к уроку, списывается с
|
||||||
|
этого счёта. Поэтому вопрос выше имеет вторую половину: **что менти платит,
|
||||||
|
чтобы это прочесть, и что получает взамен?** Ответ «вопрос — и ничего»
|
||||||
|
означает, что писать не нужно.
|
||||||
|
|
||||||
|
Оборонительный код дороже прочего. Читатель не отличает «стоит здесь, потому
|
||||||
|
что есть настоящая опасность» от «стоит здесь, потому что кто-то был
|
||||||
|
остроумен», — и предполагает первое. Лишний сторож не безобиден: он
|
||||||
|
утверждает, будто опасность заслуживает внимания, и этим врёт. Особенно это
|
||||||
|
касается защиты от ошибки, которую человек делает сам себе и тут же видит:
|
||||||
|
разбирательство с ней и есть урок, отнимать его не надо.
|
||||||
|
|
||||||
|
Ревью само по себе умеет только прибавлять: каждая находка просится в правку,
|
||||||
|
и за прогон их набираются десятки. Поэтому перед приёмкой заметной работы
|
||||||
|
идёт отдельный **проход на вычитание** — с правом только резать и с явным
|
||||||
|
списком неприкосновенного. Сомнение «резать или нет» решается в пользу
|
||||||
|
реза: вещь, которая себя не защитила, себя не защитила.
|
||||||
|
|
||||||
## Кому что поручать
|
## Кому что поручать
|
||||||
|
|
||||||
Читатель здесь не побочный потребитель, а тот, ради кого стенд существует:
|
Читатель здесь не побочный потребитель, а тот, ради кого стенд существует:
|
||||||
|
|||||||
@@ -1,6 +1,9 @@
|
|||||||
COMPOSE ?= docker compose
|
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:
|
up:
|
||||||
$(COMPOSE) up --detach --build --wait --wait-timeout 600
|
$(COMPOSE) up --detach --build --wait --wait-timeout 600
|
||||||
@@ -17,6 +20,14 @@ ps:
|
|||||||
logs:
|
logs:
|
||||||
$(COMPOSE) logs --follow
|
$(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:
|
config-test:
|
||||||
COMPOSE_BIN="$(COMPOSE)" ./scripts/config-test.sh
|
COMPOSE_BIN="$(COMPOSE)" ./scripts/config-test.sh
|
||||||
|
|
||||||
|
|||||||
@@ -9,9 +9,11 @@
|
|||||||
[«Боевой реализм стенда (v2)»](docs/specs/2026-07-30-stand-v2-realism.md).
|
[«Боевой реализм стенда (v2)»](docs/specs/2026-07-30-stand-v2-realism.md).
|
||||||
Сейчас работают кластер ClickHouse из двух шардов, отдельный
|
Сейчас работают кластер ClickHouse из двух шардов, отдельный
|
||||||
clickhouse-keeper, односерверная Kafka в режиме KRaft, Airflow 3.3, Superset,
|
clickhouse-keeper, односерверная Kafka в режиме KRaft, Airflow 3.3, Superset,
|
||||||
Prometheus, Grafana и общая база Postgres для метаданных. Начат генератор
|
Prometheus, Grafana и общая база Postgres для метаданных. Генератор
|
||||||
кликстрима: в [`generator/`](generator/) заведён контракт схемы события, из
|
кликстрима в [`generator/`](generator/) собран целиком со стороны клиента: у
|
||||||
которого собрано [описание выгрузки](docs/formats/clickstream-event.md).
|
него есть контракт схемы события, из которого собрано
|
||||||
|
[описание выгрузки](docs/formats/clickstream-event.md), модельный мир и
|
||||||
|
проигрыватель, отправляющий дни в Kafka или в файл.
|
||||||
|
|
||||||
## Быстрый старт
|
## Быстрый старт
|
||||||
|
|
||||||
@@ -117,6 +119,70 @@ Kafka по той же причине спрашивают снаружи. Её
|
|||||||
полного сброса с удалением всех именованных томов используйте `make clean`.
|
полного сброса с удалением всех именованных томов используйте `make clean`.
|
||||||
Повторный `make up` безопасен: одноразовая подготовка приложений идемпотентна.
|
Повторный `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;
|
- `clickhouse-01` — инициатор DDL и точка подключения Airflow;
|
||||||
|
|||||||
@@ -194,6 +194,25 @@ services:
|
|||||||
exit 1
|
exit 1
|
||||||
fi
|
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 по порядку имён файлов после готовности всего кластера.
|
# DDL применяется с ноды 1 по порядку имён файлов после готовности всего кластера.
|
||||||
clickhouse-init:
|
clickhouse-init:
|
||||||
image: ${CLICKHOUSE_IMAGE:-clickhouse/clickhouse-server:26.3.17.56}
|
image: ${CLICKHOUSE_IMAGE:-clickhouse/clickhouse-server:26.3.17.56}
|
||||||
|
|||||||
@@ -259,8 +259,7 @@
|
|||||||
отдельным сообщением. Контракт транспорта, не деталь реализации — сторона
|
отдельным сообщением. Контракт транспорта, не деталь реализации — сторона
|
||||||
хранилища читает топик байтами и кладёт сообщение строкой
|
хранилища читает топик байтами и кладёт сообщение строкой
|
||||||
([ADR 0005](../adr/0005-event-ingestion.md)), поэтому склейка нескольких
|
([ADR 0005](../adr/0005-event-ingestion.md)), поэтому склейка нескольких
|
||||||
событий в одно сообщение сломала бы разбор целиком. Сторожится тестом
|
событий в одно сообщение сломала бы разбор целиком.
|
||||||
приёмника.
|
|
||||||
- **Даты и время на проводе — ISO-8601.** `EventDate` уезжает как `2026-06-01`,
|
- **Даты и время на проводе — ISO-8601.** `EventDate` уезжает как `2026-06-01`,
|
||||||
`UTCEventTime` — как `2026-06-01T12:34:56Z`. Довод — читаемость сырья: весь
|
`UTCEventTime` — как `2026-06-01T12:34:56Z`. Довод — читаемость сырья: весь
|
||||||
смысл слоя STG в том, что менти открывает колонку `raw` в обычном клиенте и
|
смысл слоя STG в том, что менти открывает колонку `raw` в обычном клиенте и
|
||||||
@@ -269,9 +268,10 @@
|
|||||||
принимает ISO без плясок. Колонка `ecommerce` — строка, внутри которой лежит
|
принимает ISO без плясок. Колонка `ecommerce` — строка, внутри которой лежит
|
||||||
экранированный JSON, как отдаёт Метрика.
|
экранированный JSON, как отдаёт Метрика.
|
||||||
|
|
||||||
Форму пинит хранилище (#43) как первый потребитель, сериализатор (#41)
|
Форму реализует сериализатор (#41), хранилище (#43) читает то, что он
|
||||||
её соблюдает. Порядок тикетов обратный порядку зависимости, поэтому здесь она
|
положил: порядок тикетов развёрнут 6 августа 2026 года, и отправитель идёт
|
||||||
и записана — иначе каждый выберет своё, и разойдётся это уже после приёмки.
|
первым. Записана форма всё равно здесь — иначе каждый выберет своё, и
|
||||||
|
разойдётся это уже после приёмки.
|
||||||
- **Рабочий выбор сериализатора — orjson**: быстрее stdlib json в 5–14 раз,
|
- **Рабочий выбор сериализатора — orjson**: быстрее stdlib json в 5–14 раз,
|
||||||
numpy-массивы и datetime сериализует нативно (заметка исследования #31).
|
numpy-массивы и datetime сериализует нативно (заметка исследования #31).
|
||||||
Смена библиотеки меняет канонические байты, поэтому проходит как
|
Смена библиотеки меняет канонические байты, поэтому проходит как
|
||||||
@@ -428,6 +428,16 @@ pytest-тест с маркером `perf` и таймаутом-обрубан
|
|||||||
теперь есть.
|
теперь есть.
|
||||||
- **Этап 5 (Airflow)**: живой день со стороны хранилища — обычный ETL-даг
|
- **Этап 5 (Airflow)**: живой день со стороны хранилища — обычный ETL-даг
|
||||||
по расписанию (~раз в 24 минуты); генератор не дорабатывается.
|
по расписанию (~раз в 24 минуты); генератор не дорабатывается.
|
||||||
|
- **Этап 5, позиция на оси времени.** Проигрыватель состояния не хранит
|
||||||
|
(раздел 9), поэтому вести позицию — работа того, кто его зовёт. Живёт она
|
||||||
|
переменной Airflow: отсутствие переменной означает мир в зерновом
|
||||||
|
состоянии, заводит и двигает её только даг `next_day` и только по успеху.
|
||||||
|
Довод — генератор отдельная и заменяемая сущность, привязывать его к
|
||||||
|
хранилищу незачем, а `make clean` сносит том метаданных Airflow вместе с
|
||||||
|
данными ClickHouse, так что позиция и мир чистятся одной командой.
|
||||||
|
Расхождение переменной с данными (менти почистил партицию руками) лечится
|
||||||
|
документацией, а не сторожем: постоянная сверка была бы проверкой, красной
|
||||||
|
в норме ([ADR 0002](../adr/0002-monitoring-scope.md)).
|
||||||
- **Этап 7 (эталонный мир)**: пересборка — это манифест, не артефакт;
|
- **Этап 7 (эталонный мир)**: пересборка — это манифест, не артефакт;
|
||||||
CI-генерация на amd64 и arm64.
|
CI-генерация на amd64 и arm64.
|
||||||
- **Будущие лабы**: перезаливка дня X пакетным режимом проигрывателя —
|
- **Будущие лабы**: перезаливка дня 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 мастер-спеки):
|
- как фиксируется «зерновой» мир конца этапа 2 (раздел 9 мастер-спеки):
|
||||||
с манифестным решением напрашивается мини-манифест зернового мира —
|
с манифестным решением напрашивается мини-манифест зернового мира —
|
||||||
форма за #42;
|
форма за #42;
|
||||||
|
|||||||
@@ -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
@@ -4,9 +4,11 @@
|
|||||||
образцу облачной выгрузки Яндекс Метрики. Устройство и принятые решения —
|
образцу облачной выгрузки Яндекс Метрики. Устройство и принятые решения —
|
||||||
спека [«Генератор (этап 2)»](../docs/specs/2026-08-01-generator.md).
|
спека [«Генератор (этап 2)»](../docs/specs/2026-08-01-generator.md).
|
||||||
|
|
||||||
События уже есть: день-функция отдаёт по паре (зерно, D) упорядоченный поток
|
События уже есть и уже уезжают: день-функция отдаёт по паре (зерно, D)
|
||||||
трёх видов — просмотр страницы, корзина, покупка. Клиентская сторона на этом
|
упорядоченный поток трёх видов — просмотр страницы, корзина, покупка, —
|
||||||
целая; заказы бэкенда и запуск снаружи — за следующими тикетами.
|
канонический сериализатор превращает его в JSON, а проигрыватель гонит в файл
|
||||||
|
или в Kafka. Клиентская сторона на этом целая; заказы бэкенда — за следующими
|
||||||
|
этапами.
|
||||||
|
|
||||||
## Как это работает
|
## Как это работает
|
||||||
|
|
||||||
@@ -60,6 +62,16 @@ D0 живёт предыстория, поэтому любой день соб
|
|||||||
Своя случайность, поэтому правка торговли трафик не двигает.
|
Своя случайность, поэтому правка торговли трафик не двигает.
|
||||||
- `src/clickstream_generator/ids.py` — номера событий: неповторяющиеся и
|
- `src/clickstream_generator/ids.py` — номера событий: неповторяющиеся и
|
||||||
ниже 2^53. Обещание одно на обе половины дня, поэтому и живёт отдельно.
|
ниже 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` — контракт схемы: чистые данные о
|
- `src/clickstream_generator/schema.py` — контракт схемы: чистые данные о
|
||||||
колонках выгрузки. Собственность генератора; из него выводятся сам
|
колонках выгрузки. Собственность генератора; из него выводятся сам
|
||||||
генератор, его валидация и описание выгрузки в доках.
|
генератор, его валидация и описание выгрузки в доках.
|
||||||
@@ -68,10 +80,14 @@ D0 живёт предыстория, поэтому любой день соб
|
|||||||
из контракта. Документ руками не правят — пересобирают.
|
из контракта. Документ руками не правят — пересобирают.
|
||||||
- `tests/` — инварианты контракта, свежесть описания и обещания мира:
|
- `tests/` — инварианты контракта, свежесть описания и обещания мира:
|
||||||
чистота от зерна, приток, гарантия двухкуковых пар, форма суточной волны
|
чистота от зерна, приток, гарантия двухкуковых пар, форма суточной волны
|
||||||
и сборка визитов по задокументированным правилам.
|
и сборка визитов по задокументированным правилам. Там же побайтовое
|
||||||
|
обещание, доведённое до диска: два прогона дня в файл дают тот же файл, а
|
||||||
|
строк в нём ровно столько, сколько событий.
|
||||||
|
|
||||||
## Команды
|
## Команды
|
||||||
|
|
||||||
|
Проиграть день — [быстрый старт](../README.md#как-позвать-генератор) и `--help`.
|
||||||
|
|
||||||
Из корня репозитория:
|
Из корня репозитория:
|
||||||
|
|
||||||
- `make test` — тесты генератора;
|
- `make test` — тесты генератора;
|
||||||
|
|||||||
@@ -11,9 +11,18 @@ dependencies = [
|
|||||||
"orjson>=3.11.9",
|
"orjson>=3.11.9",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[project.optional-dependencies]
|
||||||
|
# Клиент Kafka — необязательная часть пакета: без него генератор пишет в файл
|
||||||
|
# и остаётся переносимым (спека генератора, раздел 9). Приёмник импортирует
|
||||||
|
# клиента внутри функции, поэтому пакет без этой группы живой.
|
||||||
|
kafka = ["confluent-kafka>=2.6"]
|
||||||
|
|
||||||
[dependency-groups]
|
[dependency-groups]
|
||||||
dev = [
|
dev = [
|
||||||
"pytest>=8",
|
"pytest>=8",
|
||||||
|
# Тот же клиент нужен и на машине разработчика: иначе ty не разберёт
|
||||||
|
# импорт приёмника, а ruff и тесты видели бы половину кода.
|
||||||
|
"confluent-kafka>=2.6",
|
||||||
"ruff>=0.15",
|
"ruff>=0.15",
|
||||||
# Типчекер; молод (0.0.x), но версия закреплена в uv.lock. Настройка
|
# Типчекер; молод (0.0.x), но версия закреплена в uv.lock. Настройка
|
||||||
# сверена с доками Astral через Context7: запуск `uv run ty check`,
|
# сверена с доками Astral через Context7: запуск `uv run ty check`,
|
||||||
|
|||||||
@@ -0,0 +1,9 @@
|
|||||||
|
"""Точка входа пакета: `python -m clickstream_generator`.
|
||||||
|
|
||||||
|
Код возврата уходит наружу как есть — по нему судит о прогоне и контейнер, и
|
||||||
|
даг, который его запустит (спека генератора, раздел 9).
|
||||||
|
"""
|
||||||
|
|
||||||
|
from clickstream_generator.cli import main
|
||||||
|
|
||||||
|
raise SystemExit(main())
|
||||||
@@ -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}")
|
||||||
@@ -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()
|
||||||
@@ -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]
|
||||||
Generated
+22
@@ -11,8 +11,14 @@ dependencies = [
|
|||||||
{ name = "orjson" },
|
{ name = "orjson" },
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[package.optional-dependencies]
|
||||||
|
kafka = [
|
||||||
|
{ name = "confluent-kafka" },
|
||||||
|
]
|
||||||
|
|
||||||
[package.dev-dependencies]
|
[package.dev-dependencies]
|
||||||
dev = [
|
dev = [
|
||||||
|
{ name = "confluent-kafka" },
|
||||||
{ name = "pytest" },
|
{ name = "pytest" },
|
||||||
{ name = "ruff" },
|
{ name = "ruff" },
|
||||||
{ name = "ty" },
|
{ name = "ty" },
|
||||||
@@ -20,12 +26,15 @@ dev = [
|
|||||||
|
|
||||||
[package.metadata]
|
[package.metadata]
|
||||||
requires-dist = [
|
requires-dist = [
|
||||||
|
{ name = "confluent-kafka", marker = "extra == 'kafka'", specifier = ">=2.6" },
|
||||||
{ name = "numpy", specifier = ">=2" },
|
{ name = "numpy", specifier = ">=2" },
|
||||||
{ name = "orjson", specifier = ">=3.11.9" },
|
{ name = "orjson", specifier = ">=3.11.9" },
|
||||||
]
|
]
|
||||||
|
provides-extras = ["kafka"]
|
||||||
|
|
||||||
[package.metadata.requires-dev]
|
[package.metadata.requires-dev]
|
||||||
dev = [
|
dev = [
|
||||||
|
{ name = "confluent-kafka", specifier = ">=2.6" },
|
||||||
{ name = "pytest", specifier = ">=8" },
|
{ name = "pytest", specifier = ">=8" },
|
||||||
{ name = "ruff", specifier = ">=0.15" },
|
{ name = "ruff", specifier = ">=0.15" },
|
||||||
{ name = "ty", specifier = ">=0.0.65" },
|
{ 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" },
|
{ 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]]
|
[[package]]
|
||||||
name = "iniconfig"
|
name = "iniconfig"
|
||||||
version = "2.3.0"
|
version = "2.3.0"
|
||||||
|
|||||||
@@ -13,6 +13,12 @@ trap cleanup EXIT
|
|||||||
|
|
||||||
"${COMPOSE_CMD[@]}" --project-directory "$ROOT_DIR" config --quiet
|
"${COMPOSE_CMD[@]}" --project-directory "$ROOT_DIR" config --quiet
|
||||||
config_json="$("${COMPOSE_CMD[@]}" --project-directory "$ROOT_DIR" config --format json)"
|
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"
|
grep -qx 'ARG SUPERSET_BASE_IMAGE' "$ROOT_DIR/infra/superset/Dockerfile"
|
||||||
jq -e '
|
jq -e '
|
||||||
.services.superset.image == "clickstream-superset:local" and
|
.services.superset.image == "clickstream-superset:local" and
|
||||||
|
|||||||
Reference in New Issue
Block a user