diff --git a/CONTEXT.md b/CONTEXT.md index daef449..d18bde1 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -95,6 +95,19 @@ ClickHouse. Форма файла решена, длина — нет: стро первого дня D0; реальный календарь в модели не участвует. Между прогонами живут только зерно и позиция на оси. +**Стартовый мир**: +Первые восемь дней оси (понедельник по понедельник), которыми стенд +наполняется при каждом `make up`. Маленький кусок эталонного мира, всегда +один и тот же: на нём принимаются следующие этапы. +_Избегать_: зерновой мир + +**Опись мира**: +`data/world-inventory.json` — единственное, что о мире хранится в git: +паспорт (зерно, версия генератора, хеш каталога) и по строке на день с +датой, числом событий и хешем его байтов. Сам мир в git не лежит — он +пересчитывается. Опись отвечает на один вопрос: тот ли это мир. +_Избегать_: манифест, мини-манифест + **Пошаговый режим**: Базовый способ движения по оси модельного времени: «прожить следующий день» — явное действие. diff --git a/Makefile b/Makefile index 48c96e9..03a5c18 100644 --- a/Makefile +++ b/Makefile @@ -3,10 +3,13 @@ GENERATOR_DAY ?= 0 GENERATOR_LIMIT ?= GENERATOR_SPEED ?= -.PHONY: up down clean ps logs generate-batch generate-live 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 inventory smoke check-clickhouse check-services +# Второй шаг: `--wait` дожидается служб, а приём событий асинхронный — довод +# целиком в шапке скрипта. up: $(COMPOSE) up --detach --build --wait --wait-timeout 600 + COMPOSE_BIN="$(COMPOSE)" ./scripts/wait-for-world.sh down: $(COMPOSE) down --remove-orphans @@ -44,6 +47,10 @@ docs: cd generator && uv run python -m clickstream_generator.schema_doc \ ../docs/formats/clickstream-event.md +inventory: + cd generator && uv run python -m clickstream_generator.inventory \ + ../data/world-inventory.json + smoke: COMPOSE_BIN="$(COMPOSE)" ./scripts/stand-smoke.sh diff --git a/README.md b/README.md index 5ef3fce..eaeb016 100644 --- a/README.md +++ b/README.md @@ -13,7 +13,9 @@ Prometheus, Grafana и общая база Postgres для метаданных. кликстрима в [`generator/`](generator/) собран целиком со стороны клиента: у него есть контракт схемы события, из которого собрано [описание выгрузки](docs/formats/clickstream-event.md), модельный мир и -проигрыватель, отправляющий дни в Kafka или в файл. +проигрыватель, отправляющий дни в Kafka или в файл. События доезжают до +типизированного `ods.event`, а `make up` наполняет стенд стартовым миром — +первой неделей модельного времени. ## Быстрый старт @@ -26,10 +28,12 @@ Prometheus, Grafana и общая база Postgres для метаданных. хватает, стенд одинаково хорошо живёт на недорогом VPS. Проверкам на поднятом стенде также нужны `curl`, `jq`, `awk`, -`grep`, `sed`, `tail`, `sleep` и `timeout`. По умолчанию должны быть свободны +`grep`, `sed`, `tail`, `sleep` и `timeout`. `jq` нужен и самому `make up`: им +читается опись мира, по которой он ждёт заливки. По умолчанию должны быть свободны порты `23000`, `28080`, `28088`, `28123`, `28124`, `29000`, `29001`, `29090` и `29092`. Проверкам без стенда — `make config-test`, `make lint`, -`make typecheck`, `make test` — и сборке документации `make docs` нужен `uv`. +`make typecheck`, `make test` — и сборкам `make docs` и `make inventory` нужен +`uv`. Стенд запускается без `.env`: @@ -70,6 +74,7 @@ make smoke `superset-init` завершаются с кодом 0. Первый обновляет схему Airflow, подготавливает администратора и подключение к `clickhouse-01`. Второй обновляет Superset, создаёт администратора и импортирует подключение к `clickhouse-02`. +С нуля подъём занимает около трёх минут, на живом стенде — около минуты. `make smoke` за секунды спрашивает, собран ли стенд: зависимости машины, здоровье контейнеров, устройство keeper, ответ Kafka с машины через отображённый @@ -105,7 +110,10 @@ Kafka по той же причине спрашивают снаружи. Её `make check-clickhouse` запускает отдельную глубокую проверку ClickHouse: описание кластера, макросы, связь с keeper, `ReplicatedMergeTree`, -`Distributed`, очередь распределённых DDL и очистку временных таблиц. +`Distributed`, очередь распределённых DDL и очистку временных таблиц. Последняя +из девяти проверок — единственная на настоящих данных: она подневно сверяет +события стартового мира с описью и при расхождении говорит, где искать — +в событиях или в браке. `make config-test` проверяет Compose, синтаксис файлов DAG и пробельные ошибки в diff без запуска стенда. @@ -119,6 +127,38 @@ Kafka по той же причине спрашивают снаружи. Её полного сброса с удалением всех именованных томов используйте `make clean`. Повторный `make up` безопасен: одноразовая подготовка приложений идемпотентна. +## Стартовый мир + +Стенд поднимается не пустым: разовая служба `world-init` играет в топик `hits` +первые восемь дней модельного времени — понедельник по понедельник, 401 185 +событий. Дальше их обычным путём разбирает хранилище, и к концу `make up` они +лежат в `ods.event`. Так у всякой лабы есть данные, и всегда одни и те же. + +Ждать приходится дольше, чем работает заливка: приём асинхронный, поэтому +вторым шагом `make up` зовёт `scripts/wait-for-world.sh` — тот опрашивает +ClickHouse, пока мир не доедет. Повторный `make up` заливает мир заново; это +не ошибка, а свойство: номера событий те же, и повтор схлопнет +`ReplacingMergeTree`. + +Сам мир в git не хранится — он чистая функция зерна, и держать его в +репозитории значило бы держать там кэш. Вместо него лежит **опись мира**, +[`data/world-inventory.json`](data/world-inventory.json): зерно, версия +генератора, хеш каталога товаров и по строке на каждый день — дата, число +событий и хеш его байтов. Опись отвечает на единственный вопрос: тот ли это +мир, что был вчера. + +Спрашивают её двое. `make test` сверяет опись с тем, что собирается из кода +сегодня: правка генератора меняет мир, и опись надо пересобрать — +`make inventory`. `make check-clickhouse` сверяет с описью то, что доехало до +`ods.event`, подневно. Хеш дня можно пересчитать и руками — это обычный +`sha256sum` файла, который пишет файловый приёмник: + +```bash +uv run --project generator python -m clickstream_generator batch \ + --day 0 --file tmp/day0.jsonl +sha256sum tmp/day0.jsonl +``` + ## Как позвать генератор События производит генератор из [`generator/`](generator/). Модельный день — @@ -160,6 +200,10 @@ make generate-live GENERATOR_DAY=3 GENERATOR_SPEED=1000 docker compose --profile generator build generator ``` +Тот же образ несёт заливка стартового мира, поэтому свежим его держит и +обычный `make up`: он собирает образы всего стенда, и генератор теперь среди +них. + Разведено это нарочно: собрать образ и запустить контейнер — разные действия, и цель запуска, молча пересобирающая образ, стирает между ними границу. Что образ устарел, видно по собственному прогону — это обратная связь, а не ловушка. diff --git a/compose.yaml b/compose.yaml index c361ff8..a702018 100644 --- a/compose.yaml +++ b/compose.yaml @@ -52,6 +52,28 @@ x-airflow-common: &airflow-common - airflow_logs:/opt/airflow/logs - airflow_auth:/opt/airflow/auth +x-generator-common: &generator-common + image: clickstream-generator:local + build: + context: . + dockerfile: generator/Dockerfile + restart: "no" + depends_on: + kafka-init: + condition: service_completed_successfully + # Читать топик движок Kafka начинает не при своём создании, а когда над ним + # появляется матвью приёма, и она создаётся последней по номеру файла + # (sql/ddl/40-stg-views.sql). Отправь генератор события раньше — они лягут в + # топик и молча минуют хранилище. + clickhouse-init: + condition: service_completed_successfully + environment: + # Адрес брокера и имя топика — факты стенда, и называет их стенд. + # Остальное (день, зерно, число дней, предел пачки) приходит аргументами + # от того, кто запускает: у службы нет позиции на оси мира. + KAFKA_BOOTSTRAP_SERVERS: ${KAFKA_BOOTSTRAP_SERVERS:-kafka:9092} + KAFKA_TOPIC: ${KAFKA_TOPIC:-hits} + x-superset-common: &superset-common image: clickstream-superset:local build: @@ -194,25 +216,6 @@ 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} @@ -244,6 +247,27 @@ services: volumes: - ./sql/ddl:/ddl:ro + # Стартовый мир: восемь модельных дней уезжают в топик при каждом подъёме. + # Что это за дни и каким мир обязан выйти — data/world-inventory.json. + # + # Число дней стоит здесь числом: YAML не читает Python, и одно из двух мест + # (второе — `STARTING_DAYS` в inventory.py) лишнее по построению. Разъедутся + # они — покраснеют счётчики `make check-clickhouse`. + # + # Повторный `make up` заливает мир заново, и это не оплошность: `WatchID` у + # событий те же, ReplacingMergeTree схлопнет повтор в ODS. Сырьё в STG при + # этом честно удвоится — свойство слоя, описанное в storage.md. + world-init: + <<: *generator-common + command: ["batch", "--day", "0", "--days", "8"] + + # Тот же образ для ручных прогонов: `make generate-batch`, `make generate-live`. + # Под профилем — чтобы обычный подъём стенда её не трогал. + generator: + <<: *generator-common + profiles: ["generator"] + command: ["batch"] + postgres-metadata: image: ${POSTGRES_IMAGE:-postgres:16-alpine} restart: unless-stopped @@ -285,6 +309,13 @@ services: # упавшим для --wait. clickhouse-init: condition: service_completed_successfully + # По той же причине — и ещё по одной. Заливка стартового мира идёт + # последней в цепи разовых служб, и зависимого ей взять больше негде. + # Заодно это правда про порядок: даг `next_day` этапа 5 продолжает ось с + # того дня, на котором заливка остановилась, — Airflow приходит в мир, + # который уже есть. + world-init: + condition: service_completed_successfully entrypoint: ["/bin/bash"] command: ["/opt/airflow/init.sh"] diff --git a/data/world-inventory.json b/data/world-inventory.json new file mode 100644 index 0000000..07bb7f8 --- /dev/null +++ b/data/world-inventory.json @@ -0,0 +1,55 @@ +{ + "seed": 20260601, + "generator_version": "0.1.0", + "catalog_sha256": "4fffceb3ec278c0b364ec84d936d991946ddd2a0912ed9c36471ff7a55c7483f", + "days": [ + { + "day": 0, + "date": "2026-06-01", + "events": 50626, + "sha256": "81ad7f402bdd4c8b56425f6439a024b9996e5ca95dd207d444290d20a70ed407" + }, + { + "day": 1, + "date": "2026-06-02", + "events": 49722, + "sha256": "5f7aa0301832008db612eff05899b1210af958c888ac2e9e89f4b7a552569eff" + }, + { + "day": 2, + "date": "2026-06-03", + "events": 52361, + "sha256": "10ac9ce65ca6dc171689edf78ef5cd785e8c8a80e54d236bf2b9d6dbf6997070" + }, + { + "day": 3, + "date": "2026-06-04", + "events": 53157, + "sha256": "7bf4deadd8f69a6c479fceabfe2d1e8ea8db76919b10a8b9e780006708da31c3" + }, + { + "day": 4, + "date": "2026-06-05", + "events": 50078, + "sha256": "b9e758fb72f937aa5ad1f2e90dee5ec147b16cfd89e49e51a00582f909458c29" + }, + { + "day": 5, + "date": "2026-06-06", + "events": 47164, + "sha256": "f08a9f9a439f1da9e4ad8d94b9f65b7f3329256e41446d22d900cc4825f2480c" + }, + { + "day": 6, + "date": "2026-06-07", + "events": 48683, + "sha256": "7e45d9dd59c902067fa6847d813669b1ab9a4ee88d2f192382577a5ed8d590a7" + }, + { + "day": 7, + "date": "2026-06-08", + "events": 49394, + "sha256": "286ccc2451847ac8de215b5e902c328371021ad5502bd4afca0494b6f313aef7" + } + ] +} diff --git a/docs/architecture/storage.md b/docs/architecture/storage.md index 867a293..80138c4 100644 --- a/docs/architecture/storage.md +++ b/docs/architecture/storage.md @@ -227,10 +227,10 @@ Airflow она стоит — там это обычный `SETTINGS` у зап роняет вставку и останавливает потребление до починки ([ADR 0005](../adr/0005-event-ingestion.md)). -Постоянной сверки счётчиков при этом нет и не должно быть. У сырья срок жизни -трое суток, а ODS хранит всё, поэтому равенство «сырьё = события + ошибки» -разъедется само: сырьё истечёт раньше, а перезаливка модельного дня задвоит его, -тогда как в ODS тот же повтор схлопнется. Правило, красное в норме, учит не +Постоянной сверки **слоёв между собой** при этом нет и не должно быть. У сырья +срок жизни трое суток, а ODS хранит всё, поэтому равенство «сырьё = события + +ошибки» разъедется само: сырьё истечёт раньше, а перезаливка модельного дня +задвоит его, тогда как в ODS тот же повтор схлопнется. Правило, красное в норме, учит не смотреть на оповещения ([ADR 0002](../adr/0002-monitoring-scope.md)). Равенство проверено разовым опытом при исполнении #43, на управляемой пачке: отправили N сообщений — @@ -241,6 +241,15 @@ Airflow она стоит — там это обычный `SETTINGS` у зап `ReplacingMergeTree` зависит от того, сколько мержей успело пройти, и спека это прямо запрещает (раздел 6). +**Постоянная сверка одна, и она другого рода** — не слой против слоя, а ODS +против [описи мира](../../data/world-inventory.json), лежащей в git +(`make check-clickhouse`, пришла с #42). Разъехаться сама она не может, и в +этом вся разница: опись описывает восемь дней стартового мира, счёт обрамлён +их датами, повтор заливки схлопывается под `FINAL`, а срок хранения ей не +помеха — ODS хранит всё. Обе стороны равенства зафиксированы: одна кодом +генератора, другая файлом в git. У межслойной сверки такой опоры нет ни с +одной стороны. + ## Срок жизни сырья Сырьё в STG живёт трое суток реального времени и уходит само. Трое — это окно diff --git a/docs/architecture/testing.md b/docs/architecture/testing.md index 307f973..9640668 100644 --- a/docs/architecture/testing.md +++ b/docs/architecture/testing.md @@ -23,25 +23,29 @@ заметит потерянного подключения Superset: он про ClickHouse и только. Отсюда правило для новой проверки: **спроси, кого она спрашивает.** Счётчики -против манифеста — вопрос к ClickHouse: строки в `ods.event` считает сам +против описи — вопрос к ClickHouse: строки в `ods.event` считает сам сервер и отвечает сразу, значит дом им в `check-clickhouse`, даже если по цене они подошли бы смоуку. ## Карта целей -Стенд нужен трём целям из семи. Цена — замер 6 августа 2026 года, см. «Что +Стенд нужен трём целям из семи. Цена — замер 7 августа 2026 года, см. «Что проверено». | Цель | Что утверждает | Стенд | Цена | |---|---|---|---| | `make lint` | Код генератора отформатирован и проходит ruff | не нужен | 0,4 с | | `make typecheck` | Типы генератора сходятся (ty) | не нужен | 0,5 с | -| `make test` | Генератор делает то, что обещает; схема события остаётся объявленным контрактом, а собранное из неё [описание выгрузки](../formats/clickstream-event.md) — свежим | не нужен | 40 с | +| `make test` | Генератор делает то, что обещает; схема события остаётся объявленным контрактом, а собранные из кода [описание выгрузки](../formats/clickstream-event.md) и [опись мира](../../data/world-inventory.json) — свежими | не нужен | 71 с | | `make config-test` | Compose разбирается, файлы DAG синтаксически целы, в diff нет пробельных ошибок. О работоспособности не говорит ничего | не нужен | 1 с | -| `make smoke` | Стенд **собран**: службы живы, порты отвечают, подключения настроены друг на друга. Вширь и по касательной к каждой службе. Единственная цель, которая здесь правда смоук | нужен | 8 с | -| `make check-clickhouse` | Всё, что спрашивают **у ClickHouse** и он отвечает сам: макросы, шарды, реплики, путь в keeper, ключ шардирования, очередь распределённых DDL | нужен | 7 с | +| `make smoke` | Стенд **собран**: службы живы, порты отвечают, подключения настроены друг на друга. Вширь и по касательной к каждой службе. Единственная цель, которая здесь правда смоук | нужен | 9 с | +| `make check-clickhouse` | Всё, что спрашивают **у ClickHouse** и он отвечает сам: макросы, шарды, реплики, путь в keeper, ключ шардирования, очередь распределённых DDL, счёт событий стартового мира против описи | нужен | 8 с | | `make check-services` | **Службы работают**: DAG запускается и доходит, топик создаётся и удаляется, Superset логинится и ходит в базу | нужен | 44 с | +Цена самого `make up` — 2 м 50 с с нуля (`make clean` перед ним — ещё 23 с) и +1 м 5 с на живом стенде. В цену целей она не входит, но записана здесь по той +же причине: с #42 подъём стенда заливает стартовый мир и стал заметно дороже. + ## Правило: смоук обязан оставаться быстрым `make smoke` — быстрая проверка для регулярного прогона: её гоняют не @@ -72,6 +76,14 @@ ClickHouse отвечает сразу. Правки генератора добавляют к этому `make lint`, `make typecheck` и `make test`: стенд им не нужен, а `make test` из них самая дорогая. +Тринадцать секунд из её семидесяти одной — пересборка восьми модельных дней +для сверки с описью мира. Дешевле хеши не сравнить: чтобы узнать, тот ли +получается мир, его надо получить. Место выбрано по той же оси «кого +спрашивают» — вопрос обращён к коду генератора, стенд ему не нужен, — и +краснеет проверка там, где надо: сразу после правки генератора. Не будь её, +расхождение всплыло бы получасом позже, на поднятом стенде, где выглядит +поломкой хранилища, а не забытой пересборкой. + ## Порогов по времени здесь нет Цена в таблице — замеренное число с датой замера, а не назначенный порог. @@ -117,9 +129,16 @@ ClickHouse отвечает сразу. же механизм, на котором стоит переобработка дня X. Что тогда стережёт цепочку постоянно — проверки на **настоящих** данных, тех, -что стенд произвёл сам: счётчики против манифеста зернового мира ловят -сломанный разбор (события уехали в брак — счёт разошёлся), ничего при этом не -вкладывая. +что стенд произвёл сам: подневный счёт событий против описи мира ловит и +сломанный разбор (события уехали в брак — счёт разошёлся), и потерю по дороге, +ничего при этом не вкладывая. Это единственный постоянный сторож цепочки +Kafka → STG → ODS, и живёт он в `make check-clickhouse`. + +Правило, которое такая проверка обязана выдержать и которое стоит держать в +голове для следующей: **на растущих данных утверждать можно только про +зафиксированный кусок.** Мир растёт — менти переиграет день, этап 5 добавит +следующий, — и всякое утверждение про таблицу целиком однажды покраснеет +законно, то есть впустую. Как это сделано здесь, написано в самом скрипте. ## Корректность процессов живёт в дагах DQ, а не в целях `make` @@ -145,6 +164,13 @@ ClickHouse отвечает сразу. | `make check-services` | `scripts/stand-services.sh` | | `make check-clickhouse` | `scripts/check-clickhouse.sh` | +Особняком — `scripts/wait-for-world.sh`: он не проверка, а вторая половина +`make up`. `docker compose up --wait` дожидается служб, в том числе успешно +отработавшей заливки, но «заливка кончилась» значит «события в топике», а не +«события в хранилище»: приём асинхронный. Скрипт ждёт доезда ограниченным +циклом опроса. Без него `make check-clickhouse` следом краснел бы по +устройству, а не по поломке. + Общее у смоука и `check-services` — счёт проверок, обращение к Compose и две проверки — вынесено в `scripts/stand-common.sh`; сам он не запускается. Оттуда же приходят два решения, которые видно по счёту прогонов: @@ -165,11 +191,29 @@ ClickHouse отвечает сразу. ## Что проверено -Замеры 6 августа 2026 года, стенд поднят заранее; время `make up` в цену целей -не входит. Время взято по `time` и совпадает с тем, что цель печатает сама. Оно -плавает от прогона к прогону: смоук дал 8 секунд дважды (без проверки Kafka, -до её появления, — 5 и 6), `check-services` — 43 и 44. В таблице стоит большее -из замеренных. +**Перезамер 7 августа 2026 года при исполнении #42.** Подъём стенда с нуля +(`make clean && make up`) — 2 м 50 с, из них 22,5 с занимает сама заливка: +18,7 с генерация восьми дней и 3,8 с доставка 401 185 сообщений в Kafka. +Повторный `make up` на живом стенде — 1 м 5 с: мир заливается заново, и это +не оплошность — `WatchID` те же, повтор схлопнет `ReplacingMergeTree`. +Смоук — 9 с, `check-clickhouse` — 8 с с новой девятой проверкой. `make test` +вырос до 71 с: 58 с прежних тестов плюс 13 с на пересборку восьми дней для +сверки с описью. Прежние 40 с в таблице устарели ещё до #42 — тесты добавляли +#41 и #43. + +Что новая проверка **умеет краснеть**, снято двумя поломками того же дня, и +проверялись обе ветви её диагноза. Снесли партицию `2026-06-03` в +`ods.event_rep` — проверка покраснела, показала недостающий день и назвала +адрес: «события не доехали до ODS, начните с чтеца топика». Затем положили в +сырьё заведомо негодную строку — брак появился, и проверка сменила диагноз на +«сломан разбор, начните с `ods.event_errors_dist`». Строки опыта убраны, день +переигран `make generate-batch GENERATOR_DAY=2`; счёт вернулся к 401 185, и +это заодно показало дедупликацию: повтор дня не удвоил счёт под `FINAL`. + +Замеры 6 августа 2026 года, стенд поднят заранее. Время взято по `time` и +совпадает с тем, что цель печатает сама. Оно плавает от прогона к прогону: +смоук дал 8 секунд дважды (без проверки Kafka, до её появления, — 5 и 6), +`check-services` — 43 и 44. В таблице стоит большее из замеренных. До деления `scripts/stand-smoke.sh` шёл 48 секунд на 25 проверок, из них 42 секунды съедали шесть: Kafka с машины, два запуска пробников Airflow и три diff --git a/docs/specs/2026-07-30-stand-v2-realism.md b/docs/specs/2026-07-30-stand-v2-realism.md index 5eca965..93f96ca 100644 --- a/docs/specs/2026-07-30-stand-v2-realism.md +++ b/docs/specs/2026-07-30-stand-v2-realism.md @@ -253,7 +253,7 @@ CSV в репозитории (`data/catalog/products.csv`: `sku`, `name`, `cate Пятое — **опоздание** — бесплатно даёт формат доставки: часть заказов впервые появляется в слепке D+1/D+2 («вчера не сходилось, сегодня сошлось»), ориентир ~10%. Точные доли фиксируются при пересборке эталонного мира; -манифест хранит точные счётчики по каждому классу расхождений (отмены, +опись хранит точные счётчики по каждому классу расхождений (отмены, потери, дубли). Не берём: сироту-фрод (`purchase` есть, а заказа не будет никогда) — @@ -286,7 +286,7 @@ CSV в репозитории (`data/catalog/products.csv`: `sku`, `name`, `cate `uniq(посетителей) > uniq(людей)`, менти выводит расхождение сам. Константа мира: каждый двухкуковый покупатель делает минимум по одному заказу с каждой куки — иначе вторая кука не попадает в карту соответствий - (она строится только из покупок) и лаба не воспроизводится. Манифест + (она строится только из покупок) и лаба не воспроизводится. Опись хранит число именно таких пар. - Витрины разводят имена честно: **«посетители»** (`uniq(ClientID)`) и **«известные пользователи»** (после склейки) — оба числа рядом в дашборде. @@ -464,19 +464,19 @@ Kafka день переигрывается генератором заново: пачками), Postgres остаётся только служебной базой Airflow. Смешанность ландшафта выражена режимами и частотами, а не второй трубой. -## 8. Эталонный мир и манифест +## 8. Эталонный мир и опись -В git хранится только манифест эталонного мира; сам снимок (14 модельных +В git хранится только опись эталонного мира; сам снимок (14 модельных дней) генерируется на месте — при `make up` и при проверках (решение развилки «Производительность» карты #26, подробности — -[спека генератора](2026-08-01-generator.md), раздел 5). Манифест несёт +[спека генератора](2026-08-01-generator.md), раздел 5). Опись несёт паспорт мира (каноническое зерно, версия генератора), контрольные счётчики и хеши по дням; проверки «пустой git diff» и «пересгенерируй день N — -сравни хеш» живут на нём. Политика версионирования артефакта (бывший туман -карты #10) закрыта этим же ходом: версионируется манифест. -Контрольные числа манифеста: +сравни хеш» живут на ней. Политика версионирования артефакта (бывший туман +карты #10) закрыта этим же ходом: версионируется опись. +Контрольные числа описи: -- заказная сторона: заказы и выручка по дням; манифест хранит точные +- заказная сторона: заказы и выручка по дням; опись хранит точные счётчики по каждому классу расхождений (отмены, потери, дубли, дельты сумм) — самопроверка лабы сверки; - идентичность: uniq кук, uniq известных пользователей, число двухкуковых @@ -497,7 +497,7 @@ v2 стартует пустым, поэтому объём ниже — это | SQL | DDL по слоям и ролям (ON CLUSTER, Replicated*, Distributed; раскладка файлов — в доке хранилища) + трансформации событий, заказов, identity, сверки + словарь | L — ~12–15 файлов, главная сложность | | Airflow | DAG'и по образцу v1: etl_pipeline (партиционная переобработка, ожидание дневного батча заказов — сенсор/Datasets), world_init/next_day, helpers | M — ~5–6 файлов | | Superset | датасеты + дашборд с тремя новыми сюжетами | M — 2 файла | -| Эталонный мир | пересборка снимка на месте, манифест-счётчики, чек-скрипты | M–L | +| Эталонный мир | пересборка снимка на месте, счётчики описи, чек-скрипты | M–L | | Мониторинг | дашборды Grafana «данные», «кластер», «запросы»; ClickHouse источником данных, панели на SQL; Prometheus тонким полом (ADR 0002) | M — конфиги и дашборды | | Документация | доки v2 пишутся заново (см. раздел 12) | M, в тех же PR | @@ -534,15 +534,15 @@ v2 стартует пустым, поэтому объём ниже — это 2. DDL и генератор (слиты в один этап — DDL проверяется только настоящими данными): базы и таблицы событий ON CLUSTER, приём `hits` обеими нодами; широкое событие, таксономия, анонимность, N:1 (клиентская сторона - целиком). В конце этапа фиксируется маленький «зерновой» мир для + целиком). В конце этапа фиксируется маленький стартовый мир для стабильных приёмок следующих этапов (полная пересборка эталонного мира — отдельный этап 7). 3. Заказы и каталог: генератор слепков, STG/ODS/DDS заказа, словарь. 4. Трансформации и витрины: сессии, identity_map, выручка, сверка A+C. 5. Airflow: `etl_pipeline` (партиционная переобработка, ожидание дневного батча заказов — сенсор/Datasets). -6. Расхождения B+D и опоздания; счётчики манифеста. -7. Эталонный мир: манифест и пересборка снимка, чек-скрипты; CI-генерация +6. Расхождения B+D и опоздания; счётчики описи. +7. Эталонный мир: опись и пересборка снимка, чек-скрипты; CI-генерация на amd64 и arm64. 8. Superset-дашборд v2. 9. Мониторинг и runbook «keeper упал / DDL повис в очереди». Состав дашбордов @@ -563,7 +563,7 @@ v2 стартует пустым, поэтому объём ниже — это - Лабы и курс: v2 — другой стенд, лабы для него пишутся с нуля отдельной работой после этой спеки; редизайн лаб v1 (#7) остаётся в v1 и сюда не переносится. Спека даёт будущим лабам только опорные точки — контрольные - числа манифеста (сверка, идентичность). Явное следствие: после этапа 9 + числа описи (сверка, идентичность). Явное следствие: после этапа 9 стенд работает, но учебного пути на нём ещё нет. - Инкрементальный ETL (#8) — свой issue. - Реплики (2×2), HAProxy, репликационная эксплуатация — в лекцию, не в стенд. diff --git a/docs/specs/2026-08-01-generator.md b/docs/specs/2026-08-01-generator.md index db9aa7e..b5e23ef 100644 --- a/docs/specs/2026-08-01-generator.md +++ b/docs/specs/2026-08-01-generator.md @@ -31,14 +31,14 @@ модельный день D — функция (зерно, D). Между прогонами живут только зерно и позиция на оси времени. - **Детерминизм до байта.** Одно зерно — побайтово тот же снимок; сверка — - хешами манифеста. Транспорт (офсеты Kafka, темп) — вне обещания. + хешами описи. Транспорт (офсеты Kafka, темп) — вне обещания. - **Схема — контракт генератора.** Python-модуль с чистыми данными; хранилище строится по рендеренной документации, границу сторожит строгий приём на стороне хранилища. - **Один сериализатор, глупые приёмники.** День-функция выдаёт канонические байты; приёмники — файл, Kafka пачкой, Kafka с темпом. - **Числа.** Средний день ~50 тыс. событий; эталонный снимок — 14 дней; - в git — только манифест; автоматический порог один — день ≤ 30 с. + в git — только опись мира; автоматических порогов по времени нет. ## 1. Модель мира @@ -93,7 +93,7 @@ человеку предыстории, чьё окно активности таких дней не оставляет, пара не назначается. День-функция обязана назначенные заказы реализовать; остальные покупки — вольные, их решает генератор торговых - событий (#40). Манифест считает пары, реализованные в горизонте + событий (#40). Опись считает пары, реализованные в горизонте снимка. Условность в данных не видна: дни назначены той же случайностью, просто брошенной планом один раз. - **Своя ось модельного времени.** Ось событий начинается в @@ -117,19 +117,19 @@ Отклонено с доводами: - *Мутирующее состояние мира* («мир стареет»): ломает параллельность по - дням, требует чекпоинтов, счётчики манифеста узнаваемы только постфактум; + дням, требует чекпоинтов, счётчики описи узнаваемы только постфактум; ни один урок стенда на старении не стоит. - *Чистая функция без слоя состава*: глобальные инварианты пришлось бы выводить в каждом дне заново — тот же план мира, но неявный и размазанный. - *Привязка модельного времени к реальному календарю* (T-1 с догоном): конфликтует с ускорением ×60 — за вечер мир уезжает в будущее — и делает - даты эталонного мира зависимыми от даты запуска, манифест теряет + даты эталонного мира зависимыми от даты запуска, опись теряет воспроизводимость. Отклонено при исполнении #38 (2026-08-02): - *Разгон вместо предыстории* («первые дни малы — магазин запустился»): - зерновой мир и половина снимка оказались бы на разгоне, недельная лаба + стартовый мир и половина снимка оказались бы на разгоне, недельная лаба сравнивала бы несравнимые недели, «средний день ~50 тыс.» перестал бы быть средним — пришлось бы двигать принятые числа раздела 5. - *Материализованный план на горизонт*: горизонт становится обязательным @@ -139,7 +139,7 @@ константы — потолка одновременно живущих кук, которого в жизни нет и который менти нечем объяснить. - *Вероятностная гарантия пар* («почти наверняка купит с обеих кук»): не - гарантия — однажды манифест покраснеет, а чинить нечем, кроме смены + гарантия — однажды опись покраснеет, а чинить нечем, кроме смены зерна; при этом несклеенная пара в данных неотличима от двух незнакомцев, так что реализм этой лотереи невидим. - *Покупка пары в первый визит куки*: гарантия железная и дешёвая, но узор @@ -154,14 +154,14 @@ - **Обещание — содержимое до байта.** Два прогона с одним зерном дают тот же набор событий: те же `WatchID`/`VisitID`, поля, метки модельного времени. Снимок при пересборке побайтово совпадает: канонический порядок ключей и - строк; сверка — по хешам манифеста, а два локально пересобранных снимка + строк; сверка — по хешам описи, а два локально пересобранных снимка сравнимы обычным diff — пустой означает «ничего не изменилось». Вне обещания — транспорт: офсеты и партиции Kafka, какая нода прочитала, `_load_ts`, темп живого дня. - **Условия обещания.** Детерминизм держится при зафиксированном `uv.lock` и внутри канонического контейнера — то есть везде Linux, на маке и в WSL тоже; единственная переменная — архитектура CPU. Истина — CI на Linux; сходимость любой машины проверяет - скрипт «пересгенерируй день N — сравни хеш с манифестом». Расхождение на + скрипт «пересгенерируй день N — сравни хеш с описью». Расхождение на любой платформе — баг генератора, а не допуск. - **Раздача зерна — иерархией подпотоков.** Корневое зерно → подпоток состава мира, ветвящийся по номеру дня на когорты плана (состав @@ -173,7 +173,7 @@ торговые события, расхождения, опоздания — в фиксированном порядке. По построению: параллельный прогон равен последовательному; продление истории днём N+1 не трогает дни 1…N; правка одного компонента меняет - только его часть снимка — в манифесте меняются хеши только затронутых + только его часть снимка — в описи меняются хеши только затронутых дней, дифф двух локальных пересборок читаем. - **Механизм подпотоков — `numpy.random.SeedSequence`.** Сверено через Context7 по документации numpy (2026-08-01): `spawn(n)` порождает детей @@ -188,7 +188,7 @@ представление в клиентском `purchase` (урок мастер-спеки о расхождениях). - **Канонический seed и паспорт мира.** Эталонный мир собирается одним каноническим зерном — константой репозитория; свои зёрна менти крутит без - гарантий манифеста. Манифест хранит паспорт мира — зерно и версию + гарантий описи. Опись хранит паспорт мира — зерно и версию генератора; чек-скрипты сверяют паспорт раньше счётчиков. - **Суточный профиль интенсивности задаёт день-функция.** Форма — волны: ночной провал, обеденный и вечерний пики, различие будней и выходных; @@ -207,7 +207,7 @@ - *Обещание детерминизма поверх обновления зависимостей*: numpy сознательно улучшает алгоритмы распределений между версиями (NEP 19), Faker меняет словари. Фиксация — `uv.lock`; обновление зависимостей — осознанная - пересборка манифеста одним PR. + пересборка описи одним PR. ## 3. Контракт схемы @@ -255,7 +255,7 @@ - **Один канонический сериализатор, глупые приёмники.** День-функция выдаёт упорядоченный поток канонических байтов — единственное место, где событие превращается в JSON. Приёмники не знают о содержимом: файл (локальный кэш - для пересборки и проверок манифеста), Kafka пачкой — пакетный режим, + для пересборки и проверок описи), Kafka пачкой — пакетный режим, Kafka с темпом ×60 — живой день. Новых топиков нет. - **Одно событие — одно сообщение Kafka.** «Пачкой» относится к темпу отправки, а не к упаковке: приёмник шлёт события подряд без пауз, но каждое @@ -288,7 +288,7 @@ - **Рабочий выбор сериализатора — orjson**: быстрее stdlib json в 5–14 раз, numpy-массивы и datetime сериализует нативно (заметка исследования #31). Смена библиотеки меняет канонические байты, поэтому проходит как - обновление зависимости: осознанная пересборка манифеста одним PR. + обновление зависимости: осознанная пересборка описи одним PR. - **Промежуточные файлы не хранятся.** Файл дня — кэш чистой функции: потерял — пересчитал. В git снимок не попадает (раздел 5). - **Обрыв любого режима — переигровка дня целиком**; дедуп склеивает @@ -306,7 +306,7 @@ - *Отдельный топик / Kafka как хранилище дней*: офсеты и партиции вне обещания детерминизма, retention конечен, в git топик не положишь, хеш с - манифестом не сверишь; Kafka на стенде — труба, не хранилище (раздел 7 + описью не сверишь; Kafka на стенде — труба, не хранилище (раздел 7 мастер-спеки). - *Файл как обязательная станция доставки*: доигрывание обрыва уже решено через дедуп, канон держит единственный сериализатор, а не диск; файл @@ -327,15 +327,15 @@ - **Эталонный снимок — 14 дней**: две полные календарные недели, D0 — понедельник. Самая короткая длина, при которой есть замороженная зона за окном K = 7, дышащая зона и две волны недельной сезонности. Удлинение до - месяца — дешёвый ход (пересборка манифеста), если понадобится. -- **В git — только манифест, снимок не хранится.** Снимок генерируется при + месяца — дешёвый ход (пересборка описи), если понадобится. +- **В git — только опись, снимок не хранится.** Снимок генерируется при `make up` и при проверках: артефакт — кэш чистой функции, кэш в git не - хранят. Манифест несёт паспорт мира, счётчики и хеши по дням; проверки - «пустой git diff» и «пересгенерируй день N — сравни хеш» живут на нём. + хранят. Опись несёт паспорт мира, счётчики и хеши по дням; проверки + «пустой git diff» и «пересгенерируй день N — сравни хеш» живут на ней. Каждый `make up` — живая демонстрация детерминизма. Честная потеря — страховка на случай платформенного бага: раньше менти с расходящимися байтами мог взять готовый снимок из git, теперь он упрётся в красный чек - манифеста; смягчение — CI гоняет генерацию на amd64 и arm64. + описи; смягчение — CI гоняет генерацию на amd64 и arm64. - **Живой день — ×60 по умолчанию**: модельные сутки за 24 реальные минуты, суточная волна разворачивается на глазах; темп в среднем ~35 событий/с, в пиковые часы сильных дней — до ~100. Число — @@ -348,34 +348,37 @@ Схема двухъярусная — урок ADR 0004: пороги впритык к расчёту на разном железе кончаются ритуальным удалением проверки. -Ориентиры на референсной машине — в спеке, без автоматики: +Ориентиры на референсной машине — в спеке, без автоматики. Где стоит замер, +там он с датой; остальное — порядок величины, пока не мерили: | Операция | Ориентир | |---|---| -| Генерация одного дня | секунды | -| Пересборка эталонного мира (14 дней + манифест, без транспорта) | до минуты | -| Заливка снимка при `make up` (Kafka → матвью → ODS) | минуты | +| Генерация одного дня | 1,7–2,8 с (замер 7 августа 2026 года) | +| Пересборка эталонного мира (14 дней + опись, без транспорта) | до минуты | +| Заливка стартового мира при `make up` (8 дней, Kafka → матвью → ODS) | 22,5 с (замер 7 августа 2026 года) | | Лаг живого дня | секунды | -**Автоматический порог один: полный день (50 тыс. событий) генерируется -≤ 30 с.** Расчёт по планке исследования (~2–5×10⁵ событий/с на ядро) — доли -секунды; порог держит машинный разброс ×2–5 и ловит деградацию на 1–2 -порядка: Faker в горячем цикле, случайная квадратичность. Реализация — -pytest-тест с маркером `perf` и таймаутом-обрубанием: обязателен в CI, -исключён из быстрой локальной петли, зовётся отдельной целью при правках -горячего цикла. +**Автоматических порогов нет ни одного** — решение владельца 7 августа +2026 года при исполнении #42. Здесь стоял единственный: «полный день +генерируется ≤ 30 с». Назначен он был до того, как генератор написали, а +первый замер настоящего кода дал **1,7 секунды на день в 50 626 событий** +(машина стенда, 7 августа 2026 года). Порог оказался в восемнадцать раз выше +факта: зелен при любой правдоподобной регрессии, то есть не сторож, а +украшение. Как здесь говорят о времени, к тому дню уже решила [карта +целей](../architecture/testing.md): цена — замеренное число с датой, а не +назначенный предел. -Остальное — наблюдаемость без порогов: `make up` и smoke печатают тайминги -(генерация и доставка отдельно), проигрыватель логирует лаг. Прототип-замер -до этапа 2 не нужен: числа назначены с запасом порядок и больше от планки -исследования, планка подтверждена локальной проверкой на машине стенда; -первый замер настоящего кода — порог этапа 2. +Наблюдаемость без порогов остаётся и делает всю работу: проигрыватель +печатает тайминги генерации и доставки раздельно, лаг живого дня — в его +логе, а `make up` после #42 гоняет генератор по восемь дней при каждом +подъёме стенда: замедлись он на порядок — подъём стенда встанет колом, и +заметит это первый же человек, который его поднял. Отклонено с доводами: - *Четыре жёстких CI-ворот на все бюджеты*: машинный разброс против порогов - впритык — повторение истории с памятью (ADR 0004); порог оставлен один, - грубый, между «×5 шума» и «×100 беды». + впритык — повторение истории с памятью (ADR 0004). Оставленный было один + грубый порог снят при исполнении #42 (выше). - *Снимок в git* (статус-кво раздела 8 мастер-спеки): основание из v1 — медленный генератор — съедено детерминизмом и скоростью; остаётся только раздутый репозиторий. *Снимок вложением релиза Gitea*: страховка без @@ -411,9 +414,9 @@ pytest-тест с маркером `perf` и таймаутом-обрубан - **Раздел 1.4**: «из контракта выводятся DDL и валидация» заменено на data contract — хранилище пишется по документации, границу сторожит строгий приём (раздел 3 здесь). -- **Раздел 8**: артефакт `data/startup_history/` в git заменён манифестом; +- **Раздел 8**: артефакт `data/startup_history/` в git заменён описью; снимок генерируется на месте (раздел 5 здесь). Туман «политика - версионирования артефакта» закрыт этим же ходом: версионируется манифест. + версионирования артефакта» закрыт этим же ходом: версионируется опись. - **Раздел 11**: пункт «до этапа 3 зафиксировать требования производительности» закрыт числами раздела 5. - Мелкие согласования там, где текст опирался на артефакт в git: источник @@ -444,7 +447,7 @@ pytest-тест с маркером `perf` и таймаутом-обрубан по расписанию (~раз в 24 минуты); генератор не дорабатывается. - **Этап 5, позиция на оси времени.** Проигрыватель состояния не хранит (раздел 9), поэтому вести позицию — работа того, кто его зовёт. Живёт она - переменной Airflow: отсутствие переменной означает мир в зерновом + переменной Airflow: отсутствие переменной означает мир в стартовом состоянии, заводит и двигает её только даг `next_day` и только по успеху. Довод — генератор отдельная и заменяемая сущность, привязывать его к хранилищу незачем, а `make clean` сносит том метаданных Airflow вместе с @@ -452,7 +455,7 @@ pytest-тест с маркером `perf` и таймаутом-обрубан Расхождение переменной с данными (менти почистил партицию руками) лечится документацией, а не сторожем: постоянная сверка была бы проверкой, красной в норме ([ADR 0002](../adr/0002-monitoring-scope.md)). -- **Этап 7 (эталонный мир)**: пересборка — это манифест, не артефакт; +- **Этап 7 (эталонный мир)**: пересборка — это опись, не артефакт; CI-генерация на amd64 и arm64. - **Будущие лабы**: перезаливка дня X пакетным режимом проигрывателя — готовая демонстрация идемпотентности конвейера. @@ -503,11 +506,11 @@ pytest-тест с маркером `perf` и таймаутом-обрубан возвратами предыстории), 170 пар. - **Конфигурация мира — модуль чистых данных** рядом с контрактом схемы: все числа мира в одном месте, написанном как приглашение любопытному - менти крутить. Правка модуля — смена мира: чек манифеста честно - краснеет, манифест сторожит только канон. Вне модуля — лишь то, что + менти крутить. Правка модуля — смена мира: чек описи честно + краснеет, опись сторожит только канон. Вне модуля — лишь то, что мира не меняет: своё зерно и транспортные флаги проигрывателя. Отклонено: внешний конфиг и env-переопределения — переменная мира, - которую паспорт манифеста не видит; файл-конфиг в репозитории — по + которую паспорт описи не видит; файл-конфиг в репозитории — по смыслу равен модулю, но платит загрузчиком и валидацией (довод раздела 3 против YAML). @@ -590,7 +593,7 @@ pytest-тест с маркером `perf` и таймаутом-обрубан посетителю 6% до корзины и те же 45% и 55%, помеченному покупателю — 18% и 60%. Перемеренный средний день: 9 879 визитов, 47 433 pageview, 4,80 страницы на визит, конверсия визита 2,4%; полные числа — в блоке #40, - оттуда же они лягут в манифест. + оттуда же они лягут в опись. - **Недельная волна применяется один раз.** Профиль ведёт приток, трафик наследует его через дневную аудиторию; измеренный размах трафика — ±7% против ±10% у притока. Второе умножение удвоило бы недельный размах. Если @@ -611,8 +614,8 @@ pytest-тест с маркером `perf` и таймаутом-обрубан Решено при исполнении #40 (2026-08-02) — решения владельца до реализации. Три решения выходят за границы тикета: правятся план состава (#38), воронка -(#39) и каталог товаров (#39). Сейчас это дёшево — манифест ещё не собран -(#42); позже обошлось бы его пересборкой. Часть решений принята после +(#39) и каталог товаров (#39). Сейчас это дёшево — опись ещё не собрана +(#42); позже обошлось бы её пересборкой. Часть решений принята после холодного ревью первой редакции блока. - **Корзина шире заказа.** Посетитель кладёт в корзину товары тех карточек, @@ -872,7 +875,7 @@ pytest-тест с маркером `perf` и таймаутом-обрубан цены и по медиане, потому что одни только края списка обманываются четырьмя исключениями. Хвост соседям: словарь ClickHouse над этим файлом ещё не построен (#37/#43) — колонку надо взять в его описание сразу, а не - догонять правкой; манифесту нужен хеш каталога (хвост #42) тем более, + догонять правкой; описи нужен хеш каталога (хвост #42) тем более, потому что теперь в файле живёт ещё и поведение. - **Доля просмотров карточек, доходящих до корзины, выросла с 7,3% до 8,6% и подошла к потолку вилки.** Иначе не сошлось: треть кладущих визитов вне @@ -885,7 +888,7 @@ pytest-тест с маркером `perf` и таймаутом-обрубан пока число правдоподобно, само по себе оно ничего не сторожит, и подгонять поведение под его край не надо. Настоящий предел здесь другой и считается в другой валюте — дневной бюджет событий: события корзины - входят в те самые «около 50 тыс.», на которых стоят манифест (#42) и + входят в те самые «около 50 тыс.», на которых стоят опись (#42) и порог скорости дня. Упрётся будущая правка — двигать надо бюджет и его причины, а не долю. @@ -926,7 +929,7 @@ pytest-тест с маркером `perf` и таймаутом-обрубан итог виден кодом возврата. Общее у режимов: `--seed`, `--day`, `--days` и приёмник — `--file` либо `--brokers` с `--topic`. Состояния нет: позицию на оси ведёт зовущий (хвост этапу 5 — раздел 8). `--days` играет несколько - дней подряд одним запуском — этим зальётся зерновой мир (#42), восемь + дней подряд одним запуском — этим зальётся стартовый мир (#42), восемь запусков службы внутри `make up` были бы плохим ответом. Отклонено: *режимы флагом `--speed 0`* — различие тогда прячется за числом, тогда как у режимов разное устройство: у пакетного есть пачка и нет ожидания, у @@ -986,11 +989,30 @@ pytest-тест с маркером `perf` и таймаутом-обрубан модуль его не находит. Проверено запуском в контейнере; заодно день, сыгранный в образе, совпал побайтово с днём, сыгранным на машине. -Остаётся открытым, за тикетами: +Решено при исполнении #42 (2026-08-07): -- как фиксируется «зерновой» мир конца этапа 2 (раздел 9 мастер-спеки): - с манифестным решением напрашивается мини-манифест зернового мира — - форма за #42; -- паспорт мира в манифесте (зерно и версия генератора) файл каталога не - накрывает: правка цены в CSV меняет мир молча. Манифесту нужен хеш - каталога — хвост для #42. +- **Опись мира одна, и она же растёт до эталонной.** Стартовый мир конца + этапа 2 — первые восемь дней оси, понедельник по понедельник (решение + владельца при нарезке, 2026-08-01): полная неделя с выходными и первый + замкнутый цикл окна K = 7. Отдельного «младшего» файла для них не + заводится: `data/world-inventory.json` — та самая опись, которую раздел 5 + обещает эталонному миру, пока короткая. Этап 7 продлит её до четырнадцати + дней, а не заведёт вторую. Форма — JSON: паспорт мира (зерно, версия + генератора, хеш каталога) и по строке на день с датой, числом событий и + хешем. Собирается целью `make inventory`, свежесть сторожит тест + `test_inventory.py` — тем же способом, что свежесть «описания выгрузки». + Отклонено: *два файла, «мини» и полный* — две правды об одном мире и + лишнее слово в словаре. +- **Хеш каталога — в паспорте, и работа у него объяснительная.** Правка + цены в `data/catalog/products.csv` меняет мир так же молча, как правка + кода, но поймают её и без хеша: хеши дней разойдутся. Хеш каталога + отвечает на следующий вопрос — **что** правили: разошлись дни и каталог — + CSV; разошлись только дни — код. +- **Слова.** «Манифест» по всей спеке переименован в **опись мира**, а + «зерновой мир» — в **стартовый мир** (решение владельца 2026-08-07). + Довод — правило языка репозитория: есть обычное русское слово — берётся + оно. Заодно исчезла пара «манифест / мини-манифест», в которой читателю + пришлось бы различать две сущности там, где вещь одна. + +Открытых вопросов за разделом не осталось: интерфейс запуска закрыт #41, +форма описи и хеш каталога — здесь. diff --git a/generator/README.md b/generator/README.md index a801a44..61250b4 100644 --- a/generator/README.md +++ b/generator/README.md @@ -28,7 +28,7 @@ D0 живёт предыстория, поэтому любой день соб бросок достаётся тому, кто спросил k-м. Приписать новый бросок в конец функции безопасно: у прежних он ничего не отнимает. Вставить в середину — значит сдвинуть все броски после него, а с ними и весь мир: события того же -дня станут другими, счётчики канонического мира разойдутся с манифестом, и +дня станут другими, счётчики канонического мира разойдутся с описью, и поймается это не ошибкой, а красным чеком. Ровно поэтому паспорта кук в `plan.cohort` бросаются последними. @@ -78,6 +78,10 @@ D0 живёт предыстория, поэтому любой день соб - `src/clickstream_generator/schema_doc.py` — сборка «описания выгрузки» ([`docs/formats/clickstream-event.md`](../docs/formats/clickstream-event.md)) из контракта. Документ руками не правят — пересобирают. +- `src/clickstream_generator/inventory.py` — сборка описи мира + ([`data/world-inventory.json`](../data/world-inventory.json)): паспорт мира + и хеши восьми дней, которыми наполняется стенд. Руками не правят — + пересобирают целью `make inventory`. - `tests/` — инварианты контракта, свежесть описания и обещания мира: чистота от зерна, приток, гарантия двухкуковых пар, форма суточной волны и сборка визитов по задокументированным правилам. Там же побайтовое diff --git a/generator/src/clickstream_generator/cli.py b/generator/src/clickstream_generator/cli.py index 45fcb1c..6dd719a 100644 --- a/generator/src/clickstream_generator/cli.py +++ b/generator/src/clickstream_generator/cli.py @@ -8,7 +8,7 @@ Зовущих трое, и все трое видны в форме команд: - даги `world_init` и `next_day` этапа 5 — по дню за запуск, приёмник Kafka; -- заливка зернового мира (#42) — восемь дней подряд одним запуском: `--days`; +- заливка стартового мира — восемь дней подряд одним запуском: `--days`; - проверки хранилища (#43) — ограниченная пачка в файл: `--limit` и `--file`. **Режимы разведены командами, а не флагом**, потому что различаются не темпом diff --git a/generator/src/clickstream_generator/inventory.py b/generator/src/clickstream_generator/inventory.py new file mode 100644 index 0000000..57a58fa --- /dev/null +++ b/generator/src/clickstream_generator/inventory.py @@ -0,0 +1,99 @@ +"""Опись мира: чем стенд наполняется при подъёме и каким это обязано выйти. + +Мир — чистая функция зерна (спека генератора, раздел 2), поэтому в git лежит не +он сам, а опись: паспорт мира, число событий по дням и хеш байтов каждого дня. +Сам мир пересчитывается когда угодно, а опись отвечает на единственный вопрос — +**тот ли это мир, что был вчера**. Разошлись хеши — мир уехал, и дальше уже +неважно, чего от него ждали проверки. + +Дней в описи восемь: столько заливается в стенд при `make up`. Понедельник по +понедельник — полная неделя с выходными и первый замкнутый цикл окна K = 7. +Эталонный снимок в четырнадцать дней придёт на этапе 7 и станет продолжением +этой же описи, а не вторым файлом. + +**Сторожат мир хеши, а не паспорт.** Паспорт отвечает на другой вопрос — «чем +это сделано»: зерно и версия генератора. Поменяй кто-нибудь код так, что мир +сдвинется, — версия останется прежней, а хеши покраснеют; наоборот не бывает. + +Хеш каталога стоит здесь по третьему основанию — ни сторожить, ни описывать, а +**объяснять**. Правка цены в `data/catalog/products.csv` меняет мир так же +молча, как правка кода, и по одним хешам эти два случая неразличимы. С хешем +каталога различимы: разошлись хеши дней и каталога — правили CSV; разошлись +только дни — правили код. + +Хеш дня — sha256 тех самых байтов, что уезжают в Kafka, с переводом строки +после каждого события. Это ровно то, что пишет файловый приёмник, поэтому +пересчитывается он и обычным `sha256sum` по сыгранному в файл дню (как +именно — в README репозитория). + +Собирается опись из корня репозитория целью `make inventory`, а свежесть её +сторожит тест — как и у «описания выгрузки». +""" + +import argparse +import hashlib +import json +from datetime import timedelta +from importlib.metadata import version +from pathlib import Path +from typing import Any + +from clickstream_generator import day as day_module +from clickstream_generator import serialize, world +from clickstream_generator.catalog import CATALOG_PATH +from clickstream_generator.seeds import CANONICAL_SEED + +# Сколько дней оси заливается в стенд при подъёме. То же число стоит у службы +# `world-init` в compose.yaml: YAML не читает Python, и одно из двух мест — +# лишнее по построению. Расхождение поймают счётчики make check-clickhouse. +STARTING_DAYS = 8 + + +def build() -> dict[str, Any]: + """Опись целиком: паспорт мира и по строке на каждый его день.""" + return { + "seed": CANONICAL_SEED, + "generator_version": version("clickstream-generator"), + "catalog_sha256": _digest(CATALOG_PATH.read_bytes()), + "days": [_day(number) for number in range(STARTING_DAYS)], + } + + +def render() -> str: + """Опись текстом файла: отступы в два пробела, кириллица как есть.""" + return json.dumps(build(), ensure_ascii=False, indent=2) + "\n" + + +def _day(number: int) -> dict[str, Any]: + """Строка описи: номер дня, его дата, число событий и хеш байтов. + + Дата считается от D0 арифметикой, а не берётся из событий: ось модельного + времени так и определена (`world.ORIGIN`), и по этой же дате счётчики + стенда обрамляют счёт в `ods.event`. Соври она — подневная сверка это и + покажет, каждый день сразу. + """ + payloads = serialize.events(day_module.stream(CANONICAL_SEED, number)) + return { + "day": number, + "date": (world.ORIGIN + timedelta(days=number)).isoformat(), + "events": len(payloads), + "sha256": _digest(b"".join(payload + b"\n" for payload in payloads)), + } + + +def _digest(payload: bytes) -> str: + return hashlib.sha256(payload).hexdigest() + + +def main() -> None: + parser = argparse.ArgumentParser( + description="Собирает опись мира: паспорт, счётчики и хеши дней." + ) + parser.add_argument("output", type=Path, help="путь к файлу описи") + output = parser.parse_args().output + output.write_text(render(), encoding="utf-8") + print(f"Опись мира собрана: {output}") + + +if __name__ == "__main__": + main() diff --git a/generator/src/clickstream_generator/player.py b/generator/src/clickstream_generator/player.py index 95aeaea..7453cbb 100644 --- a/generator/src/clickstream_generator/player.py +++ b/generator/src/clickstream_generator/player.py @@ -15,7 +15,7 @@ Тайминги печатаются раздельно — генерация и доставка, как требует спека (раздел 5): это разные машины разной природы, и сложенные в одно число они перестают что-либо говорить. Сериализация считается частью генерации: она -рождает те самые байты, которые сторожит манифест. +рождает те самые байты, которые сторожит опись. """ import logging diff --git a/generator/src/clickstream_generator/seeds.py b/generator/src/clickstream_generator/seeds.py index 0da6d53..31439d3 100644 --- a/generator/src/clickstream_generator/seeds.py +++ b/generator/src/clickstream_generator/seeds.py @@ -25,8 +25,8 @@ from enum import IntEnum import numpy as np -# Каноническое зерно эталонного мира — константа репозитория; манифест хранит -# его в паспорте мира. Свои зёрна менти крутит без гарантий манифеста. +# Каноническое зерно эталонного мира — константа репозитория; опись хранит +# его в паспорте мира. Свои зёрна менти крутит без гарантий описи. CANONICAL_SEED = 20260601 diff --git a/generator/src/clickstream_generator/world.py b/generator/src/clickstream_generator/world.py index 1ade6fb..872e9fd 100644 --- a/generator/src/clickstream_generator/world.py +++ b/generator/src/clickstream_generator/world.py @@ -2,8 +2,8 @@ Модуль — приглашение крутить: поменяйте число, пересоберите снимок и посмотрите, что стало с данными. Правка любой константы здесь — смена мира, -поэтому чек манифеста честно покраснеет: манифест сторожит только канонический -мир, свои миры менти собирает без его гарантий (спека генератора, раздел 9). +поэтому чек описи честно покраснеет: опись сторожит только канонический мир, +свои миры менти собирает без её гарантий (спека генератора, раздел 9). Числа решены спекой и связаны между собой; связки сторожат тесты `test_world.py`, чтобы правка одного числа не рассыпала вывод соседнего. @@ -26,7 +26,7 @@ COUNTER_TIMEZONE_MINUTES = 240 # D0 — первый день оси модельного времени, понедельник. Реальный календарь в # модели не участвует: дата нужна лишь затем, чтобы дни оси легли в # `EventDate`/`UTCEventTime` конкретными числами. От даты запуска мир не -# зависит — иначе манифест перестал бы быть воспроизводимым. +# зависит — иначе опись перестала бы быть воспроизводимой. ORIGIN = date(2026, 6, 1) # Приток: сколько новых людей приходит в мир в средний день. Каждый приводит diff --git a/generator/tests/test_inventory.py b/generator/tests/test_inventory.py new file mode 100644 index 0000000..32f2e47 --- /dev/null +++ b/generator/tests/test_inventory.py @@ -0,0 +1,27 @@ +"""Проверка описи мира: та ли она, что собирается из кода сегодня. + +Опись собирается из кода, значит разойтись они могут только одним способом — +код правили, опись не пересобрали. Ровно это здесь и сторожится, тем же +способом, что свежесть «описания выгрузки». + +Проверка дорогая — она пересчитывает восемь модельных дней целиком, и это +единственный способ сравнить хеши: дешевле мир не пересобрать. Зато краснеет +она там, где надо, — сразу после правки генератора, а не через полчаса на +поднятом стенде, где расхождение счётчиков выглядит поломкой хранилища. +""" + +import json +from pathlib import Path + +from clickstream_generator.inventory import build + +INVENTORY_PATH = Path(__file__).resolve().parents[2] / "data" / "world-inventory.json" + + +def test_inventory_is_up_to_date(): + stored = json.loads(INVENTORY_PATH.read_text(encoding="utf-8")) + assert stored == build(), ( + "опись мира отстала от кода — пересоберите: make inventory." + " Разошлись хеши дней и каталога — правили data/catalog/products.csv;" + " разошлись только дни — правили генератор" + ) diff --git a/generator/tests/test_player.py b/generator/tests/test_player.py index 0022699..273ce11 100644 --- a/generator/tests/test_player.py +++ b/generator/tests/test_player.py @@ -98,7 +98,7 @@ def _run_apart(path) -> None: def test_days_play_in_a_row(tmp_path, monkeypatch): """Дни идут подряд от названного, а пачка считается на весь прогон. - Восемь дней одним запуском — то, чем зальётся зерновой мир (#42), поэтому + Восемь дней одним запуском — то, чем заливается стартовый мир, поэтому порядок дней проверяется, а не предполагается. День здесь подменён коротким: проверяется ход проигрывателя, а не содержимое дня, и платить за полсотни тысяч событий трижды незачем. diff --git a/generator/tests/test_serialize.py b/generator/tests/test_serialize.py index 9f2927c..be03286 100644 --- a/generator/tests/test_serialize.py +++ b/generator/tests/test_serialize.py @@ -28,7 +28,7 @@ def test_every_event_carries_every_column(events): «Пусто» по контракту — пустое значение, а не отсутствие ключа: пропавший ключ уводит событие в брак целиком (ADR 0005). Порядок ключей — часть - канона: от него зависят байты, а значит и хеши манифеста. + канона: от него зависят байты, а значит и хеши описи. """ names = [column.name for column in schema.COLUMNS] for event in events: diff --git a/scripts/check-clickhouse.sh b/scripts/check-clickhouse.sh index 94e6357..85ae97c 100755 --- a/scripts/check-clickhouse.sh +++ b/scripts/check-clickhouse.sh @@ -2,6 +2,7 @@ set -Eeuo pipefail readonly ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +readonly INVENTORY="${ROOT_DIR}/data/world-inventory.json" readonly CLUSTER="clickstream_cluster" readonly LOCAL_TABLE="smoke_replicated_local" readonly DISTRIBUTED_TABLE="smoke_distributed" @@ -76,6 +77,58 @@ on_signal() { exit 130 } +# Единственный постоянный сторож цепочки Kafka → STG → ODS. +# +# Проверка спрашивает не «есть ли данные», а «те ли это данные»: подневный счёт +# в `ods.event` против описи мира, лежащей в git. Поэтому она умеет краснеть на +# сломанном разборе — уехали события в брак, счёт разошёлся, — и на потере по +# дороге, а не только на пустой таблице. +# +# **Счёт обрамлён датами стартового мира.** Мир растёт: менти переиграет день, +# этап 5 добавит следующий, — и голый счёт по таблице разойдётся с описью +# законно, без всякой поломки. +# +# **Счёт идёт через FINAL** — правило репозитория для ODS, довод в storage.md: +# голый `count()` по ReplacingMergeTree зависит от числа прошедших мержей. +# +# Таблица брака здесь не второе утверждение, а объяснение первого. Утверждай мы +# «брака нет», проверка краснела бы навсегда после первого же урока, где менти +# нарочно отправил в топик мусор, — и краснела бы не о том. +check_starting_world() { + local expected actual broken first_date last_date + + expected="$(jq -r '.days[] | "\(.date)\t\(.events)"' "$INVENTORY")" + first_date="$(jq -r '.days | first | .date' "$INVENTORY")" + last_date="$(jq -r '.days | last | .date' "$INVENTORY")" + + actual="$(query clickhouse-01 " + SELECT EventDate, count() + FROM ods.event_dist FINAL + WHERE EventDate BETWEEN '${first_date}' AND '${last_date}' + GROUP BY EventDate + ORDER BY EventDate + FORMAT TSV")" + + if [[ "$actual" == "$expected" ]]; then + printf 'ЗЕЛЁНО: события всех дней стартового мира на месте, счёт сходится с описью.\n' + return + fi + + broken="$(query clickhouse-01 " + SELECT error_class, count() + FROM ods.event_errors_dist + GROUP BY error_class + ORDER BY count() DESC + FORMAT TSV")" + printf 'Опись мира ожидает (дата, событий):\n%s\n' "$expected" >&2 + printf 'В ods.event лежит:\n%s\n' "${actual:-— ничего —}" >&2 + if [[ -n "$broken" ]]; then + printf 'В ods.event_errors по классам брака:\n%s\n' "$broken" >&2 + fail 'счёт разошёлся с описью, и в таблице брака есть строки: сломан разбор — начните с ods.event_errors_dist и матвью ods.event_mv' + fi + fail 'счёт разошёлся с описью, а таблица брака пуста: события не доехали до ODS — начните с чтеца топика stg.hits_raw_kafka и матвью приёма stg.hits_raw_mv' +} + assert_ddl_queue_completed() { local service="$1" local phase="$2" @@ -93,7 +146,7 @@ ensure_stand_running trap on_exit EXIT trap on_signal INT TERM -printf 'Проверка 1/8: описание кластера одинаково на обеих нодах...\n' +printf 'Проверка 1/9: описание кластера одинаково на обеих нодах...\n' cluster_sql="SELECT cluster, shard_num, replica_num, host_name, port FROM system.clusters WHERE cluster = '${CLUSTER}' ORDER BY shard_num, replica_num FORMAT TSV" cluster_01="$(query clickhouse-01 "$cluster_sql")" cluster_02="$(query clickhouse-02 "$cluster_sql")" @@ -102,7 +155,7 @@ assert_equal "$expected_cluster" "$cluster_01" "неверная тополог assert_equal "$expected_cluster" "$cluster_02" "неверная топология на второй ноде" printf 'ЗЕЛЁНО: обе ноды видят ожидаемые два шарда: clickhouse-01 и clickhouse-02.\n' -printf 'Проверка 2/8: у нод разные макросы shard и replica...\n' +printf 'Проверка 2/9: у нод разные макросы shard и replica...\n' macros_sql="SELECT macro, substitution FROM system.macros WHERE macro IN ('shard', 'replica') ORDER BY macro FORMAT TSV" macros_01="$(query clickhouse-01 "$macros_sql")" macros_02="$(query clickhouse-02 "$macros_sql")" @@ -111,14 +164,14 @@ assert_equal $'replica\tclickhouse-02\nshard\t02' "$macros_02" "неверные [[ "$macros_01" != "$macros_02" ]] || fail "макросы нод не должны совпадать" printf 'ЗЕЛЁНО: clickhouse-01=(shard 01, replica clickhouse-01), clickhouse-02=(shard 02, replica clickhouse-02).\n' -printf 'Проверка 3/8: keeper отвечает обеим нодам...\n' +printf 'Проверка 3/9: keeper отвечает обеим нодам...\n' query clickhouse-01 "SELECT name FROM system.zookeeper WHERE path = '/' ORDER BY name FORMAT Null" query clickhouse-02 "SELECT name FROM system.zookeeper WHERE path = '/' ORDER BY name FORMAT Null" printf 'ЗЕЛЁНО: system.zookeeper доступна с обеих нод.\n' cleanup_tables || fail "не удалось очистить объекты предыдущего запуска" -printf 'Проверка 4/8: ReplicatedMergeTree создаётся через ON CLUSTER...\n' +printf 'Проверка 4/9: ReplicatedMergeTree создаётся через ON CLUSTER...\n' query clickhouse-01 " CREATE TABLE default.${LOCAL_TABLE} ON CLUSTER ${CLUSTER} ( @@ -136,13 +189,13 @@ assert_equal $'smoke_replicated_local\tReplicatedMergeTree' "$(query clickhouse- assert_equal $'smoke_replicated_local\tReplicatedMergeTree' "$(query clickhouse-02 "$tables_sql")" "локальная таблица не создана на второй ноде" printf 'ЗЕЛЁНО: ReplicatedMergeTree видна в system.tables на обеих нодах.\n' -printf 'Проверка 5/8: путь в keeper собран из макроса shard...\n' +printf 'Проверка 5/9: путь в keeper собран из макроса shard...\n' path_sql="SELECT zookeeper_path, replica_name FROM system.replicas WHERE database = 'default' AND table = '${LOCAL_TABLE}' FORMAT TSV" assert_equal "/clickhouse/tables/01/${LOCAL_TABLE}"$'\t'"clickhouse-01" "$(query clickhouse-01 "$path_sql")" "неверные путь или имя реплики на первой ноде" assert_equal "/clickhouse/tables/02/${LOCAL_TABLE}"$'\t'"clickhouse-02" "$(query clickhouse-02 "$path_sql")" "неверные путь или имя реплики на второй ноде" printf 'ЗЕЛЁНО: пути собраны из shard (/01/ и /02/), имя реплики собрано из макроса replica.\n' -printf 'Проверка 6/8: Distributed создаётся ON CLUSTER и передаёт данные между нодами...\n' +printf 'Проверка 6/9: Distributed создаётся ON CLUSTER и передаёт данные между нодами...\n' query clickhouse-01 " CREATE TABLE default.${DISTRIBUTED_TABLE} ON CLUSTER ${CLUSTER} AS default.${LOCAL_TABLE} @@ -164,12 +217,12 @@ sharding_sql="SELECT countIf(_shard_num != cityHash64(ClientID) % 2 + 1), uniqEx assert_equal $'0\t2' "$(query clickhouse-02 "$sharding_sql")" "Distributed использует неверный ключ шардирования" printf 'ЗЕЛЁНО: локальная строка первой ноды читается со второй; ключ cityHash64(ClientID) разложил строки по двум шардам.\n' -printf 'Проверка 7/8: в очереди распределённых DDL нет незавершённых заданий...\n' +printf 'Проверка 7/9: в очереди распределённых DDL нет незавершённых заданий...\n' assert_ddl_queue_completed clickhouse-01 'после CREATE' assert_ddl_queue_completed clickhouse-02 'после CREATE' printf 'ЗЕЛЁНО: очередь содержит задания CREATE, незавершённых среди них нет.\n' -printf 'Проверка 8/8: временные таблицы удаляются через ON CLUSTER...\n' +printf 'Проверка 8/9: временные таблицы удаляются через ON CLUSTER...\n' query clickhouse-01 "DROP TABLE default.${DISTRIBUTED_TABLE} ON CLUSTER ${CLUSTER} SYNC" >/dev/null query clickhouse-01 "DROP TABLE default.${LOCAL_TABLE} ON CLUSTER ${CLUSTER} SYNC" >/dev/null remaining_sql="SELECT count() FROM system.tables WHERE database = 'default' AND name IN ('${LOCAL_TABLE}', '${DISTRIBUTED_TABLE}') FORMAT TSVRaw" @@ -179,4 +232,10 @@ assert_ddl_queue_completed clickhouse-01 'после DROP' assert_ddl_queue_completed clickhouse-02 'после DROP' trap - EXIT INT TERM printf 'ЗЕЛЁНО: временные таблицы удалены; проверены завершённые задания CREATE и DROP.\n' -printf 'ИТОГ: все 8 проверок кластера ClickHouse прошли.\n' + +# После снятия ловушек: своих объектов эта проверка не заводит и прибирать за +# собой ей нечего — она только смотрит на то, что стенд произвёл сам. +printf 'Проверка 9/9: стартовый мир в ods.event сходится с описью...\n' +check_starting_world + +printf 'ИТОГ: все 9 проверок кластера ClickHouse прошли.\n' diff --git a/scripts/wait-for-world.sh b/scripts/wait-for-world.sh new file mode 100755 index 0000000..df73bd3 --- /dev/null +++ b/scripts/wait-for-world.sh @@ -0,0 +1,79 @@ +#!/usr/bin/env bash +# Ждёт, пока стартовый мир доедет до ODS. Второй шаг `make up`, и вот зачем он. +# +# `docker compose up --wait` дожидается служб, в том числе успешно отработавшей +# заливки. Но «заливка кончилась» — это «события лежат в топике», а не «события +# в хранилище»: приём асинхронный. Движок Kafka копит блок и отдаёт его по +# размеру либо по `stream_flush_interval_ms` (умолчание 7,5 с), а стартовый мир — +# около 400 тыс. сообщений. Верни `make up` управление сразу — и `make +# check-clickhouse` следом покраснел бы по устройству, а не по поломке. +# +# Ждём ограниченным циклом с потолком, а не паузой наугад: пауза либо коротка на +# медленной машине, либо ворует минуту на быстрой. +# +# Условие ожидания — «всё доехало», а не «всё разобралось»: сумма событий и +# брака. Сломайся разбор — события уедут в `*_errors`, сумма сойдётся, ожидание +# кончится, и поломку назовёт `make check-clickhouse`. Ждать здесь одних годных +# событий значило бы висеть пять минут вместо внятного ответа. +set -Eeuo pipefail + +readonly ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +readonly INVENTORY="${ROOT_DIR}/data/world-inventory.json" +readonly ATTEMPTS=100 +readonly PAUSE_SECONDS=3 +# Как часто отчитываться о ходе: молчащая минуту команда выглядит зависшей. +readonly REPORT_EVERY=5 + +read -r -a COMPOSE_CMD <<<"${COMPOSE_BIN:-docker compose}" + +fail() { + printf 'ОШИБКА: %s\n' "$1" >&2 + exit 1 +} + +query() { + "${COMPOSE_CMD[@]}" --project-directory "$ROOT_DIR" exec -T clickhouse-01 \ + clickhouse-client --query "$1" /dev/null || true)" + if [[ "$arrived" =~ ^[0-9]+$ ]] && ((arrived >= expected)); then + if ((arrived > expected)); then + # Повторный `make up` заливает мир заново. Промолчи мы об этом — + # удвоенное число под «ЗЕЛЁНО» выглядело бы поломкой. + printf 'ЗЕЛЁНО: стартовый мир доехал; строк %s при ожидаемых %s —' \ + "$arrived" "$expected" + printf ' мир заливали повторно, повтор схлопнет ReplacingMergeTree.\n' + else + printf 'ЗЕЛЁНО: стартовый мир доехал: %s строк.\n' "$arrived" + fi + exit 0 + fi + if ((attempt % REPORT_EVERY == 0)); then + printf 'Доехало %s из %s, ждём дальше (%d с)...\n' \ + "${arrived:-0}" "$expected" "$((attempt * PAUSE_SECONDS))" + fi + sleep "$PAUSE_SECONDS" +done + +fail "стартовый мир не доехал за $((ATTEMPTS * PAUSE_SECONDS)) с: + доехало ${arrived:-0} из ${expected}. Смотрите журнал заливки + (docker compose logs world-init) и чтеца топика: + SELECT * FROM system.kafka_consumers"