- Зачем:
- world_next_day перечитывал всю историю топиков Kafka ради
накопительных счётчиков — время прогона росло с возрастом мира
(issue #5, находка F9).
- Что:
- счётчики засеваются при import из уже прочитанного артефакта и при
backfill из потока; next-day продвигает их только событиями нового
дня, полного чтения Kafka больше нет;
- катящаяся контрольная сумма — сумма SHA-256 событий по модулю 2^256
(инкремент равен полному пересчёту), старый формат артефакта
принимается без изменений;
- точные множества click_id/user_domain_id вынесены из manifest в
цепочку контент-адресуемых фрагментов (<=10 000 ID, SHA-256-цепочка,
отдельный топик counter_chunks) — потолок сообщения Kafka не грозит,
предел 900 000 байт проверяется явно с понятной ошибкой;
- порядок записи всюду: фрагменты -> manifest -> state; старое локальное
состояние отклоняется с подсказкой перезапустить import;
- документация manifest/state обновлена (ARCHITECTURE, OPERATIONS,
runbook startup-history).
- Проверка:
- make test (216+31) и make lint зелёные;
- живая приёмка на чистом стенде: import 235 с; три прогона
world_next_day — 716/718/716 с (плоское время, O(нового дня));
мир 3->6 дней, 561 942 события; make generated-history-chain-check —
все порции и стыки однородны;
- тест равенства инкремента и полного пересчёта:
test_incremental_counters_equal_full_recompute.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
14 KiB
Runbook: стартовая история стенда
Этот runbook нужен, чтобы один раз создать стартовую историю генератора, сохранить её в файл и быстро восстановить на чистом стенде.
За устройством генератора см. generator/README.md
и спеку модельного времени.
Что дёшево
- Перезапуск без очистки volumes: Kafka хранит
generator_stateи manifest. - Live-продолжение после импорта: генератор стартует с
T_end, если настройки совпадают. - Восстановление чистого стенда из готового файла: события повторно пишутся в Kafka, ClickHouse наполняется штатным путём.
Что требует нового артефакта
- Другая длительность истории или другой профиль запуска.
- Другой
GEN_SEED,GEN_MODEL_T0, часовой пояс, скорость или настройки генерации. - Осознанный новый мир после несовместимого state: сначала сбросьте state через
GEN_STATE_RESET=trueили чистые volumes.
Профили запуска
Основной способ — глагол плюс профиль:
| Профиль | Для чего | Длительность | Live-ход |
|---|---|---|---|
ci |
Быстрая проверка и CI | 6h |
×1, тик 60 с |
daily-wave |
История с видимой суточной волной | 3d |
×60, тик 1 с |
Длительность можно переопределить через GEN_HISTORY_DURATION, например 2d.
Команда сама считает GEN_MODEL_T_END от GEN_MODEL_T0.
Эталонный мир
В репозитории хранится готовый мир за три модельных дня:
data/startup_history/reference-world.json.xz. Это обычный JSON-артефакт,
сжатый xz. Его добавляют прямо в Git, без Git LFS.
Сопровождающий проекта собирает файл одним запуском из корня репозитория:
ARTIFACT="$PWD/data/startup_history/reference-world.json.xz" \
make startup-history-export
Команда выполняет backfill с профилем daily-wave и сразу пишет xz-файл.
Несжатый JSON занимает около 850 МБ, сжатый файл — около 32 МБ. После проверки
файл можно добавить обычной командой:
git add data/startup_history/reference-world.json.xz
При импорте DAG дважды читает и распаковывает артефакт: во время предпроверки и перед записью в Kafka. Это увеличивает время импорта, но не меняет результат.
Менти запускает world_init с пустой формой. Тогда читается эталонный мир из
репозитория. Следующий модельный день добавляет отдельный беспараметрный DAG
world_next_day. Из консоли тот же импорт запускается без указания пути:
make startup-history-import
Пульт в Airflow
Основной учебный путь в Airflow UI:
- Поднимите стенд:
make up. - Если DDL ещё не применён, запустите
ddl_init. - Снимите паузу с
etl_pipeline, если он ещё paused:docker compose exec -T airflow-webserver airflow dags unpause etl_pipeline. - Запустите
world_initс пустой формой. По умолчанию он импортирует эталонный мир. - Когда нужен ещё один модельный день, запустите
world_next_dayс пустой формой.
world_next_day имеет расписание каждые 30 минут, но по умолчанию стоит на паузе.
Один запуск добавляет один модельный день; max_active_runs=1 не допускает
параллельных доливок.
Операции сопровождающего
В форме world_init сопровождающему дополнительно доступны операции:
backfill— создать стартовую историю. После записи в Kafka DAG сам запускаетetl_pipeline, ждёт завершения и выполняетcheck.check— сверить ClickHouse с manifest из Kafka.
Поля формы:
profileберётся из профилей генератора.durationможно оставить пустым, тогда берётся длительность профиля.seedиmodel_time_speed— необязательные переопределения мира.artifact_path: дляbackfill— куда сохранить файл; пусто — не сохранять. Дляimport— что читать; пусто — эталонный мир из репозитория:/opt/airflow/data/startup_history/reference-world.json.xz.
Backfill и import работают только на чистом стенде. Если Kafka-топики данных или
STG уже непустые, DAG упадёт до записи и подскажет make clean. Консольные
команды make generator-backfill и make startup-history-import делают такую же
предпроверку с хоста. Это защита от смешивания разных миров.
Границы пульта:
make up,make cleanи live-продолжение остаются в консоли.- Операции
continueв DAG нет намеренно: live — долгоживущий сервис, а пульт управляет разовыми пакетными операциями. - Airflow не получает доступ к жизненному циклу контейнеров; таски выполняют обычный Python-код генератора.
Если backfill сохраняет файл в ./data, он создаётся пользователем Airflow
внутри контейнера. Чтение работает из Airflow и консольных команд, но перезапись
чужого файла может потребовать удалить старый файл вручную.
Экспорт
По умолчанию создаётся трёхсуточный артефакт daily-wave с суточной волной:
make startup-history-export
Файл по умолчанию: /tmp/clickstream-startup-history.json. В live-продолжении
профиль daily-wave проживает модельные сутки примерно за 24 настенные минуты.
Для быстрой автоматической проверки явно задайте служебный профиль ci:
PROFILE=ci ARTIFACT=/tmp/clickstream-startup-history-ci.json \
make startup-history-export
Команда делает чистый backfill и пишет в файл один связный набор: события Kafka, state и manifest.
Импорт на чистый стенд
make clean
docker compose up -d clickhouse kafka
make ddl
make startup-history-import
sleep 10
make transform
ARTIFACT=data/startup_history/reference-world.json.xz make startup-history-check
CHECK_LIVE_SEAM=0 make generated-history-check
Это путь только импорта: live-строк после T_end ещё нет. Стык backfill/live
проверяйте через make generated-history-runtime-check или после
make generator-continue и повторного make transform.
Импорт не пишет напрямую в ClickHouse. Он воспроизводит события и служебные compact-топики в Kafka. ClickHouse читает данные через свои Kafka-таблицы и Materialized View, затем batch строит ODS, DDS и DM.
Во время импорта генератор за один проход по событиям создаёт локальное
накопительное состояние manifest. Эталонный файл не меняется: в нём этого
служебного раздела нет. Основная запись manifest в Kafka хранит суммы, катящуюся
контрольную сумму и ссылку на цепочку точных множеств. click_id и
user_domain_id лежат небольшими неизменяемыми фрагментами в отдельном
compact-топике. State хранит SHA-256-ссылку на основную запись.
backfill создаёт такое состояние сразу по ходу генерации.
world_next_day восстанавливает эти числа и добавляет только события нового
дня, не читая старую историю Kafka и не перезаписывая прежние фрагменты.
Контрольная сумма складывает 256-битные
SHA-256-отпечатки записей по модулю 2^256, поэтому результат не зависит от
границ пакетов и совпадает с полным пересчётом. Если старый локальный state не
содержит ссылку на накопительные числа, очистите стенд и повторите import.
Импорт запускайте с тем же профилем, на котором создан артефакт. Для служебного
артефакта ci профиль нужно задать явно:
PROFILE=ci ARTIFACT=/tmp/clickstream-startup-history-ci.json \
make startup-history-import
make startup-history-check сверяет контрольные числа DM-витрины с manifest
артефакта: события, визиты, пользователей и диапазон event_timestamp. Если
data-топики Kafka уже непустые, импорт остановится до публикации событий.
В текущем стеке kafka-python не даёт транзакционный producer для нескольких
топиков. Поэтому импорт остаётся clean-stand операцией: при ошибке записи он
удаляет import-топики Kafka, чтобы повторный импорт не дописал дубли. Если
ClickHouse уже успел прочитать частичные сообщения, очистите стенд через
make clean и повторите импорт.
Live-продолжение
После импорта запускайте live с тем же профилем, что был в артефакте. Для
обычного артефакта daily-wave достаточно команды:
make generator-continue
У daily-wave скорость ×60. Долгий простой стенда создаёт большую дыру в
модельном времени: ночь простоя может стать десятками модельных суток без
событий. Для чистой демонстрации лучше очистите стенд, затем запустите
make generator-reset или повторите импорт стартовой истории через
make startup-history-import.
make generator-reset начнёт новый live-мир только на чистом стенде; если в
Kafka data-топиках или STG уже есть строки, команда попросит make clean.
Если читаемый state есть, но настройки не совпадают, генератор падает с перечнем
полей. Это защита от смешения разных миров. Для намеренного нового мира
сначала очистите стенд через make clean, затем запускайте нужный сценарий.
После обновления кода старый state может оказаться в старом формате. При
GEN_STATE_RESET=false это теперь громкий отказ, а не тихий старт с нуля поверх
старой истории. Оператору нужно выбрать одно из двух: очистить стенд через
make clean и заново создать стартовую историю, либо осознанно начать новый
live-мир через make generator-reset на чистом стенде.
После нестандартного мира make generator-continue нужно запускать с теми же
настройками, что были у backfill/import. При расхождении генератор громко
покажет поля, которые не совпали.
Старые артефакты daily-wave, созданные до перехода на ×60 и тик 1 с, с новым
профилем несовместимы. Это ожидаемо: защита от смешения миров должна остановить
такой запуск.
Старые переменные GEN_RUN_MODE, GEN_STATE_RESET и GEN_MODEL_T_END остаются
низкоуровневым способом для отладки и прямого docker compose run.