Files
ddadminandClaude Fable 5 103ac021c8 feat(generator): инкрементальные счётчики manifest без перечитки Kafka
- Зачем:
  - 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>
2026-07-22 23:47:02 +03:00

14 KiB
Raw Permalink Blame History

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:

  1. Поднимите стенд: make up.
  2. Если DDL ещё не применён, запустите ddl_init.
  3. Снимите паузу с etl_pipeline, если он ещё paused: docker compose exec -T airflow-webserver airflow dags unpause etl_pipeline.
  4. Запустите world_init с пустой формой. По умолчанию он импортирует эталонный мир.
  5. Когда нужен ещё один модельный день, запустите 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.