Compare commits

..
13 Commits
Author SHA1 Message Date
ddadminandClaude Fable 5 5f72e63b14 docs(specs): модель DDS заказа перенесена из этапа 3 в этап 4
Зачем: черновик этапов держал DDS заказа в этапе 3, а принятый набор
docs/architecture/orders/ оставляет модель DDS проектированию этапа 4
(ingestion.md «Не входит», тикет #85 перевешен на #6) — расхождение
всплыло на холодном ревью нарезки этапа (#89).

Что: в разделе 9 этап 3 сужен до STG/ODS и дага приёма, модель DDS
заказа названа работой этапа 4.

Проверка: чтение; нарезка этапа 3 (#90–#96) согласована с этой строкой.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-17 22:57:20 +03:00
ddadminandClaude Fable 5 cf1c6c4636 docs(orders): спека превращена в принятый живой набор architecture/orders
Зачем: датированная спека — событие истории, а устройство компонента — живой
документ; большое полотно плохо грузится и агентом, и человеком (ADR 0011).

Что: вычитание #84 слито с переустройством формы: набор
docs/architecture/orders/ — индекс README и семь файлов по частям устройства
(нарезка по правилу «семь плюс-минус два»); датированные файлы удалены,
ссылки перенацелены, AGENTS.md дополнен правилом подпапки. Приёмка
владельцем #88 пройдена, черновой статус снят из README.

Проверка: холодная сверка миграции свежим тредом — потерь решений нет;
обход относительных ссылок набора — битых нет.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-17 22:05:52 +03:00
ddadmin 3b4827d819 chore(git): каталог .scratch добавлен в игнор
Зачем: локальные рабочие файлы агентов (хэндоффы, обменные артефакты ревью)
живут в репозитории, чтобы переживать перезагрузку, но в git им не место.

Что: .scratch/ в .gitignore.

Проверка: git status не показывает .scratch/.
2026-08-17 22:05:52 +03:00
ddadminandClaude Fable 5 905ca5bbaf docs(orders): спека второго источника собрана из решений развилок
Зачем: решения развилок #70–#74, #80, #81 карты #69 разошлись по резолюциям
тикетов — цельная картина второго источника нужна в одном месте.

Что: черновик спеки заказов docs/specs/2026-08-16-orders.md и спека приёма
docs/specs/2026-08-16-order-ingestion.md; согласующие правки соседних
документов; правки горячего ревью и двух слепых холодных линий — дефекты
сборки закрыты, назван порядок строк внутри слепка.

Проверка: перепроверка находок теми же ревьюерами — ALL_CLOSED.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-17 22:05:23 +03:00
ddadmin ede1df765c docs(orders): зафиксирован версионный приём заказов
- Зачем:
  - отменённая подмена партиции snapshot_date противоречила порционному чтению Kafka и могла обучать потере ранее принятых версий.
- Что:
  - добавлены спецификация приёма заказов и ADR о версионном ODS с ods.order_v.
  - согласованы мастер-спека, дока хранилища, ADR 0008 и исследование формата.
  - зафиксированы граница приёма, координаты загрузки, диагностические повторы и отложенное проектирование DDS.
- Проверка:
  - git diff --cached --check.
  - горячее ревью по правилам репозитория и принятому решению.
  - два прохода холодного ревью.
2026-08-16 21:42:11 +03:00
ddadmin 8d54a3caba docs(orders): уточнён формат слепка и временные поля
- Зачем:
  - аудит источника нельзя смешивать с бизнес-временем и `_load_ts` ClickHouse.
- Что:
  - зафиксированы JSON-контракт слепка, порядок ключей и намеренное различие `Array(Float64)` и `Decimal`.
  - разведены бизнес-время, аудит источника и загрузка; уточнены слой ODS, `snapshot_date` и технический ключ партиции.
  - добавлены термины словаря и исследование точного миллисекундного формата с проверками ClickHouse.
- Проверка:
  - выполнен `git diff --cached --check`.
2026-08-16 17:56:46 +03:00
ddadminandClaude Opus 5 c17c5ef995 docs(domain): опись мира — вернуть счёт событий и паспорт
Зачем: определение описи разошлось с файлом, который называет. Формула
«счётчики закрытых окон с названным читателем» не покрывает счёт событий по
дням, а он в описи лежит и сегодня, и читатель у него самый живой: make
check-clickhouse сверяет по нему приход в ods.event.

Что: счёт событий вернулся в определение строкой дня; «паспорт» снова
раскрыт скобкой — двумя термами выше стоит «Паспорт куки», и голое слово с
ним сталкивалось; вернулось «единственное, что о мире хранится в git» —
определяющее свойство описи, а не следствие. Строка про слепок и правило
«счётчик только с названным читателем» сохранены.

Проверка: определение сверено с data/world-inventory.json и с резолюцией
«Места заказов в стартовом мире» (#74).

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-16 12:30:39 +03:00
ddadmin 86e2d6dbb2 docs(domain): уточнены слепки и опись мира
- Зачем:
  - документация не должна обещать постоянные числа для движущегося мира.
- Что:
  - уточнено появление опоздавших заказов в слепках.
  - опись ограничена неизменяемым выходом и счётчиками с читателем.
  - счётчик пар оставлен плану, но исключён из описи.
- Проверка:
  - CONTEXT.md перечитан целиком; git diff --check.
2026-08-16 12:26:10 +03:00
ddadminandClaude Opus 5 34267e8494 docs(spec): рамка — чем меряется генератор
Зачем: взгляд владельца на генератор жил только в разговорах, и каждый
следующий агент восстанавливал его заново или не восстанавливал вовсе.
Без него решения по генератору принимаются не по тому критерию: внутренняя
стройность вместо правдоподобия данных на выходе.

Что: новый раздел «Чем меряется генератор» сразу после «Зачем» — пять
утверждений рамки. Генератор меряется выходом, а не устройством; менти
внутрь не смотрит, поэтому конструкция оправдана только наблюдаемым
эффектом; повторимы эффекты, а не числа, и хеши описи обещают видимость
сдвига, а не постоянство чисел; цена ошибки мала, поэтому оборона дороже
ошибки не заводится; сложность изолирована службой ради дешевизны правки.

Рамка шире этапа 2, поэтому спеки следующих этапов ссылаются на неё, а не
копируют.

Проверка: правка только текстовая, кода не касается.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-16 12:09:11 +03:00
ddadmin c656ca198f docs(domain): уточнены роли человека и посетителя
- Зачем:
  - генератору нужен единый язык для скрытой личности и наблюдаемых идентификаторов.
- Что:
  - определены человек, посетитель и пользователь магазина.
  - закреплена граница между ClientID кликстрима и user_id заказа.
- Проверка:
  - документ перечитан целиком.
2026-08-16 11:56:57 +03:00
ddmitry f51048dd23 Merge pull request 'feat(airflow): пульт мира — работники, выключатель, контейнер' (#83) from feat/79-world-control into main
Reviewed-on: #83
2026-08-13 15:09:19 +03:00
ddadminandClaude Opus 5 292302e161 fix(airflow): правки пульта мира по двум линиям ревью
Зачем: линия дефектов нашла опору на дефект провайдера, линия постановки —
незаписанный ответ на вопрос ADR 0009 и переменную образца, чьё значение на
чужой машине неверно, а узнаёт об этом читатель через двести строк.

Что: `auto_remove` у прогона генератора переведён с `success` на `force`.
Значение `success` тоже убирало контейнер в любом исходе, но случайно —
удаление у провайдера 4.5.7 стоит в `finally` вопреки собственной
документации; `force` то же поведение называет прямо и переживёт починку.
У выключателя назван `catchup=False`: умолчание Airflow 3 то же самое, но
у дага с тиком в 25 минут и `start_date` в январе это первый вопрос
читателя. «Быстрый старт» предупреждает про `DOCKER_GID` — единственное
значение образца, неверное вне этой машины. В образе Airflow записано, что
провайдер docker приходит с базой (4.5.7 к 3.3.0), — это ответ на вопрос,
который ADR 0009 оставил тикету. Список томов планировщика получил ту же
пометку «правя одно, правьте второе», что стоит у числа дней. Проход на
вычитание срезал три комментария, пересказывавших ADR.

Проверка: `make lint`, `make config-test`, `make smoke` (20/0),
`make check-clickhouse` (9/9) — зелёные. Живой пульт на чистом стенде:
снятый с паузы `world_live` сыграл два дня подряд без нажатия (24 м 07 с и
24 м 06 с, пауза между ними 55 секунд), дочерний прогон виден ссылкой из
задачи выключателя; пауза остановила мир на границе суток — тик прошёл,
третьего прогона нет. Контейнер после прогона с `force` не остался.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-13 13:57:59 +03:00
ddadminandClaude Opus 5 35dc933d29 feat(airflow): пульт мира — работники, выключатель, контейнер
Зачем: модельный день прогонялся только руками, и владельцу нечем было
проверять процессы стенда вживую. Пульт нужен раньше этапа 3 и независимо
от него: он обкатывает то, что приёму заказов понадобится готовым — вызов
генератора из задачи Airflow.

Что: `dags/world_control.py` — три дага по ADR 0009. Работники
`world_next_day` (день пачкой, «сколько дней» параметром) и `world_live_day`
(день в темпе) живут без расписания и без паузы; выключатель `world_live`
создаётся на паузе, тикает раз в 25 минут и дёргает работника живого дня
с ожиданием конца. Генератор зовётся `DockerOperator` в каноническом
контейнере: сокет докера отдан планировщику, потому что при LocalExecutor
задачи исполняет он, а GID группы `docker` уехал в `.env` как локальная
настройка. Позицию на оси ведёт переменная `world_position` — её ставит
сыгравший день работник и только по успеху. Факты стенда — образ, сеть,
брокер, топик, размер стартового мира — даги получают окружением от compose;
внутри compose они названы по разу якорями, иначе разошлись бы с разовой
службой генератора. README получил раздел про пульт с названной вслух платой
за сокет.

Проверка: `make lint`, `make config-test`, `make smoke` (20/0), `make
check-services` (7/0), `make check-clickhouse` (9/9) — зелёные. На чистом
стенде: два прогона `world_next_day` подряд двигают позицию на два дня,
«дней = 3» — на три, все пять дней доехали в ODS; обрыв контейнера позицию
не двигает, повторный запуск играет тот же день с тем же счётом событий.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-13 13:09:14 +03:00
23 changed files with 1400 additions and 108 deletions
+5
View File
@@ -15,6 +15,11 @@ GRAFANA_PORT=23000
AIRFLOW_PORT=28080
SUPERSET_PORT=28088
# GID группы `docker` на этой машине: без него планировщик Airflow не достучится
# до сокета докера, и пульт мира не поднимет контейнер генератора.
# Подсмотреть свой — `getent group docker`.
DOCKER_GID=127
# Учётные данные и ключи учебного стенда. Для VPS замените значения образца.
GRAFANA_ADMIN_USER=admin
GRAFANA_ADMIN_PASSWORD=admin
+3
View File
@@ -44,3 +44,6 @@ Thumbs.db
# Локальные настройки Claude Code
.claude/*
!.claude/agents/
# Локальные рабочие файлы агентов: хэндоффы, обменные артефакты ревью
.scratch/
+4 -1
View File
@@ -158,4 +158,7 @@ Airflow) и названия из кода. Если для понятия ес
- Имена файлов в `docs/architecture/` — слаг строчными латинскими буквами через
дефис (`storage.md`). Здесь живут рабочие справочники по зонам
ответственности: не событие истории и не решение, а текущее устройство —
один файл на зону, правится по мере постройки.
один файл на зону, правится по мере постройки. Крупный компонент живёт
подпапкой: индекс `README.md` и файл на каждую часть устройства
(`architecture/orders/`), имена по тому же правилу
([ADR 0011](docs/adr/0011-component-docs.md)).
+41 -9
View File
@@ -34,6 +34,20 @@ _Избегать_: состояние мира
Люди, впервые пришедшие в мир в один день, со всеми их куками и днями
активности. Единица плана состава: когорта — функция зерна и номера дня.
**Человек**:
Устойчивая личность в составе мира, которой принадлежат одна или две куки.
Не поле источника: кликстрим видит куку, а заказ — пользователя магазина.
**Посетитель**:
Кука, наблюдаемая в кликстриме и обозначенная `ClientID`; единица
`uniq(ClientID)`. Один человек может быть представлен несколькими посетителями.
_Избегать_: «посетитель» про человека после склейки
**Пользователь магазина**:
Представление человека со стороны заказов; в записи заказа обозначается
`user_id`. Через заказы связывает разные куки человека, в кликстрим напрямую
не попадает.
**Приток**:
Появление новых кук на всём протяжении оси модельного времени; единица —
кука (`ClientID`). Из-за притока накопленная аудитория растёт с
@@ -102,10 +116,11 @@ ClickHouse. Форма файла решена, длина — нет: стро
_Избегать_: зерновой мир
**Опись мира**:
`data/world-inventory.json` — единственное, что о мире хранится в git:
паспорт (зерно, версия генератора, хеш каталога) и по строке на день с
датой, числом событий и хешем его байтов. Сам мир в git не лежит — он
пересчитывается. Опись отвечает на один вопрос: тот ли это мир.
`data/world-inventory.json` — единственное, что о мире хранится в git: паспорт
мира (зерно, версия генератора, хеш каталога), строка на каждый день событий
(число и хеш байтов) и строка на каждый отправленный слепок (хеш байтов).
Счётчик сверх этого заводится только с названным читателем. Сам мир не
хранится, а пересчитывается; опись отвечает на один вопрос: тот ли это мир.
_Избегать_: манифест, мини-манифест
**Пошаговый режим**:
@@ -114,7 +129,8 @@ _Избегать_: манифест, мини-манифест
**Живой день**:
Проигрывание текущего модельного дня в реальном времени с ускорением;
включается по требованию, не постоянный фон.
включается по требованию постоянным фоном идёт, только пока включён
выключатель пульта.
**Пакетный режим**:
Проигрывание готового дня пачкой, без темпа: заливка снимка при старте
@@ -146,8 +162,8 @@ Python-модуль с описателями колонок события —
_Избегать_: описание схемы, документация контракта
**Канонический сериализатор**:
Единственное место, где событие превращается в байты. Фиксированный порядок
ключей и строк — основа побайтовой воспроизводимости.
Единственная граница, где запись генератора превращается в байты. Фиксированный
порядок ключей и строк — основа побайтовой воспроизводимости.
**Проигрыватель**:
Компонент доставки готового потока дня в приёмник. Два режима: пакетный
@@ -155,8 +171,24 @@ _Избегать_: описание схемы, документация кон
**Слепок**:
Полная выгрузка заказов окна изменяемости, снятая бэкендом на границе суток:
состояние заказов на этот момент, а не поток их изменений. Один заказ
приезжает в стольких слепках, сколько дней окна он прожил.
состояние заказов на этот момент, а не поток их изменений. Опоздавший заказ
может отсутствовать в ранних слепках; после первого появления ездит до конца
своего окна.
**Дата слепка**:
Дата завершившегося модельного дня, состояние которого снято на исходящей
границе суток. Одна для всех записей выгрузки; не дата заказа и не время
отправки сообщения.
**Время аудита источника**:
Время создания или последнего изменения строки заказа по часам базы источника.
В записи слепка это `created_at` и `updated_at`; бизнес-время покупки живёт в
событии `purchase` отдельно.
**Судьба заказа**:
Исход, моменты изменений, вычеркнутая позиция и опоздание заказа — всё
решается бросками при его рождении и укладывается в окно изменяемости целиком.
Слепок любого дня — чтение готовой судьбы, а не накопление состояния.
**Окно изменяемости**:
Сколько модельных дней заказ ещё может измениться и потому продолжает ездить
+31 -1
View File
@@ -50,7 +50,10 @@ make check-services
Дальше, в рабочей петле, обычно хватает `make smoke`.
`.env.example` — образец обязательной локальной настройки. В нём живут имя
экземпляра, внешние порты, учебные учётные данные и ключи. `compose.yaml` не
экземпляра, внешние порты, учебные учётные данные и ключи. Одно значение образца
верно не везде — `DOCKER_GID`: у группы `docker` на каждой машине свой номер, и
свой подскажет `getent group docker`. Пока он не сойдётся, стенд работает весь,
кроме пульта мира — о нём ниже. `compose.yaml` не
дублирует их умолчаниями: если значения нет в `.env`, Compose сразу предложит
скопировать образец. Внутренние адреса, имена топиков и версии образов описывают
сам стенд, поэтому записаны литералами в `compose.yaml` и Dockerfile.
@@ -219,6 +222,33 @@ uv run --project generator python -m clickstream_generator batch \
код возврата, ещё до первого события. Поэтому умолчание дня и живёт в
`Makefile`: день называет тот, кто запускает, а не тот, кого запускают.
### Пульт мира
Считать дни самому необязательно. В Airflow живёт **пульт мира** — три дага,
которые зовут тот же образ генератора и ведут позицию на оси за вас.
| Даг | Что делает |
| --- | --- |
| `world_next_day` | играет следующие дни пачкой; сколько — параметр запуска, по умолчанию один |
| `world_live_day` | играет следующий день в темпе модельного времени, около двадцати четырёх минут |
| `world_live` | выключатель: пока снят с паузы, дёргает `world_live_day` день за днём |
Позицию хранит переменная Airflow `world_position` — номер первого несыгранного
дня. Ставит её работник, сыгравший день, и только по успеху: оборванный прогон
позицию не двигает, и следующий запуск играет тот же день заново. Нет
переменной — мир в стартовом состоянии.
Выключатель создаётся на паузе. Снимите — мир поедет сам; поставите обратно —
встанет на границе модельных суток, доиграв начатый день. Форма пульта и доводы
целиком — [ADR 0009](docs/adr/0009-world-control.md).
**Плата названа вслух: планировщику Airflow отдан сокет докера** — иначе
контейнер генератора ему не поднять. Доступ к сокету равен праву root на
машине; для локального учебного стенда размен принят, но знать о нём надо.
Открывает дверь не монтирование, а членство в группе: GID группы `docker` у
каждой машины свой, живёт в `.env` и подсматривается командой
`getent group docker`.
## Состав и доступ
- `clickhouse-01` — инициатор DDL и точка подключения Airflow;
+38 -5
View File
@@ -29,6 +29,17 @@ x-clickhouse-common: &clickhouse-common
retries: 30
start_period: 10s
# Факты стенда, у которых стало по два потребителя: разовая служба генератора и
# пульт мира. Названы по одному разу — иначе однажды разойдутся, и заметит это
# не проверка, а менти с пустым топиком.
x-generator-image: &generator-image clickstream-generator:local
x-kafka-target: &kafka-target
KAFKA_BOOTSTRAP_SERVERS: kafka:9092
KAFKA_TOPIC: hits
x-world-starting-days: &world-starting-days "8"
x-airflow-common: &airflow-common
image: clickstream-airflow:local
build:
@@ -50,6 +61,13 @@ x-airflow-common: &airflow-common
AIRFLOW_ADMIN_USER: ${AIRFLOW_ADMIN_USER:?Скопируйте .env.example в .env}
AIRFLOW_ADMIN_PASSWORD: ${AIRFLOW_ADMIN_PASSWORD:?Скопируйте .env.example в .env}
CLICKHOUSE_ETL_PASSWORD: ${CLICKHOUSE_ETL_PASSWORD:?Скопируйте .env.example в .env}
# Пульт мира поднимает генератор сам, отдельным контейнером, — значит те же
# факты стенда, что compose даёт разовой службе генератора, нужны и дагам.
# Имя сети собирается из имени проекта: у второй копии стенда оно другое.
<<: *kafka-target
GENERATOR_IMAGE: *generator-image
STAND_NETWORK: ${COMPOSE_PROJECT_NAME}_default
WORLD_STARTING_DAYS: *world-starting-days
volumes:
- ./dags:/opt/airflow/dags:ro
- ./infra/airflow/init.sh:/opt/airflow/init.sh:ro
@@ -57,7 +75,7 @@ x-airflow-common: &airflow-common
- airflow_auth:/opt/airflow/auth
x-generator-common: &generator-common
image: clickstream-generator:local
image: *generator-image
build:
context: .
dockerfile: generator/Dockerfile
@@ -75,8 +93,7 @@ x-generator-common: &generator-common
# Адрес брокера и имя топика — факты стенда, и называет их стенд.
# Остальное (день, зерно, число дней, предел пачки) приходит аргументами
# от того, кто запускает: у службы нет позиции на оси мира.
KAFKA_BOOTSTRAP_SERVERS: kafka:9092
KAFKA_TOPIC: hits
<<: *kafka-target
x-superset-common: &superset-common
image: clickstream-superset:local
@@ -261,14 +278,15 @@ services:
# Число дней стоит здесь числом: YAML не читает Python, и одно из двух мест
# (второе — `STARTING_DAYS` в inventory.py) лишнее по построению. Правя одно,
# правьте второе — на страже тут никто не стоит: залей эта служба лишний
# день, он лёг бы за рамкой дат описи и остался бы незамеченным.
# день, он лёг бы за рамкой дат описи и остался бы незамеченным. Внутри YAML
# число одно на всех: то же говорит дагам пульта, где кончается стартовый мир.
#
# Повторный `make up` заливает мир заново, и это не оплошность: `WatchID` у
# событий те же, ReplacingMergeTree схлопнет повтор в ODS. Сырьё в STG при
# этом честно удвоится — свойство слоя, описанное в storage.md.
world-init:
<<: *generator-common
command: ["batch", "--day", "0", "--days", "8"]
command: ["batch", "--day", "0", "--days", *world-starting-days]
# Тот же образ для ручных прогонов: `make generate-batch`, `make generate-live`.
# Под профилем — чтобы обычный подъём стенда её не трогал.
@@ -350,6 +368,21 @@ services:
airflow-init:
condition: service_completed_successfully
command: scheduler
# Пульт мира зовёт генератор отдельным контейнером, а при LocalExecutor
# задачи исполняет сам планировщик — значит сокет докера нужен ему одному.
# Плата названа вслух в README и ADR 0009: доступ к сокету равен праву root
# на машине. Открывает дверь не монтирование, а группа: у сокета права 660
# и группа `docker`, чей GID на каждой машине свой и живёт в `.env`.
group_add:
- ${DOCKER_GID:?Скопируйте .env.example в .env}
# Тома перечислены заново: список службы общий не дополняет, а заменяет.
# Появится новый том у остальных служб Airflow — добавьте и сюда.
volumes:
- /var/run/docker.sock:/var/run/docker.sock
- ./dags:/opt/airflow/dags:ro
- ./infra/airflow/init.sh:/opt/airflow/init.sh:ro
- airflow_logs:/opt/airflow/logs
- airflow_auth:/opt/airflow/auth
mem_limit: 640m
healthcheck:
test: ["CMD-SHELL", "curl -sf http://127.0.0.1:8974/health >/dev/null"]
+183
View File
@@ -0,0 +1,183 @@
"""Пульт мира: даги, которыми двигают ось модельного времени.
Их три, и они двух родов. **Работники** играют день и запускаются руками:
`world_next_day` пачкой, без пауз, сколько дней попросили; `world_live_day`
один день в темпе модельного времени. **Выключатель** `world_live` своей работы
не делает: он тикает по расписанию и дёргает работника живого дня, дожидаясь
конца. Снят с паузы мир едет день за днём; поставлен на паузу встал на
границе модельных суток.
Разделение не косметическое. Расписание на самом работнике заставило бы кнопку
паузы значить две вещи разом «мир не едет сам» и «даг выключен», а работник
при этом выглядел бы в списке выключенным, хотя нажимают его каждый день.
Календарь Airflow к оси мира отношения не имеет: какой день играть, работник
спрашивает у переменной, а не у логической даты прогона.
Решения и доводы целиком ADR 0009.
"""
from __future__ import annotations
import datetime
import os
from airflow.providers.docker.operators.docker import DockerOperator
from airflow.providers.standard.operators.trigger_dagrun import TriggerDagRunOperator
from airflow.sdk import Param, Variable, dag, get_current_context, task
# Факты стенда — образ генератора, сеть, адрес брокера, топик и размер
# стартового мира — приходят окружением, и называет их compose: тот же, что
# называет их разовой службе генератора.
GENERATOR_IMAGE = os.environ["GENERATOR_IMAGE"]
STAND_NETWORK = os.environ["STAND_NETWORK"]
GENERATOR_ENVIRONMENT = {
"KAFKA_BOOTSTRAP_SERVERS": os.environ["KAFKA_BOOTSTRAP_SERVERS"],
"KAFKA_TOPIC": os.environ["KAFKA_TOPIC"],
}
STARTING_DAYS = int(os.environ["WORLD_STARTING_DAYS"])
# Позиция на оси: номер первого несыгранного дня. Переменной нет — мир в
# стартовом состоянии, и играть надо сразу за ним.
#
# Позиция ставится, а не увеличивается. Наложись один прогон на другой, худшее
# при таком правиле — сыгранный дважды день, а повтор схлопнет
# ReplacingMergeTree. Увеличение молча съело бы день, и в мире осталась бы
# дыра, которой никто не заметит.
WORLD_POSITION = "world_position"
# Тик выключателя. Каденцию задаёт не он, а сама длина живого дня — около
# двадцати четырёх минут: тик только спрашивает «не пора ли снова».
LIVE_TICK = datetime.timedelta(minutes=25)
START_DATE = datetime.datetime(2026, 1, 1, tzinfo=datetime.UTC)
TAGS = ["пульт мира"]
@task
def first_unplayed_day() -> int:
"""Номер дня, с которого играть."""
return int(Variable.get(WORLD_POSITION, default=STARTING_DAYS))
def _play(task_id: str, command: list[str]) -> DockerOperator:
"""Задача, играющая дни в каноническом контейнере генератора.
Внутрь образа Airflow генератор не поставить: он требует Python 3.14, а
образ несёт 3.13. Да и обещание побайтовой воспроизводимости дано для
зафиксированного образа генератора держится оно только там.
"""
return DockerOperator(
task_id=task_id,
image=GENERATOR_IMAGE,
command=command,
network_mode=STAND_NETWORK,
environment=GENERATOR_ENVIRONMENT,
# Контейнер убирается в любом исходе: вывод генератора оператор уже
# перелил в журнал задачи, а мёртвые контейнеры копить незачем.
auto_remove="force",
# По умолчанию оператор монтирует контейнеру временный каталог. Здесь
# это ловушка: путь он заводит внутри Airflow, а монтирует демон с
# хоста, где такого пути нет. Генератору временный каталог не нужен.
mount_tmp_dir=False,
)
@dag(
dag_id="world_next_day",
schedule=None,
start_date=START_DATE,
is_paused_upon_creation=False,
max_active_runs=1,
tags=TAGS,
params={"days": Param(1, type="integer", minimum=1, title="Сколько дней прожить")},
)
def world_next_day():
"""Прожить следующие дни пачкой, без пауз.
Запускается руками. День по умолчанию один, но разгон вперёд идёт одним
нажимом, а не десятью: сколько дней играть параметр запуска.
"""
@task
def remember_played(first_day: int) -> None:
"""Позиция ставится по сыгранным дням и только по успеху."""
days = get_current_context()["params"]["days"]
Variable.set(WORLD_POSITION, str(first_day + days))
first_day = first_unplayed_day()
played = _play(
"play_days",
[
"batch",
"--day",
"{{ ti.xcom_pull(task_ids='first_unplayed_day') }}",
"--days",
"{{ params.days }}",
],
)
first_day >> played >> remember_played(first_day)
@dag(
dag_id="world_live_day",
schedule=None,
start_date=START_DATE,
is_paused_upon_creation=False,
max_active_runs=1,
tags=TAGS,
)
def world_live_day():
"""Прожить следующий день в темпе модельного времени.
Ускорение ×60: модельные сутки укладываются примерно в двадцать четыре
реальные минуты, и суточная волна разворачивается на глазах. Запускается
руками; чтобы мир жил так день за днём сам, есть выключатель `world_live`.
"""
@task
def remember_played(first_day: int) -> None:
"""Позиция ставится по сыгранному дню и только по успеху."""
Variable.set(WORLD_POSITION, str(first_day + 1))
first_day = first_unplayed_day()
played = _play(
"play_day",
["live", "--day", "{{ ti.xcom_pull(task_ids='first_unplayed_day') }}"],
)
first_day >> played >> remember_played(first_day)
@dag(
dag_id="world_live",
schedule=LIVE_TICK,
start_date=START_DATE,
catchup=False,
is_paused_upon_creation=True,
max_active_runs=1,
tags=TAGS,
)
def world_live():
"""Выключатель: пока включён, мир живёт день за днём.
Своей работы у выключателя нет он дёргает `world_live_day` и ждёт конца.
Ожидание тут несущая конструкция, а не вежливость: без него тик шёл бы
независимо от хода дня, лишние прогоны скопились бы очередью, и мир потом
промчался бы по ней без всякого темпа.
Ждём триггером, а не сенсором: оператор опрашивает тот прогон, который сам
и создал, и ссылка на дочерний прогон видна прямо отсюда. Упал день
краснеет и выключатель.
"""
TriggerDagRunOperator(
task_id="trigger_live_day",
trigger_dag_id="world_live_day",
wait_for_completion=True,
)
world_next_day()
world_live_day()
world_live()
+19 -4
View File
@@ -1,6 +1,13 @@
# ADR 0008. Приём заказов: пакетный забор слепка, инициируемый Airflow
Дата: 12 августа 2026 года. Статус: принято. Реализация — отдельным тикетом.
Дата: 12 августа 2026 года. Статус: частично заменено
[ADR 0010](0010-order-versions-in-ods.md). Реализация — отдельным тикетом.
Сохраняются пакетный забор, `RawBLOB`, один чтец на `clickhouse-01`, одна
партиция топика и отсутствие матвью. Отменены граница «одно чтение — полный
слепок», замена партиции `snapshot_date`, ODS без дедупликации и готовая схема
`dds.order` с `argMax`. Действующая форма STG → ODS описана в
[спецификации приёма заказов](../architecture/orders/ingestion.md).
## Решение
@@ -10,6 +17,11 @@
прочитанное в `stg.orders_raw_dist`, разбирает его в типизированный слепок и
заменяет партицию дня в `ods.order_snapshot`. Тот же даг проигрывает модельный
день генератором, поэтому переливается ровно то, что он положил в топик.
Уточнение со сдвигом отправки, решённым позже (#71, [слепок и его
доставка](../architecture/orders/snapshot.md)): даг, играющий день D, в
штатном прогоне кладёт и забирает слепок дня D−1 — слепок предыдущего дня, а
не сыгранного. После падения между шагами в топике может ждать и хвост
прежних слепков; забор принимает всё приехавшее.
Прямое чтение из Kafka-движка требует двух настроек, и вторая не очевидна:
`stream_like_engine_allow_direct_select = 1` разрешает читать чтеца запросом,
@@ -128,7 +140,9 @@
## Что проверено
По документации ClickHouse через MCP Context7, 12 августа 2026 года.
Документация ClickHouse проверена через MCP Context7 12 августа 2026 года.
Предел порции одного опроса Kafka дополнительно снят 16 августа на локальном
ClickHouse `26.3.17.56`.
- Прямое чтение из движков-очередей (Kafka, RabbitMQ, FileLog) запрещено
начиная с версии 21.12 и открывается настройкой
@@ -138,8 +152,9 @@
- Прямое чтение офсеты по умолчанию **не** коммитит; коммит включается
настройкой `kafka_commit_on_select` на самой таблице. Это тот подводный
камень, который молчит на первом прогоне и вылезает на втором.
- Сколько строк отдаёт одно чтение, задаёт `kafka_max_block_size`. При слепке
порядка полутора тысяч строк это один блок с запасом.
- Прямое чтение возвращает одну порцию, полученную одним опросом Kafka. При
настройках стенда её предел — 65 409 сообщений, поэтому слепок порядка
полутора тысяч строк помещается с запасом.
- Про `Distributed` поверх Kafka документация не говорит ничего — ни
поддержки, ни запрета.
+37
View File
@@ -0,0 +1,37 @@
# ADR 0010. Заказы в ODS: версии сущности вместо подмены слепка
Дата: 16 августа 2026 года. Статус: принято. Частично заменяет
[ADR 0008](0008-order-ingestion.md).
## Решение
Пакетный забор заказов остаётся прямым чтением байтового Kafka-чтеца по
команде Airflow, но порция чтения больше не считается полным слепком и не
публикуется заменой партиции `snapshot_date`. Прочитанные строки получают
`_load_id` запуска, разбираются из одного среза STG в годные строки и ошибки,
а `ods.order_snapshot` хранит принятые версии заказа в
`ReplacingMergeTree(updated_at)`. Ключ сущности — `order_id`, партиция строится
от неизменного `created_at`, все версии ключа направляются на один шард.
`snapshot_date` остаётся датой наблюдения источника, `_load_id` — координатой
запуска приёма, `_load_ts` — временем прибытия строки. Ни одна из них не
заменяет бизнес-версию `updated_at`. Физическая таблица может показывать
несколько версий; точное текущее состояние ODS открывает `ods.order_v`, которое
скрывает `FINAL` или равносильный способ выбора последней версии.
Модель заказа в DDS этим решением не задаётся. DDS получает устойчивую
типизированную поверхность ODS и отдельно решает зерно, связи и способ
материализации своей модели.
## Почему отменена подмена партиции
Прямой `SELECT` Kafka Engine заканчивается после одной полученной порции, а
протокол не несёт признака конца слепка. Поэтому `snapshot_date` не доказывает,
что в STG собрана полная партиция, и её подмена могла бы удалить уже принятые
версии прошлого дня. Маркер конца, опись или фиксация конечных офсетов сделали
бы границу настоящей, но добавили бы новый протокол без нужного стенду урока.
Малый объём позволяет оставить одно чтение на запуск как проверяемое
эксплуатационное допущение, а не границу полноты. Полный контракт разбора,
граница брака и поведение повторов заданы в
[спецификации приёма заказов](../architecture/orders/ingestion.md).
+42
View File
@@ -0,0 +1,42 @@
# ADR 0011. Устройство компонента — связный набор живых документов
Дата: 16 августа 2026 года. Статус: принято.
## Решение
Детальное устройство крупного компонента живёт в `docs/architecture/`
подпапкой: индекс `README.md` — целевая картина и указатели — и отдельный
файл на каждую часть устройства. Имена — слаги без дат; набор правится по
мере постройки, как и остальные справочники этой папки.
Первый такой набор — [заказы бэкенда](../architecture/orders/README.md),
собранный из спеки этапа 3 и спеки приёма. Датированные файлы
`2026-08-16-orders.md` и `2026-08-16-order-ingestion.md` удалены, ссылки на
них перенацелены; историю держит git.
## Почему
**Спека и базовая документация — разные жанры, а файл был один.** Спека —
событие: проект изменения с датой в имени, замерзающий после приёмки. Но
собранная спека заказов сразу стала и детальным устройством сервиса — тем
документом, по которому этап 3 будут строить и с которым потом сверяться.
Живому устройству дата в имени врёт, а замерзать ему нельзя.
**Один большой документ плохо читается обоими читателями.** Агент, строящий
судьбу заказа, вынужден везти в контексте формат провода и приём; человек
листает пятьсот строк ради одного раздела. Индекс с файлами по частям даёт
обоим одно и то же: загружается только нужная часть, а карта целого — один
экран.
**Долговечно только то, что в git.** Резолюции развилок живут в трекере, а
трекер долговечным хранилищем не считается. Поэтому решения — вместе с
отклонёнными вариантами при каждом правиле — лежат в файлах набора; ссылки
на тикеты остаются вежливостью, не записью.
## Следствия
- `docs/specs/` остаётся событиям: мастер-спека и спека генератора живут как
есть; переводить ли их в живую форму — отдельное решение, когда встанет.
- Раздел «Структура» в AGENTS.md обновлён тем же коммитом.
- Ссылки из ADR 0008 и 0010, мастер-спеки, спеки генератора, исследования
формата и доки хранилища перенацелены на набор.
+75
View File
@@ -0,0 +1,75 @@
# Заказы бэкенда
Устройство второго источника стенда: раз в модельный день бэкенд магазина
выгружает полный слепок заказов окна изменяемости, и деньги в витринах
считаются по нему, а не по трекеру. Набор собран картой #69 (этап 3) и
правится по мере постройки.
Границы уже решены мастер-спекой [«Боевой реализм стенда
(v2)»](../../specs/2026-07-30-stand-v2-realism.md): поля слепка, окно K = 7
как константа мира, три статуса, приоритет классов расхождений, правило
«поведение и атрибуцию считаем по трекеру, деньги — по бэкенду». Рамка, перед
которой отвечает каждое решение, — «Чем меряется генератор» в [спеке
генератора](../../specs/2026-08-01-generator.md): конструкция внутри
оправдана только наблюдаемым эффектом на выходе.
## Целевая картина одним взглядом
- **Заказ — проекция, не порождение.** Торговая половина дня-функции уже
посчитала корзину, цены, купон и номер заказа; заказная половина навешивает
судьбу и собирает слепок. Ни одного нового броска в торговом подпотоке —
[откуда берётся заказ](snapshot.md).
- **Слепок дня D — чистая функция (зерно, D)**: состояние заказов, рождённых
в дни D−6…D, снятое на границе суток D|D+1. Отправляет его следующий
прогон — ночная выгрузка бэкенда за вчера; на проводе — один JSON-документ
на заказ — [слепок и его доставка](snapshot.md).
- **Судьба заказа решается при рождении** и обязана уложиться в окно K либо
не случиться вовсе. На выходе из окна заказ либо `paid`, либо `cancelled`
[судьба заказа](fate.md).
- **Расхождения и опоздания — часть мира, а не грязь**: два подпотока —
заказная и событийная стороны; броски независимы, пересечения выходят
арифметикой, приоритет классов работает по-настоящему —
[классы расхождений и опоздание](fate.md).
- **Опись хранит только то, чего движение мира не меняет**: хеш байтов
каждого слепка и счётчики наблюдаемых классов —
[что хранит опись](inventory.md).
- **Мост к склейке**: план состава владеет человеком; его непрозрачный
`person_id` заказ показывает как `user_id`, кликстрим остаётся анонимным —
[мост к склейке](identity.md).
- **Приём — пакетный забор**: одно прямое чтение Kafka в STG, два
`INSERT SELECT` в типизированный ODS и таблицу ошибок; `ods.order_snapshot`
принимает версии заказа на `ReplacingMergeTree(updated_at)`
[приём из Kafka в ODS](ingestion.md).
- **Стартовый мир отправляет семь слепков** (дни 0…6): у последнего прожитого
дня клики есть, а заказов нет, и график выручки дозаполняется по ходу
мира — [заказы в стартовом мире](start-world.md).
Правила, обязательные для кода этапа 3, собраны в
[правилах кода](code-rules.md).
## Открытые решения
- **Имена подпотоков сторон судьбы** — при реализации; из мёртвых имён никто
не бросает, переименование ничего не сдвигает (#72).
- **Конкретные веса и доли** — таблицы исходов, моментов, задержки опоздания,
доли классов, стоимость доставки — калибровка при реализации; финальная
фиксация чисел — пересборка эталонного мира, этап 7. При пересборке правки
потребуют только числа, не устройство.
- **Проверки приёма** — разовая приёмка допущения «один запуск — одно чтение»
и опыты из [«Рисков и проверки»](ingestion.md) — тикет реализации приёма.
- **`_load_id` выше ODS** — вместе с устройством `dds.order` (#85).
- **Контур проверок качества для расхождений** (даг DQ, `dm.dq_summary`) —
остаётся в тумане карты #69; естественное место разговора — этап 4.
- **Каноническое чтение событий `ods.event_v`** — отдельный тикет #86, к
механике заказов не привязан.
## Родословная
Собрано тикетом #87 по резолюциям развилок карты #69: приём (#70), генератор
слепков (#71), расхождения и опоздания (#72), мост к склейке (#73), место в
стартовом мире (#74), брак и версии в ODS (#80), форма записи на проводе
(#81). Решения о приёме — [ADR 0008](../../adr/0008-order-ingestion.md) и
[ADR 0010](../../adr/0010-order-versions-in-ods.md); контракт провода —
мастер-спека, раздел 2, и [исследование формата
слепка](../../research/2026-08-16-order-snapshot-wire-format.md). Форма
набора — [ADR 0011](../../adr/0011-component-docs.md).
+17
View File
@@ -0,0 +1,17 @@
# Правила кода этапа 3
Хвосты резолюций, обязательные для реализации:
- **Броски заказной стороны — на полную длину дня**, а не на отобранных
заказах (правило формы, [судьба заказа](fate.md)); моментов всегда два.
- **Целочисленная случайность** наследуется правилом кода этапа 2: таблицы
целых весов, никаких плавающих распределений; деньги — в целых копейках.
- **Подпотоки — по позиции в дереве**: стороны судьбы ветвятся по дню рождения
заказа; в `COMMERCE` новых бросков нет; `person_id` — последний бросок
когорты.
- **Окно K = 7 — константа мира** в конфигурации мира, рядом с D0 и поясом
([исследование
формата](../../research/2026-08-16-order-snapshot-wire-format.md)).
- **Один сериализатор**: запись слепка собирает явная функция `serialize.py`,
прямых `json.dumps` по коду нет.
- **Значения идентификаторов — ниже 2^53** (`user_id` наравне с прочими).
+140
View File
@@ -0,0 +1,140 @@
# Судьба заказа и расхождения
Резолюция развилки [«Расхождения A–D и опоздания: механика, доли и что
обещано»](https://git.dementev.space/ddmitry/clickstream-data-platform/issues/72).
## Судьба заказа
**Рамка.** Расхождение — не грязь и не шум, а часть мира: судьба заказа,
решённая при его рождении и уложенная в окно K целиком.
**Два подпотока дня.** Заказная сторона — судьба заказа: исход, моменты,
дельта, опоздание. Событийная сторона — порча событийного потока: потеря и
дубль. Довод за разделение — различимость по описи: правка заказной механики
не двигает хеши событий, правка событийной не двигает байты слепка, и по
покрасневшим хешам видно, какую сторону трогали. Прежние имена `DISCREPANCIES`
и `LATECOMERS` решения не переживают — они названы по классам витрины, а
компонент называет часть мира; новые стороны занимают те же позиции, имена —
при реализации. Отклонено: *один компонент на всю судьбу* — правка событийной
механики молча меняла бы байты слепка; *компонент на класс* — пять имён под
ручки калибровки, которые крутятся разом.
**Правило формы, без которого разделение не работает: броски заказной стороны
делаются на полную длину дня, а не на отобранных заказах.** Иначе длина броска
становится функцией доли, и правка одной доли перебрасывает весь подпоток
после себя. Изоляцию даёт форма броска, а не число подпотоков.
**Путь по статусам.** Три исхода, все внутри окна: **оплачен**; **оплачен и
отменён**; **не оплачен и отменён**. Инвариант на выходе из окна: заказ либо
`paid`, либо `cancelled`; `created` — только промежуточное состояние. За окном
будущего у заказа нет, а заказ, навсегда застрявший в `created`, — модельная
небрежность, которой в выгрузке живого магазина соответствия нет. Поэтому
нового значения `mismatch_class` не нужно: `cancelled` покрывает обе дороги
отмены (шестое значение занято `awaiting_order`).
**Форма броска.**
1. **Исход** — таблица долей из трёх строк; доля неоплаченных пишется явной
строкой, а не оставляется читателю складывать хвост в уме.
2. **Моменты** — таблица целых весов «сколько часов от рождения — с каким
весом», строки 0…143, плюс равномерная секунда внутри часа — чтобы разности
времён аудита не давали точных равенств (урок правки #50).
3. Моментов бросается **всегда два, на полную длину дня**; у одномоментных
исходов второй выбрасывается. Где их два по существу, ранний считается
оплатой — порядок выходит сортировкой, условной точки отсчёта не нужно.
143 часа — самый узкий край окна: у заказа, рождённого в конце суток, до
последнего его слепка 144 часа. Таблица кончается там, где кончается окно у
самого невезучего: вылезти нечему, сторожа не нужно. Цена — заказ, рождённый в
начале суток, не использует почти сутки своего окна; в хвосте таблицы веса
мизерные, в данных это не видно. Форма из #71 — «вес за краем окна означает
„не оплачен никогда“» — этим отменена: такой заказ теперь отменяется, а край
окна не выражается числом часов — иначе доля неоплаченных стала бы функцией
часа покупки, и менти нашёл бы этот наклон первым же разрезом.
**Доли.** Ориентир мастер-спеки ~5% читается как доля отменённых вообще; как
она делится между двумя дорогами — строки таблицы исходов. Точные числа —
калибровка этапа 7; проверок вида «отмен от 4 до 6 процентов» не заводим.
**Что из двух дорог видно.** В `dds.order` дороги неразличимы — там последняя
версия; различает их история версий: сырьё STG и физические версии
`ods.order_snapshot` до фоновых слияний, а в витринах — выручка дня, которая
сначала выросла, потом убыла. Полное различение не обещано: слепок — состояние
на границе суток, и оплата с отменой в один день в сырье неразличимы; то же у
сильно опоздавших, приехавших уже терминальными. Точную форму этого урока
решает этап 4. Отклонено: *отмена только после оплаты* — неоплаченному некуда
деться, кроме как остаться брошенным; *мгновенная отмена при рождении*
«дыхание» окна на отменах исчезает; *отмен нет вовсе* — страховочный срез 1,
он в резерве.
## Классы расхождений и опоздание
Классы и ориентиры долей — [мастер-спека,
раздел 4](../../specs/2026-07-30-stand-v2-realism.md); здесь — механика
каждого.
**Дельта суммы (C) — вычеркнутая позиция.** Товара не оказалось в наличии,
позицию сняли: у заказа на одну позицию меньше, чем в клиентских массивах, а
`items_total` меньше на её стоимость. Момента у неё нет — заказ приезжает
урезанным во всех своих слепках: первый слепок снимается на границе суток,
когда склад заказ уже собрал. Заказ из одной позиции дельты не получает —
пустых заказов не бывает. Позиция выбирается равновероятно: корреляция со
спросом на доле 1–2% статистически ненаблюдаема — менти платил бы за неё
таблицей чисел мира, а увидеть не мог бы ничем. Доводы за вычёркивание:
остаток — сотни рублей, он торчит в витрине сверки сам; расхождение
объясняется сравнением позиций — разбором вложенного JSON и `ARRAY JOIN`,
ровно тем навыком, ради которого позиции разбираются; история рассказывается
словами без легенды про генератор. Отклонено: *переоценка позиции* и *другое
количество* — дельта в десятки рублей, её надо захотеть заметить; *чистая
дельта без истории* — тупик, объяснить нечем; *врёт клиент, а не бэкенд*
заказ у нас проекция той же корзины.
**Потеря события (B) — точечная.** Уходит строка `purchase`, просмотр
`/confirmation` остаётся: события уезжают разными запросами, потерять один и
сохранить другой — обычное дело. Единственный вариант, при котором потеря
видна со стороны трекера: до подтверждения дошли сто, покупок девяносто семь.
**Дубль события (D) — сюжетный.** Обновление страницы шлёт и просмотр, и
покупку. Довод не в связности легенды: точечный дубль ломал бы урок соседнего
класса — разрыв воронки, на котором держится потеря, сжался бы втрое; при
сюжетном разрыв снова равен доле потерь. `purchaseID`, суммы и позиции у дубля
один в один — посчитал наивно, удвоил выручку. Дубль случается только там, где
до следующего визита куки остаётся запас сверх таймаута: иначе сборка сессий у
менти разошлась бы с `VisitID` — сломался бы эталон, ради которого `VisitID`
в потоке лежит. Задержка дубля — секунды-минуты, короче таймаута визита,
поэтому `VisitID` тот же. Полночь режет дубль парой — просмотр вместе с
покупкой, по тому же правилу, что у подтверждения с торговым хвостом.
Отклонено: *обе точечные* — дубль затирает урок потери; *обе сюжетные*
потеря перестаёт быть видна со стороны трекера.
**Гарантия моста сильнее порчи.** Назначенные планом покупки не теряются, и
назначенные планом заказы не опаздывают: `dds.identity_map` строится из моста
«`purchase` ↔ заказ», и выброшенное событие — как и заказ, не попавший ни в
один снятый слепок, — уносит куку из карты. Это был бы отказ лабы склейки, а
не расхождение в данных; менти различить не может.
**Опоздание — заказ прячется от ранних слепков.** `created_at` не
подделывается — строка создана, когда заказ родился, — но в слепках дней
d…d+δ−1 её нет, а с d+δ она появляется в том состоянии, до которого заказ
дожил: отменённый на второй день и опоздавший на третий приедет в первом же
своём слепке как `cancelled` — «выгрузка догоняет жизнь». Задержка — таблица
весов из трёх строк: 0 на подавляющем весе, 1 и 2 — это и есть «D+1/D+2»
мастер-спеки. Меряется она в днях снятия слепка, а не отправки: иначе сдвиг
отправки удвоился бы, и обещанные D+1/D+2 стали бы D+2/D+3. Дальше таблица не
идёт: заказ с δ = 6 приехал бы ровно в одном слепке, и обещание «пропущенный
день ничего не ломает» на нём перестало бы быть верным; при δ ≤ 2 у всякого
заказа слепков не меньше пяти. Легенда: заказ ушёл в ручную обработку и попал
в выгрузку позже. Отклонено: *сдвиг `created_at`* — подделка аудита источника:
день создания строки разошёлся бы с днём покупки, чьё равенство держит
синхронная модель ([откуда берётся заказ](snapshot.md)), заказ уехал бы в
чужую партицию, и «выручка дня D» перестала бы отвечать покупкам дня D —
сломалась бы та самая сверка, ради которой всё строится.
**Пересечения.** Броски независимы, пересечения выходят арифметикой, приоритет
мастер-спеки работает по-настоящему. Исключений два, и оба названы выше: дубль
решается только у выживших покупок, а назначенное планом не теряется и не
опаздывает — дельта и отмена ему разрешены, моста они не рвут.
Следствие для калибровки: брошенная доля и наблюдаемая в сверке — разные числа
(часть заказов забирают победители по приоритету, часть у класса C недоступна —
однопозиционных заказов больше половины). Отклонено: *один класс на заказ*
приоритет в SQL стал бы мёртвой веткой, которую менти читает как живую.
+37
View File
@@ -0,0 +1,37 @@
# Мост к склейке: человек, `person_id`, `user_id`
Резолюция развилки [«Мост к склейке: user_id и двухкуковые
пары»](https://git.dementev.space/ddmitry/clickstream-data-platform/issues/73).
Причинная модель: план состава порождает **человека**; ему принадлежат одна
или две куки, каждая наблюдается в кликстриме как отдельный посетитель. Визит
может оформить заказ — тогда заказная сторона представляет того же человека
как **пользователя магазина**. Это не модель аккаунта: регистрации нет, люди
без визитов не порождаются, внутреннее знание наружу не выдаётся.
План владеет устойчивой личностью и отношением «кука принадлежит человеку».
Минимальная форма — непрозрачный внутренний `person_id`, выровненный по кукам:
у двух кук пары он одинаков. Заказ выводит то же значение под родным именем
`user_id` (`UInt64`); в кликстрим ни `person_id`, ни `user_id` не попадает —
анонимность формата держится формой, а не забывчивостью сериализатора.
`person_id` — последний обычный бросок потока случайности когорты, после уже
принятых свойств: добавление личности не сдвигает куки, пары, возвраты и
паспорта. Для второй куки повторяется ID её человека. Раздельные прогоны
ничего не хранят и не согласуют: генератор событий и команда слепка заново
спрашивают один план и получают тот же `person_id`.
Наблюдаемое обещание — четыре свойства: назначенные заказы пары с разных
`ClientID` несут один `user_id`; событие Метрики не раскрывает `user_id`;
отдельные прогоны одного мира дают то же соответствие; принятый мир
воспроизводит эффект склейки без случайных ложных объединений — конкретные
числа пар контрактом описи не являются. Межкогортные столкновения принимаются
по той же дисциплине, что у случайных `ClientID`: принимается конкретный
канонический мир по внешнему результату.
Отклонено: *заказная сторона сама назначает личность* — родство кук всё равно
пришлось бы спрашивать у плана; *материализованный реестр кука↔пользователь*
состояние между прогонами без нового внешнего эффекта; *первая кука как ID
человека* — магазин оказался бы замаскированным продолжением трекера;
*структурная координата, хеш или глобальная последовательность* — больше
механики при тех же данных; *полная модель аккаунтов* — менти её не наблюдает.
+170
View File
@@ -0,0 +1,170 @@
# Приём заказов из Kafka в ODS
Учебный результат: менти различает версию бизнес-сущности, наблюдение источника
и запуск загрузки, а затем читает физические версии через явную поверхность
текущего состояния.
## Проблема
Заказы приезжают полным слепком окна изменяемости, но Kafka передаёт его
отдельными сообщениями и не сообщает потребителю, где слепок закончился. Прямое
чтение Kafka Engine возвращает одну порцию. Поэтому прежняя публикация через
`REPLACE PARTITION snapshot_date` приравнивала дату наблюдения к отсутствующей
транспортной границе и могла заменить день неполным набором строк.
Одновременно типизированный ODS не должен принимать правдоподобные значения по
умолчанию из грязного JSON или останавливать весь пакет из-за одной строки.
## Цели
- сохранить пакетный забор как контраст потоковому приёму событий;
- один раз принять байты в STG и независимо разложить строки на годные и брак;
- хранить в ODS типизированные версии заказов, не привязывая идемпотентность к
`snapshot_date`;
- дать следующим слоям один корректный способ прочитать текущее состояние;
- оставить код проверки коротким и ограничить его контрактом провода.
## Не входит
- модель заказа в DDS: её зерно, связи, материализация и способ наполнения;
- проверка полей внутри `items`, переходов статуса и равенств денежных сумм;
- удаление заказа по отсутствию в следующем слепке;
- маркер конца слепка, опись ожидаемых строк и транзакция между целями ODS.
## Поток данных
После завершения генератора Airflow один раз читает байтовый Kafka-чтец и
записывает полученную порцию в `stg.orders_raw`. Все строки получают `_load_id`,
равный `run_id` Airflow. `_load_ts` вычисляется при этой записи и дальше
переносится без пересчёта.
Один следующий `task_id` отвечает за весь переход STG → ODS. Внутри него два
последовательных `INSERT SELECT` читают неизменный срез по `_load_id`: первый
пишет годные строки в `ods.order_snapshot`, второй — брак в
`ods.order_snapshot_errors`. Транзакции между запросами нет. При частичном сбое
Airflow повторяет весь `task_id`; одинаковые исходные строки и служебные метки
не вычисляются заново.
Условия запросов взаимодополняющие: один общий предикат определяет брак, а
годная ветвь использует его буквальное отрицание. Все функции предиката
возвращают результат без исключения, а сам предикат всегда заканчивается в
`true` или `false`, не в `NULL`. Постоянный классификатор между STG и ODS для
этого не нужен.
## Граница строгого приёма
Единица решения — одна строка `stg.orders_raw`. Корень должен быть
JSON-объектом с точным набором ключей: `order_id`, `user_id`, `status`,
`created_at`, `updated_at`, `items_total`, `discount`, `delivery`, `total`,
`items`, `snapshot_date`.
Скалярные поля проверяются по типу JSON. Деньги дополнительно обязаны быть
строками с ровно двумя знаками после точки, времена — строками RFC 3339 в UTC с
обязательными миллисекундами, дата слепка — строкой `YYYY-MM-DD`. `items`
проверяется только как JSON-массив. Каноническая форма и основания выбора
зафиксированы в
[исследовании формата](../../research/2026-08-16-order-snapshot-wire-format.md).
Проверять все верхнеуровневые поля здесь уместно: их одиннадцать, и десять
скалярных значений непосредственно образуют типизированную строку заказа. У
события из 47 полей проверяются только пять опорных; переносить то сокращение на
малый контракт заказа нет причины. Граница строгости заканчивается на форме
провода: содержимое позиций и бизнес-инварианты намеренно остаются ниже.
Класс брака выбирается первым совпадением:
1. `not_an_object`;
2. `keyset_mismatch`;
3. `field_invalid`.
Имя отдельного поля в класс не включается. В таблице ошибок остаются сырой
текст, метаданные доставки и `_load_id`, поэтому единичный случай можно разобрать
без постоянной детализации предиката.
## Роль ODS
`ods.order_snapshot_rep` хранит физически принятые версии в
`ReplacingMergeTree(updated_at)`. Ключ сортировки — `order_id`, партиция — день
неизменного `created_at`. `ods.order_snapshot_dist` шардирует по
`cityHash64(order_id)`: только так все версии заказа попадают на один шард и
`FINAL` даёт корректный результат через распределённую таблицу.
Четыре координаты отвечают на разные вопросы:
- `updated_at` — какая бизнес-версия заказа новее;
- `snapshot_date` — в слепке какого модельного дня источник показал строку;
- `_load_id` — какой запуск Airflow принял строку;
- `_load_ts` — когда строка приехала в хранилище.
Обычное чтение `_dist` показывает физически сохранившиеся версии и нужно для
диагностики. Их число зависит от фоновых слияний: ODS не служит архивом истории.
`ods.order_v` сохраняет те же источник-ориентированные поля без обогащения
данными модели, но возвращает одну актуальную версию на `order_id`. Сначала оно
может быть простым представлением над `_dist FINAL`; способ выбора можно
заменить, не меняя потребителей.
Это представление остаётся ответственностью ODS: оно скрывает механику чтения
версий, но не строит бизнес-модель. DDS читает `ods.order_v` и отдельно решает,
какие сущности, связи и производные признаки ему нужны. Нужен ли `_load_id`
выше ODS, решается вместе с DDS, а не здесь.
## Одно чтение Kafka
Стандартный слепок содержит около полутора тысяч строк, тогда как предел одной
порции на стенде — десятки тысяч сообщений. Поэтому один запуск Airflow делает
один прямой `SELECT`, без цикла до пустоты и без фиксации конечных офсетов.
Это допущение о размере стенда, а не доказательство полноты слепка. Если Kafka
вернёт короткую порцию, непрочитанный хвост останется в топике и приедет в один
из следующих запусков. После отказа от замены партиции это задержка, а не потеря
или публикация неполного дня.
## Отклонённые варианты
- Партиционная идемпотентность — `REPLACE PARTITION snapshot_date` или
`ReplacingMergeTree` по `(snapshot_date, order_id)`: у потребителя нет
признака полноты партиции, а дата наблюдения становится частью ключа
сущности.
- Обычный `MergeTree` в ODS с дедупликацией только в DDS: навсегда сохраняет
технические повторы там, где семантика версии уже известна.
- Маркер, опись, чтение до пустоты или конечные офсеты: добавляют протокол ради
объёма, который с большим запасом помещается в одну порцию.
- Два `task_id` или материализованный классификатор: дробят один короткий
переход слоя, не добавляя транзакционности.
- Представление текущего состояния в DDS и готовая схема `dds.order` в этой
задаче: перекладывают механику ODS на следующий слой и преждевременно задают
модель данных.
## Риски и проверка
- На стандартном мире сверить число отправленных заказов с числом строк,
принятых одним прямым чтением. Это разовая приёмка допущения, не постоянный
сторож.
- На малой управляемой порции дать по одной строке каждого класса брака и две
годные версии одного `order_id`. Две цели должны сохранить все непустые
сообщения, а `ods.order_v` — вернуть новую версию независимо от фонового
слияния.
- Повторить переход с тем же `_load_id`: строка в `ods.order_v` и её `_load_ts`
не должны измениться; версии одного заказа должны остаться на одном шарде и
в одной партиции. Таблица ошибок может снова записать тот же брак: совпавшие
`_load_id` и Kafka-координаты показывают повтор задачи.
## Что проверено
MCP Context7 в сессии проектирования был недоступен. На локальном ClickHouse
`26.3.17.56` проверено, что прямой `SELECT` Kafka Engine завершается после
одной порции, а `FINAL` через `Distributed` исполняется на таблицах шардов.
Поэтому версии одного `order_id` направляются на один шард. Фоновое схлопывание
`ReplacingMergeTree` и необходимость точного чтения сверены с
[официальной документацией](https://clickhouse.com/docs/reference/engines/table-engines/mergetree-family/replacingmergetree),
поведение чтения — с исходниками той же версии
[`StorageKafka.cpp`](https://github.com/ClickHouse/ClickHouse/blob/v26.3.17.56-lts/src/Storages/Kafka/StorageKafka.cpp) и
[`KafkaSource.cpp`](https://github.com/ClickHouse/ClickHouse/blob/v26.3.17.56-lts/src/Storages/Kafka/KafkaSource.cpp).
## Связанные решения
- [ADR 0008](../../adr/0008-order-ingestion.md) сохраняет выбор пакетного
забора, `RawBLOB`, одного чтеца и одной партиции топика.
- [ADR 0010](../../adr/0010-order-versions-in-ods.md) заменяет публикацию
слепка версионным ODS.
- Вопрос `_load_id` выше ODS оставлен проектированию DDS в тикете #85.
+30
View File
@@ -0,0 +1,30 @@
# Что хранит опись
Резолюции развилок [«Расхождения A–D и опоздания: механика, доли и что
обещано»](https://git.dementev.space/ddmitry/clickstream-data-platform/issues/72)
и [«Места заказов в стартовом
мире»](https://git.dementev.space/ddmitry/clickstream-data-platform/issues/74).
Словарь описи — растяжка-хеш против опоры-счётчика — задан «Чем меряется
генератор» в [спеке генератора](../../specs/2026-08-01-generator.md) и
термином «Опись мира» в [CONTEXT.md](../../../CONTEXT.md).
> Опись хранит только то, чего движение мира не меняет.
> Хеш кладём всегда, счётчик — только когда назван его читатель.
- **У каждого отправленного слепка — своя строка с хешем байтов.** Байты
слепка не покрыты хешами дней ни при каком раскладе подпотоков — это второй
артефакт мира. Побайтовое обещание («слепок переснимается и даёт те же
байты») опись начинает сторожить.
- **Счётчики классов расхождений — наблюдаемых, после приоритета**, по
итоговой судьбе заказов дня; единица счёта — заказ, и опись называет её
словом. Сойтись с запросом менти они могут только на днях с закрытым окном:
при N сыгранных днях таких N − 7 (день d закрывается слепком d + 6, а
последний отправленный слепок несёт день N − 2). Читатель у них придёт
этапом 4 — проверка сверки; не окажется читателя — та же бритва режет и их.
- **Контрольные числа идентичности не заводятся вовсе**: uniq известных
пользователей и число двухкуковых пар растут, пока мир едет, — опоры из них
не выходит; генератор сторожит хеш, транспорт — счёт событий.
- **Опись описывает мир, а не доставку.** Работа генератора кончается на
Kafka: дошли ли байты до `ods.order_snapshot` — вопрос стенда и его
проверок. Поэтому счёта строк у слепка в описи нет; понадобится проверка
приёма заказов — число заведётся вместе с ней.
+96
View File
@@ -0,0 +1,96 @@
# Заказ и его слепок
Резолюция развилки [«Генератор слепков: где живёт и чем связан с
событиями»](https://git.dementev.space/ddmitry/clickstream-data-platform/issues/71).
## Откуда берётся заказ
Заказ и событие `purchase` — не два порождения, а две проекции одного факта
мира. Корзина, цены, купон и номер заказа посчитаны торговой половиной
дня-функции; заказная половина берёт заказы дня готовой структурой — вторым
выходом `commerce`, — навешивает на них жизнь заказа и собирает слепок. В
подпоток `COMMERCE` не добавляется ни одного нового броска: мир не сдвигается,
согласованность двух источников не удерживается, а получается по построению.
Следствия, которые уже решены соседями и здесь только связываются:
- у всякого заказа изначально ровно одно событие `purchase`; заказ без события
в трекере — не отдельная порода, а класс B, и делает его событийная сторона
выбрасыванием события после присвоения номера
([классы расхождений](fate.md));
- номер заказа общий у обеих проекций: `order_id` = клиентский `purchaseID`,
читаемый номер «день и порядковый номер покупки» ([спека
генератора](../../specs/2026-08-01-generator.md), раздел 9); нумеруются все
покупки, дошедшие до потока дня, — до всяких потерь;
- скидка заказа выводится из промокода события по таблице «код → скидка» —
числу мира, которое этап 3 берёт готовым (спека генератора, разделы 8 и 9).
В модели строка заказа в базе источника создаётся синхронно с покупкой,
поэтому день рождения заказа и день создания строки совпадают.
Отклонено: *выводить заказ разбором собственного вывода* (`purchaseID`, сырой
`ecommerce`) — бэкенд стал бы читателем трекера ровно там, где стенд учит, что
это разные источники; *независимая модель бэкенда* (заказ первичен, событие —
эхо) — кто купил, решает воронка, а воронка — это трафик, то есть опрокидывание
всего генератора; *слепок собирает SQL стенда из событий* — второй источник
исчезает вместе с уроком «две версии правды».
## Слепок и его доставка
**Запуск.** Третья команда того же пакета — `snapshot --day D [--days N]`,
свой приёмник, топик `orders`. Два источника — два запуска: трекер и бэкенд
видны глазами как два производителя, каждый со своим топиком. Один прогон с
двумя выходами отклонён: экономии он не даёт (со сдвигом отправки окно слепка
и сыгранный день не пересекаются вовсе), а правило «приёмник выбирается тем,
что для него назвали» ломает. Отклонены также: *отдельный пакет и образ*
библиотека на двоих ради одной команды; *генератор пишет слепок файлом, в
топик льёт даг* — второй путь доставки и второй сериализатор; *слепок едет
топиком `hits`* — убивает два режима приёма.
**Сборка окна.** Слепок дня D несёт заказы, рождённые в дни D−6…D, и
собирается переигровкой этих семи дней: заказы дня — производная всей воронки
дня, дешёвого пути к ним нет. Цена — семь проигрышей дня (~14 с) на слепок; у
начала оси окно усекается само. Отклонено: *кэш заказов на томе* — состояние
между прогонами; *окно держит хранилище* — топик перестаёт нести слепок;
*K = 1* — это страховочный срез 1 мастер-спеки, он в резерве.
**Отправка.** Слепок **снимается на границе суток, а отправляется следующим
прогоном**: даг, играющий день D, отправляет слепок дня D−1 — ночная выгрузка
бэкенда за вчера, как в бою. Содержимое слепка — чистая функция (зерно, D), от
момента отправки не зависит. Следствия:
- живой день перестаёт быть особым случаем: своего дага у него нет, слепок
живого дня отправит следующий прогон;
- пропущенный день лечится окном: слепок переснимается и даёт те же байты,
отдельного механизма самовосстановления нет;
- покупки текущего дня в сверке всегда `awaiting_order` — сюжет «вчера не
сходилось, сегодня сошлось», ради которого мастер-спека этот класс завела;
- цена — один лишний проигрыш дня на прогон (окно и сыгранный день не
пересекаются, проигрышей всегда восемь).
На старте оси дня −1 нет, поэтому прогон дня 0 не отправляет ничего; первый
слепок — дня 0 — уезжает прогоном дня 1 ([исследование
формата](../../research/2026-08-16-order-snapshot-wire-format.md)).
**Случайность.** Судьбу заказов бросает свой подпоток, ветвящийся по дню
рождения заказа: слепок несёт семь дней рождения сразу, и судьбу каждого
заказа обязан читать из его собственного дня. Вся судьба решается при
рождении, поэтому слепок любого дня — чтение готовой судьбы, а не накопление
состояния. Отклонено: *дописывать броски в конец `COMMERCE`* — правка заказа
и правка торгового поведения стали бы одним рычагом; *бросать состояние в
подпотоке дня слепка* — траектория заказа зависела бы от того, какие слепки
снимали.
## Запись на проводе
Контракт провода — на заказ один JSON-документ: деньги строками с двумя
знаками, времена RFC 3339 в UTC с миллисекундами, `items` обычным массивом —
целиком описан мастер-спекой (раздел 2); основания, отклонённые варианты и
проверка разбора — в [исследовании
формата](../../research/2026-08-16-order-snapshot-wire-format.md) (резолюция
развилки
[«Форма записи слепка на проводе»](https://git.dementev.space/ddmitry/clickstream-data-platform/issues/81)).
Сверх контракта здесь живёт одно правило: **порядок строк внутри слепка —
порядок рождения заказов, он же возрастание `order_id`**. Детерминизм даёт
его даром, а хешу слепка в описи нужен именно названный порядок.
+34
View File
@@ -0,0 +1,34 @@
# Заказы в стартовом мире
Резолюция развилки [«Места заказов в стартовом
мире»](https://git.dementev.space/ddmitry/clickstream-data-platform/issues/74).
Стартовый мир — не полный мир, а мир, остановленный на границе суток 7|8.
Генератор играет два источника с разными темпами — поток Метрики и ночную
выгрузку магазина; разные темпы дают всё остальное.
`world-init` играет восемь дней одним прогоном `batch --day 0 --days 8`;
слепки отправляет второй запуск — команда `snapshot` того же диапазона (два
источника — два запуска, [слепок и его доставка](snapshot.md)), и со сдвигом
отправки уезжают **слепки дней 0…6**. Слепок дня 7 снят на границе суток и
уедет первым же ходом мира. Особого режима у стартового мира нет: правило
отправки живёт в одном месте — в проигрывателе; генератору это решение не
стоит ничего.
Что видно снаружи: у последнего прожитого дня клики есть, а заказов нет;
глубже — день 0 виден дожившим до конца окна, день 6 — только что родившимся.
**График выручки заваливается к правому краю** и дозаполняется, пока мир едет.
Это не издержка стенда, а главный наблюдаемый эффект второго источника: ночная
выгрузка отстаёт, свежие дни предварительны. На свежем стенде лаба сверки
видит целый день `awaiting_order` — норма, а не поломка; как это назвать
менти — за витринами этапа 4.
Цена принята с открытыми глазами: у части пар стартового мира поздний
назначенный заказ падает на день 7, и до первого хода мира этих пар в
`dds.identity_map` нет. Опись пар не считает
([что хранит опись](inventory.md)), поэтому красной проверки из этого не
выходит. Отклонено: *дослать восьмой слепок* — исчезает `awaiting_order` на
свежем стенде, в CLI заводится рычаг, стартовый мир становится особым случаем
ровно там, где #71 его убирал; *счётчик пар учится спрашивать про слепки*
число верно ровно до первого хода мира; *отменить сдвиг отправки*
`awaiting_order` пропал бы навсегда; *подогнать план под горизонт* — мир
перестал бы быть чистой функцией зерна.
+58 -23
View File
@@ -78,15 +78,16 @@ keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/`
Ключи ко-локации названы заранее, потому что на них стоит политика соединений из
раздела 6 спеки: обычное соединение разрешено только по ключу ко-локации, всё
прочее — через `GLOBAL`. Значит `dds.session` и `dds.identity_map` шардируются по
`cityHash64(ClientID)`, а `dds.order` и производные от заказа — по
`cityHash64(order_id)`. Ключи витрин появятся вместе с самими витринами.
`cityHash64(ClientID)`. Объекты DDS с зерном заказа должны сохранять ко-локацию
по `cityHash64(order_id)`; ключи остальных частей будущей модели и витрин
появятся вместе с ними.
Открытый вопрос на будущее — не сама замена партиций: операции с ними по
локальным таблицам правило разрешает прямо. Вопрос в шаге до неё. Партиция-донор
должна быть уже разложена по шардам по тому же ключу, а разложить её можно
только вставкой через распределённую таблицу — значит у каждой пакетной сущности
появится вторая пара объектов, и имени для неё конвенция пока не даёт. Решать
это вместе со сборкой DDS, а не задним числом.
Замена партиций не используется для `ods.order_snapshot`: версии прошлых дней
доливаются, а прямое чтение Kafka не задаёт границы полного слепка ([ADR
0010](../adr/0010-order-versions-in-ods.md)). Если партиционная пересборка
понадобится будущим объектам DDS или DM, партиция-донор должна быть заранее
разложена по шардам по тому же ключу. Форму донора следует решать вместе с таким
объектом, а не переносить на ODS заранее.
## Служебные колонки
@@ -116,8 +117,10 @@ keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/`
колонка молча отвечала бы на другой вопрос.
Само сообщение лежит в колонке `raw` тем, чем пришло: чтец читает байты и ничего
не проверяет, поэтому там оказываются и целые события, и мусор. Разбирается всё
это ниже, в матвью ODS — см. [ADR 0005](../adr/0005-event-ingestion.md).
не проверяет, поэтому там оказываются и целые сообщения, и мусор. События ниже
разбирают матвью ODS ([ADR 0005](../adr/0005-event-ingestion.md)), заказы —
пакетный шаг ([ADR 0008](../adr/0008-order-ingestion.md), [ADR
0010](../adr/0010-order-versions-in-ods.md)).
Движок таблицы сырья — обычный `ReplicatedMergeTree`, `ORDER BY (kafka_partition,
kafka_offset)`: разбор полётов идёт от «какое сообщение», другого ключа у сырья и
@@ -125,15 +128,23 @@ kafka_offset)`: разбор полётов идёт от «какое сооб
которого он заведён: повтор доставки в сырье обязан быть виден.
Метка времени загрузки зовётся `_load_ts`, тип `DateTime64(3, 'UTC')`. Ставится
она один раз, в матвью приёма, и дальше переносится из STG в ODS как есть:
колонка отвечает на вопрос «когда строка приехала в хранилище», а не «когда её
разобрали». В ODS она же служит колонкой версии `ReplacingMergeTree`, и работа у
этой версии ровно одна — схлопнуть повтор доставки. Содержимое у повтора то же
самое, отличается только метка, поэтому какая из двух строк переживёт мерж,
безразлично. Пакетной переобработки у ODS нет: слой наполняет матвью, а не
задание Airflow, и работа с партициями начинается выше. Переделать разобранное
руками можно — вставкой из сырья с фильтром по `_load_ts`, в пределах
трёхсуточного окна; ничья по версии разрешается в пользу вставленного позже.
она один раз при записи в STG: для событий — матвью приёма, для дневного слепка
заказов — пакетным шагом. Дальше метка переносится в ODS как есть и отвечает на
вопрос «когда строка приехала в хранилище», а не «когда её разобрали». В
`ods.event` она же служит колонкой версии `ReplacingMergeTree` и схлопывает
повтор доставки. У `ods.order_snapshot` версию задаёт `updated_at` источника;
`_load_ts` только показывает, когда конкретная строка приехала.
`created_at` и `updated_at` заказа к служебным колонкам хранилища не относятся.
Они приезжают в сообщении как аудит строки в БД источника и в ODS разбираются в
`DateTime64(3, 'UTC')`; ClickHouse их не создаёт и добавляет рядом собственную
`_load_ts`. Совпадение слов «техническое время» не делает эти часы одной осью.
Пакетной переобработки у событий в ODS нет: слой наполняют матвью, а не задание
Airflow. Переделать разобранное руками можно вставкой из сырья с фильтром по
`_load_ts` в пределах трёхсуточного окна; ничья по версии разрешается в пользу
вставленного позже. Заказы разбирает пакетный шаг из неизменного среза STG по
`_load_id`; дневные партиции ODS он не заменяет.
Имя согласовано с каноном служебных полей соседнего учебного стенда на
Greenplum, чтобы словарь был общим у двух хранилищ; ведущее подчёркивание у
@@ -141,10 +152,28 @@ Greenplum, чтобы словарь был общим у двух хранил
её же используют Fivetran, Airbyte и Stitch. С правилом выше это не спорит:
запрещено совпадать с именами виртуальных колонок, а не носить подчёркивание.
Идентификатора пачки загрузки (`_load_id`) пока нет. В STG и ODS данные приезжают
потоком через матвью, у которого нет ни батча, ни `run_id`, и колонка была бы
пустой формальностью. В слоях, которые наполняет Airflow, `run_id` появится
по-настоящему — тогда и заведём, тем же стилем имени.
Общего идентификатора пачки загрузки (`_load_id`) нет. У потока событий нет ни
пачки, ни `run_id`, и колонка была бы пустой формальностью. У заказов читатель
назван: `_load_id` равен `run_id` Airflow и переносится из STG в годную строку
ODS и в таблицу ошибок. Нужен ли он выше ODS, решается вместе с моделью DDS.
## Версии заказов
`ods.order_snapshot_rep` хранит принятые версии в
`ReplacingMergeTree(updated_at)`: ключ сортировки — `order_id`, партиция —
`toDate(created_at)`. `ods.order_snapshot_dist` шардирует по
`cityHash64(order_id)`. Все версии заказа лежат в одной партиции, чтобы их могли
схлопывать фоновые слияния, и на одном шарде, чтобы распределённый `FINAL`
выбрал одного победителя.
Физическая пара нужна для загрузки и диагностики. Обычное чтение показывает
версии, которые ещё не убрали фоновые слияния, и не является архивом истории.
`ods.order_v` служит поверхностью точного текущего состояния для следующих
слоёв. Представление сохраняет язык источника и не решает, какой станет модель
DDS. `snapshot_date` в нём остаётся датой наблюдения строки, а не ключом
публикации. Полное решение — в
[ADR 0010](../adr/0010-order-versions-in-ods.md) и
[спецификации приёма заказов](orders/ingestion.md).
## Часовые пояса
@@ -358,6 +387,12 @@ D0 и к реальному календарю не привязана; паке
kafka_offset)`: смотрят такую таблицу от класса, а внутри класса — по координатам
доставки.
`ods.order_snapshot_errors` держит тот же диагностический минимум и `_load_id`
запуска. У заказов три класса по приоритету: `not_an_object`,
`keyset_mismatch`, `field_invalid`. Сырой текст остаётся рядом, поэтому класс не
разрастается до имени отдельного поля. Точная граница приёма — в
[спецификации заказов](orders/ingestion.md).
## Раскладка DDL
Файлы лежат в `sql/ddl/` и применяются по порядку имён. Сначала все статичные
@@ -0,0 +1,212 @@
# Формат дневного слепка заказов на проводе
Дата исследования: 2026-08-16.
Учебный результат: менти различает бизнес-время, время источника, доставки и
загрузки, не разбирая ради этого лишнюю инфраструктуру. Цена — несколько явных
правил контракта; новых полей и универсального сериализатора не требуется.
## Короткий вывод
- `created_at` и `updated_at` — аудит строки в источнике, а не время покупки.
Оба поля передаются в UTC с настоящей точностью до миллисекунд:
`2026-06-03T14:21:07.123Z`.
- В ClickHouse им соответствует `DateTime64(3, 'UTC')`. Неверная строка даёт
`NULL` и уходит в `*_errors`, а не превращается в правдоподобную дату.
- `snapshot_date` — дата завершившегося модельного дня, состояние которого
снято на исходящей границе суток. Внутри одной выгрузки она одинакова, на
следующем модельном дне меняется.
- `items` на проводе — обычный массив JSON. Тип `String` в ODS означает, что
из внешнего JSON извлекли сырой фрагмент массива, а не что источник дважды
сериализовал JSON.
- Генератору достаточно собрать один словарь с вложенным списком и один раз
вызвать `orjson.dumps`. Отдельная иерархия кодеков урока не добавляет.
## Оси времени
В потоковой обработке время события принадлежит самой записи и не зависит от
часов обработчика; время обработки отвечает на другой вопрос
([Apache Flink: Event Time и Processing Time](https://nightlies.apache.org/flink/flink-docs-stable/docs/concepts/time/)).
Debezium проводит ту же границу внутри одного сообщения: время изменения в
исходной БД хранится отдельно от времени обработки коннектором, а их разность
можно использовать как задержку
([документация коннектора PostgreSQL](https://debezium.io/documentation/reference/stable/connectors/postgresql.html#postgresql-create-events)).
| Поле | Чьи часы | На какой вопрос отвечает | Форма |
|---|---|---|---|
| `UTCEventTime` события `purchase` | бизнес-событие, трекер | когда покупатель подтвердил покупку | отдельный контракт кликстрима; связь с заказом по `purchaseID = order_id` |
| `created_at` | база источника | когда строка заказа впервые создана в источнике | RFC 3339 UTC с тремя знаками долей секунды |
| `updated_at` | база источника | когда эта строка в последний раз изменена в источнике | тот же формат; версия состояния заказа |
| `snapshot_date` | модельный календарь | состояние какого завершившегося дня снято на границе суток | `YYYY-MM-DD`, без времени |
| `kafka_timestamp` | транспорт | когда брокер пометил доставленное сообщение | служебная колонка хранилища |
| `_load_ts` | хранилище | когда строка впервые приехала в хранилище | `DateTime64(3, 'UTC')` |
`updated_at` как метка последнего изменения исходной строки совпадает с
рекомендованным смыслом `updated_at` в timestamp-стратегии dbt snapshots;
время выполнения самого слепка dbt хранит отдельно
([официальная документация dbt](https://docs.getdbt.com/docs/build/snapshots#timestamp-strategy-recommended)).
Служебные метки ETL также являются отдельными метаданными процесса, а не
бизнес-фактами
([Kimball Group: Audit Dimension](https://www.kimballgroup.com/data-warehouse-business-intelligence-resources/kimball-techniques/dimensional-modeling-techniques/audit-dimension/)).
В этом проекте транспортная и складская оси уже разведены в
[конвенции хранилища](../architecture/storage.md): `_load_ts` ставится один раз,
а миллисекундный `kafka_timestamp` не округляется.
`created_at` и `updated_at` приезжают в сообщении источника; ClickHouse их не
создаёт и добавляет рядом собственную `_load_ts`.
Следствие для контракта: `created_at` нельзя называть временем покупки, а
`updated_at - created_at` — длительностью бизнес-перехода. Это время между
созданием и последним изменением строки в источнике. Бизнес-время живёт в
событии `purchase` и связывается с заказом по уже существующему ключу.
### Партиция и бизнес-время
`created_at` и `updated_at` — аудит строки источника. Поэтому
`toDate(created_at)` в
[спецификации приёма заказов](../architecture/orders/ingestion.md)
используется как стабильный технический ключ партиции `ods.order_snapshot`:
это день создания строки
источника, а не доказательство дня бизнес-события.
Минимальное решение — не добавлять `ordered_at` на всякий случай. Пока модель
создаёт исходную строку синхронно с покупкой, существующий ключ партиции можно
оставить, но в описании называть его днём создания строки. Время покупки для
сверки берётся из `purchase.UTCEventTime`. Отдельное поле в заказе понадобится
только тогда, когда появится самостоятельный учебный запрос к бизнес-времени
заказа или источник начнёт сохранять заказ асинхронно. Так различие остаётся
честным, но не порождает поле без потребителя.
## Точность и строгий разбор
RFC 3339 — профиль ISO 8601 для обмена датой и временем. Он разрешает дробную
часть секунды переменной длины и как `Z`, так и числовое смещение
([RFC 3339, §5.6](https://www.rfc-editor.org/rfc/rfc3339#section-5.6)). Значит
миллисекунды не следуют из названия стандарта сами по себе. Наш более узкий
контракт фиксирует ровно три цифры и UTC:
```text
YYYY-MM-DDTHH:mm:ss.SSSZ
```
Например: `2026-06-03T14:21:07.123Z`. Одинаковое число цифр дробной части и
одинаковая зона дают хронологическую сортировку таких строк в лексикографическом
порядке
([RFC 3339, §5.1](https://www.rfc-editor.org/rfc/rfc3339#section-5.1)). Три
цифры выбраны потому, что источник моделирует миллисекунды. Сериализатор всегда
выводит все три цифры, в том числе `.000` для значения точно на границе секунды.
Нельзя только выдавать секундную модель за миллисекундную простым дополнением
нулей.
В ClickHouse `DateTime64(3, 'UTC')` хранит три десятичных знака долей секунды,
то есть миллисекунды; пояс колонки используется при разборе и показе значения
([DateTime64](https://clickhouse.com/docs/sql-reference/data-types/datetime64)).
Для этого узкого формата подходит обнуляемый разбор по точному шаблону:
```sql
parseDateTime64InJodaSyntaxOrNull(
value,
'yyyy-MM-dd\'T\'HH:mm:ss.SSS\'Z\'',
'UTC'
)
```
`OrNull` возвращает `NULL` при несовпадении, а три `S` задают точность
`DateTime64(3)`
([документация функции](https://clickhouse.com/docs/sql-reference/functions/type-conversion-functions#parsedatetime64injodasyntaxornull),
[исходный код ClickHouse](https://github.com/ClickHouse/ClickHouse/blob/master/src/Functions/parseDateTime.cpp)).
Это строже, чем `parseDateTime64BestEffortOrNull`: функция Best Effort по
назначению принимает несколько представлений даты, тогда как здесь форма сама
является частью учебного контракта
([документация Best Effort](https://clickhouse.com/docs/sql-reference/functions/type-conversion-functions#parsedatetime64besteffortornull)).
На проектном ClickHouse 26.3.17.56 это выражение локально проверено. Оно
возвращает `Nullable(DateTime64(3, 'UTC'))` для строки с `.123Z` и `NULL` для
строки без миллисекунд, с четырьмя цифрами, со смещением `+00:00` вместо `Z`
или с хвостовым мусором. Поэтому один и тот же результат разбора можно
использовать и для типизированной строки, и для маршрутизации ошибки; нулевая
дата не нужна.
`updated_at` допустим как колонка версии `ReplacingMergeTree`: ClickHouse
явно разрешает для `ver` тип `DateTime64` и оставляет строку с максимальной
версией
([ReplacingMergeTree](https://clickhouse.com/docs/engines/table-engines/mergetree-family/replacingmergetree)).
Если две версии одного заказа имеют одинаковый `updated_at`, среди них
побеждает вставленная позже. Для учебной модели достаточно гарантировать
монотонные миллисекундные `updated_at` на один заказ; отдельный счётчик версий
без такого сценария был бы лишним.
## Дата слепка и константы мира
Периодический слепок имеет зерно заранее заданного периода — например, дня, —
а не отдельной транзакции
([Kimball Group: Periodic Snapshot Fact Tables](https://www.kimballgroup.com/data-warehouse-business-intelligence-resources/kimball-techniques/dimensional-modeling-techniques/periodic-snapshot-fact-table/)).
Поэтому `snapshot_date` — значение пачки: дата завершившегося модельного дня D,
состояние которого снято на границе D|D+1. Это не UTC-дата отправки и не
глобальная константа.
На один проход генератор вычисляет дату один раз и кладёт её во все записи; на
следующем модельном дне значение меняется. Настоящие константы мира — начало
модельной оси `ORIGIN`, пояс счётчика `COUNTER_TIMEZONE_MINUTES` и окно K = 7.
Первые две уже заданы в
[модели времени](../../generator/src/clickstream_generator/world.py), окно
добавится в конфигурацию мира вместе с заказами. Глобальная `SNAPSHOT_DATE`
смешала бы правило календаря с результатом его вычисления.
На старте оси дня −1 нет, поэтому прогон дня 0 ничего не отправляет. Первый
слепок с `snapshot_date = ORIGIN` уезжает прогоном дня 1.
`as_of_date` тоже могло бы означать дату, по состоянию на которую показаны
данные. Но `snapshot_date` уже является языком спеки и ADR о приёме заказов.
Переименование не добавляет урока и может спутать дату
выгрузки с периодом бизнес-действия записи. Для этого стенда оставляем
`snapshot_date`.
## JSON и сериализация
В JSON массив и строка — разные типы значения: массив содержит значения
непосредственно, а строка содержит последовательность символов
([RFC 8259, §§3, 5 и 7](https://www.rfc-editor.org/rfc/rfc8259)). Поэтому форма
на проводе такая:
```json
{
"order_id": "20260603-0001",
"user_id": 42,
"status": "paid",
"created_at": "2026-06-03T14:21:07.123Z",
"updated_at": "2026-06-03T14:24:18.456Z",
"items_total": "1299.90",
"discount": "0.00",
"delivery": "199.00",
"total": "1498.90",
"items": [
{"sku": "sku-17", "qty": 1, "price": "1299.90"}
],
"snapshot_date": "2026-06-07"
}
```
ClickHouse `JSONExtractRaw(raw, 'items')` возвращает выбранный фрагмент JSON
неразобранной строкой
([официальная документация](https://clickhouse.com/docs/sql-reference/functions/json-functions#jsonextractraw)).
На проектной версии локальная проверка обычного внешнего JSON показала
`JSONType(..., 'items') = 'Array'`, а `JSONExtractRaw` вернул компактный текст
массива, пригодный для колонки ODS `String`. Строка с JSON внутри потребовала
бы экранировать массив при первой сериализации и разбирать его второй раз, не
меняя результат в ODS.
`orjson.dumps` умеет сериализовать вложенные словари и списки напрямую и
возвращает JSON в UTF-8
([официальный репозиторий orjson](https://github.com/ijl/orjson)). Поэтому
KISS-вариант для [существующего модуля сериализации](../../generator/src/clickstream_generator/serialize.py)
— подготовить канонические строки времени и денег, положить `items` списком в
общий словарь и сделать один внешний `dumps` на запись. Класс кодеков, реестр
схем и повторный `dumps` для `items` здесь ничего не учат.
Граница этого решения: `toDecimal64OrNull(..., 2)` проверяет числовую
преобразуемость, но не лексическое правило «ровно два знака» — локально строки
`1299.90`, `1299.9` и `1299.900` дали одно значение. Проверка денежного формата
не нужна сериализатору этого слепка: он сам выпускает ровно два знака. Приём ODS
проверяет ту же каноническую форму и считает остальные формы браком; граница
строгого приёма зафиксирована в
[спецификации заказов](../architecture/orders/ingestion.md).
+90 -61
View File
@@ -190,36 +190,52 @@ Ecommerce (заполнены только у торговых событий):
выручка дня D «дышит» K дней, потом замерзает. Боевой аналог окна есть и у
трекеров: лог Метрики «доформировывается» ещё около трёх дней.
- Запись слепка — состояние заказа на момент выгрузки, «родной» экспорт
бэкенда в snake_case:
бэкенда в snake_case. Таблица задаёт тип после разбора в ODS:
| Поле | Тип | Комментарий |
| Поле | Тип в ODS | Комментарий |
|---|---|---|
| `order_id` | String | номер заказа; равен клиентскому `purchaseID` |
| `user_id` | UInt64 | пользователь магазина — мост к склейке |
| `status` | String | `created``paid``cancelled` |
| `created_at`, `updated_at` | DateTime | |
| `status` | String | `created``paid``cancelled`; неоплаченный отменяется прямым переходом `created``cancelled` |
| `created_at`, `updated_at` | DateTime64(3, 'UTC') | аудит строки в БД источника: создание и последнее изменение; не бизнес-время покупки и не время загрузки в ClickHouse |
| `items_total`, `discount`, `delivery`, `total` | Decimal(18,2) | деньги бэкенда — в Decimal |
| `items` | String | позиции вложенным JSON: `[{sku, qty, price}]` |
| `snapshot_date` | Date | дата слепка (день выгрузки) |
| `items` | String | сырой текст массива позиций `[{sku, qty, price}]`, извлечённый из внешнего JSON |
| `snapshot_date` | Date | завершившийся модельный день, состояние которого снято на исходящей границе суток |
- Приём идемпотентный, но дедуп расщеплён на два слоя:
- `ods.order_snapshot` — партиция по `snapshot_date`, **без дедупа**,
хранит «как приехало»; идемпотентность повторного прогона — заменой
партиции дня слепка, а не ReplacingMergeTree.
- Дедуп до последней версии — **argMax** в трансформации при сборке
`dds.order`. `dds.order` — единственная дедуплицированная таблица:
партиция по дню заказа (`toDate(created_at)`),
ReplacingMergeTree(`updated_at`), `ORDER BY order_id` — заказ всегда
лежит в одной партиции, дедуп работает.
На проводе один заказ — один документ JSON. Деньги, включая `items[].price`,
передаются строками с ровно двумя знаками после точки; `created_at` и
`updated_at` — строками RFC 3339 в UTC с обязательными миллисекундами
(`YYYY-MM-DDTHH:mm:ss.SSSZ`); `snapshot_date` — строкой `YYYY-MM-DD`;
`items` — обычным массивом JSON, не строкой с JSON внутри. Порядок внешних
ключей совпадает с порядком полей в таблице контракта, у позиции — `sku`,
`qty`, `price`: так байты воспроизводимы без сортировки ключей. Разница с
кликстримом намеренна: там значения `purchaseRevenue` приезжают JSON-числами
и разбираются как `Array(Float64)`, а бэкенд передаёт деньги строками для
точного `Decimal`. Это показывает расхождение представлений денег в двух
источниках.
Запись собирает явная функция существующего канонического сериализатора:
один словарь с вложенным списком и один `orjson.dumps`, без универсального
слоя кодеков.
Обоснование и проверка разбора — в
[исследовании формата](../research/2026-08-16-order-snapshot-wire-format.md).
- `ods.order_snapshot` принимает версии заказа в
`ReplacingMergeTree(updated_at)`: `ORDER BY order_id`, партиция по дню
неизменного `created_at`, шардирование по `cityHash64(order_id)`.
`snapshot_date` остаётся датой наблюдения источника, но не задаёт публикацию
или идемпотентность. Физическое чтение может видеть несколько версий;
`ods.order_v` возвращает точное текущее состояние. Устройство модели DDS и
способ её материализации решаются отдельно.
Пропущенный день ничего не ломает, следующий слепок самовосстанавливает.
- Разбор JSON-позиций — **один раз**, в трансформации ODS → DDS; дальше
витрины работают с плоскими массивами `dds.order`: `item_sku`
Array(String), `item_qty` Array(UInt64), `item_price` Array(Decimal(18,2))
— одной длины, порядок как в JSON. Это единственный носитель навыка
«вложенный JSON в ClickHouse» на стенде.
- `items` остаётся сырой строкой JSON в ODS. Разбирать позиции следует на
границе ODS → DDS, но их представление определяется вместе с будущей моделью
заказов. Это остаётся носителем навыка «вложенный JSON в ClickHouse», не
превращая приём в преждевременную модель данных.
- Статусы держим все три: смена `created``paid` и есть причина «дыхания»
выручки внутри окна; сужение до двух — резервный срез 1.
выручки внутри окна; сужение до двух — резервный срез 1. На выходе из окна
заказ либо `paid`, либо `cancelled`: неоплаченного отменяют, «навсегда
`created`» не бывает ([судьба заказа](../architecture/orders/fate.md)).
## 3. Каталог товаров
@@ -242,9 +258,9 @@ CSV в репозитории (`data/catalog/products.csv`: `sku`, `name`, `cate
| | Расхождение | Механика в генераторе | Ориентир доли |
|---|---|---|---|
| A | Отмена | заказ дошёл до `cancelled`, `purchase` остался | ~5% заказов |
| A | Отмена | заказ дошёл до `cancelled`, `purchase` остался | ~5% заказов — отменённые вообще, обе дороги отмены вместе |
| B | Потерянное событие | заказ есть, `purchase` не доехал | ~3% |
| C | Дельта суммы | сверка приведена к сравнимой базе (`items_total`, не `total`); `amount_delta`только необъяснённый остаток после этого, и создаёт его генератор намеренно: деньги считаются целыми копейками, поэтому Float64 сам по себе не плывёт | ~12% |
| C | Дельта суммы | вычеркнутая позиция: товара не оказалось в наличии, у заказа на одну позицию меньше, чем в клиентских массивах. Сверка приведена к сравнимой базе (`items_total`, не `total`); `amount_delta`остаток, не объяснённый этим приведением, и сравнением позиций он объясняется — на этом стоит урок класса | ~12% |
| D | Дубль события | повторный `purchase` от обновления `/confirmation`: новый `WatchID` с тем же `purchaseID` — бизнес-дубль, не технический; дедуп ReplacingMergeTree его не съедает и не должен | ~2% |
Классы пересекаются — приоритет: `cancelled` > `lost_event` >
@@ -253,8 +269,8 @@ CSV в репозитории (`data/catalog/products.csv`: `sku`, `name`, `cate
Пятое — **опоздание** — бесплатно даёт формат доставки: часть заказов
впервые появляется в слепке D+1/D+2 («вчера не сходилось, сегодня сошлось»),
ориентир ~10%. Точные доли фиксируются при пересборке эталонного мира;
опись хранит точные счётчики по каждому классу расхождений (отмены,
потери, дубли).
опись хранит счётчики наблюдаемых классов после приоритета — в заказах и
только по дням с закрытым окном (раздел 8).
Не берём: сироту-фрод (`purchase` есть, а заказа не будет никогда) —
механически дублирует B.
@@ -286,8 +302,9 @@ CSV в репозитории (`data/catalog/products.csv`: `sku`, `name`, `cate
`uniq(посетителей) > uniq(людей)`, менти выводит расхождение сам.
Константа мира: каждый двухкуковый покупатель делает минимум по одному
заказу с каждой куки — иначе вторая кука не попадает в карту соответствий
(она строится только из покупок) и лаба не воспроизводится. Опись
хранит число именно таких пар.
(она строится только из покупок) и лаба не воспроизводится. Число таких
пар растёт, пока мир едет, и в опись не кладётся
([что хранит опись](../architecture/orders/inventory.md)).
- Витрины разводят имена честно: **«посетители»** (`uniq(ClientID)`) и
**«известные пользователи»** (после склейки) — оба числа рядом в дашборде.
@@ -342,18 +359,18 @@ CSV в репозитории (`data/catalog/products.csv`: `sku`, `name`, `cate
— только явный GLOBAL; `NOT IN` — только `GLOBAL NOT IN`. Сверка
`purchase`↔заказ —
легитимная GLOBAL-витрина (заказы малы).
- **Конвейер без TRUNCATE**: поток — append-only в ReplacingMergeTree (дедуп
через argMax); батчевая переобработка — по дневным партициям
- **Конвейер без TRUNCATE**: поток версий — в ReplacingMergeTree, точное чтение
через `FINAL` или равносильный выбор последней версии; батчевая
переобработка нижележащих объектов — по дневным партициям
(`DROP/REPLACE PARTITION ON CLUSTER`); `TRUNCATE ... ON CLUSTER` в конвейере не
применяется вовсе — полный сброс стенда делается `make clean && make up`, то
есть вместе с томами. `DROP/REPLACE PARTITION` работает только по
**локальным** таблицам ON CLUSTER, не по Distributed; замена через
DROP+INSERT неатомарна — дашборд в середине прогона честно моргает (это
осознанная цена, не баг).
- **Поздние заказы поглощает только ODS** (`ods.order_snapshot` — новая
партиция дня слепка, без переделки старого); материализованное ниже —
нет. Каждый прогон ETL перестраивает партиции последних K+1 дней у
заказозависимых объектов (`dds.order` и производные, `dm.dq_summary`).
- **Поздние заказы поглощает ODS** как новую версию `order_id`. Как их
подхватывают материализованные объекты DDS и DM, решается вместе с их моделью,
а не при проектировании приёма.
Сессии перестраиваются только за текущий день: правило мира — сессия
режется по границе модельных суток, дневная партиция самодостаточна.
- Для ETL-вставок — `distributed_foreground_insert = 1` (раньше называлась
@@ -391,10 +408,10 @@ README.
| STG | `stg.hits_raw_kafka`, `stg.hits_raw` + MV | сырые строки событий, Kafka Engine на обеих нодах |
| STG | `stg.orders_raw_kafka`, `stg.orders_raw`, без MV | сырые строки слепка; чтец на ноде 1, забирает пакетный шаг |
| ODS | `ods.event` (+`_errors`) | типизированное широкое событие, ReplacingMergeTree |
| ODS | `ods.order_snapshot` (+`_errors`) | слепки заказов как приехали, партиция по `snapshot_date`, без дедупа |
| ODS | `ods.order_snapshot` (+`_errors`), `ods.order_v` | типизированные версии заказов, брак и точное текущее состояние |
| DDS | `dds.session` | сборка сессий из событий (наследник `dds.click`) |
| DDS | `dds.event_v` | представление над `ods.event`: snake_case-имена, расшифровка кодов `DeviceCategory`; витрины DM читают его, а не ODS напрямую |
| DDS | `dds.order` | единственная дедуплицированная таблица заказа: партиция по дню заказа (`toDate(created_at)`), ReplacingMergeTree(`updated_at`), `ORDER BY order_id`, дедуп до последней версии — argMax в трансформации при сборке |
| DDS | модель заказов | зерно, связи и материализация проектируются на этапе DDS |
| DDS | `dds.identity_map` | карта кука↔пользователь |
| DDS | словарь `products` | каталог из CSV |
| DM | витрины `dm.*_v`, `dm.dq_summary` | см. ниже |
@@ -404,19 +421,22 @@ README.
см. [доку хранилища](../architecture/storage.md).
Заказы принимаются **пакетным забором** ([ADR
0008](../adr/0008-order-ingestion.md)): чтец топика байтовый, как у событий, но
0008](../adr/0008-order-ingestion.md), [ADR
0010](../adr/0010-order-versions-in-ods.md)): чтец топика байтовый, как у событий, но
матвью к нему не привязана, и сырьё забирает шаг, которым управляет Airflow —
он же вставляет прочитанное в `stg.orders_raw`, разбирает в типизированный
слепок и заменяет партицию дня в `ods.order_snapshot`. Тот же даг проигрывает
модельный день генератором, поэтому переливается ровно то, что он положил в
топик. Слой сырья у заказов остаётся: без него `ods.order_snapshot` повторил бы
его роль, а обещание идемпотентности повисло бы — матвью партиций не заменяет.
одно прямое чтение вставляет порцию в `stg.orders_raw` с `_load_id` запуска.
Один следующий `task_id` двумя последовательными запросами пишет годные строки
и ошибки из того же среза. `ods.order_snapshot` хранит версии по `updated_at`,
а не публикует партицию `snapshot_date`. Слой сырья остаётся точкой повтора и
разбора одной принятой порции.
Стенд получает от этого два режима приёма рядом, поток и слепок, и сравнение
из опорных точек раздела 12 переформулировано под них.
Состав служебных колонок задаёт дока хранилища. Спеке важны два следствия:
`ods.event` и `ods.order_snapshot` получают метку загрузки `_load_ts`, и у
`ods.event` она же служит колонкой версии ReplacingMergeTree; а таблицы
`ods.event` и `ods.order_snapshot` получают метку загрузки `_load_ts`, но у
заказа бизнес-версию задаёт `updated_at`; `_load_id` проходит через STG и обе
цели ODS. Таблицы
`stg.*_raw` хранят метаданные доставки Kafka вместе с именем читавшей ноды —
без них урок «какая нода читала топик» ненаблюдаем.
Модельного дня в STG нет: `EventDate` — свойство содержимого, а содержимое
@@ -445,8 +465,9 @@ Kafka день переигрывается генератором заново:
`total`: промокод и доставка клиенту не видны), `status`, `mismatch_class`
(`match` / `cancelled` / `lost_event` / `duplicate_event` / `amount_delta`,
в порядке приоритета — классы пересекаются, побеждает более ранний).
`match` — большинство строк; `amount_delta`только необъяснённый остаток
после приведения к сравнимой базе; его создаёт генератор намеренно
`match` — большинство строк; `amount_delta`остаток, не объяснённый
приведением к сравнимой базе: сравнением позиций он объясняется, и на этом
стоит урок класса C; создаёт его генератор намеренно
(~1–2% заказов, см. раздел 4).
Строка «`purchase` без заказа» внутри живого окна — опоздание, ждущее
слепка, а не расхождение: она получает служебный класс `awaiting_order`
@@ -484,11 +505,15 @@ Kafka день переигрывается генератором заново:
карты #10) закрыта этим же ходом: версионируется опись.
Контрольные числа описи:
- заказная сторона: заказы и выручка по дням; опись хранит точные
счётчики по каждому классу расхождений (отмены, потери, дубли, дельты сумм) —
самопроверка лабы сверки;
- идентичность: uniq кук, uniq известных пользователей, число двухкуковых
покупателей — лаба склейки получает самопроверку.
- по строке на каждый отправленный слепок — хеш байтов: побайтовое обещание
«слепок переснимается и даёт те же байты» опись сторожит наравне с днями;
- заказная сторона: заказы, выручка и счётчики наблюдаемых классов
расхождений после приоритета, в заказах — по дням с закрытым окном
(самопроверка лабы сверки; при N сыгранных днях таких дней N − 7);
- контрольные числа идентичности не заводятся: uniq известных пользователей
и число двухкуковых покупателей растут, пока мир едет, а счётчик без
названного читателя в опись не кладётся
([что хранит опись](../architecture/orders/inventory.md)).
Снимок вырастет против v1 (ecommerce-массивы, заказы) — размер проверить
при пересборке.
@@ -545,10 +570,12 @@ v2 стартует пустым, поэтому объём ниже — это
целиком). В конце этапа фиксируется маленький стартовый мир для
стабильных приёмок следующих этапов (полная пересборка эталонного мира —
отдельный этап 7).
3. Заказы и каталог: генератор слепков, STG/ODS/DDS заказа, словарь и первый
настоящий даг — проигрыш модельного дня плюс переливка слепка ([ADR
0008](../adr/0008-order-ingestion.md)).
4. Трансформации и витрины: сессии, identity_map, выручка, сверка A+C.
3. Заказы и каталог: генератор слепков, STG/ODS заказа, словарь и даг приёма
слепка ([ADR 0008](../adr/0008-order-ingestion.md)). Модель заказа в DDS
уехала этапу 4 — у служебной колонки и зерна читатель появляется там
(нарезка этапа, #89).
4. Трансформации и витрины: модель DDS заказа, сессии, identity_map,
выручка, сверка A+C.
5. Airflow: `etl_pipeline` (партиционная переобработка).
6. Расхождения B+D и опоздания; счётчики описи.
7. Эталонный мир: опись и пересборка снимка, чек-скрипты; CI-генерация
@@ -618,9 +645,11 @@ Kafka Engine на двух нодах снят с этого списка при
третьей версии (Datasets → Assets) — актуальные операторы проверить через
Context7. Первый даг приходит этапом 3, а не 5 ([ADR
0008](../adr/0008-order-ingestion.md)).
- Пакетный забор слепка из Kafka: коммит офсетов прямым чтением, хватает ли
одного чтения на слепок дня, доступны ли при нём виртуальные колонки
доставки — список и ответы в [ADR 0008](../adr/0008-order-ingestion.md).
- Пакетный забор слепка из Kafka: прямое чтение коммитит офсеты и возвращает
одну порцию. На стандартном мире около полутора тысяч заказов помещаются в
неё с запасом; это проверяемое допущение стенда, а не граница полноты слепка
([ADR 0008](../adr/0008-order-ingestion.md), [ADR
0010](../adr/0010-order-versions-in-ods.md)).
- Генератор: рабочее решение — Python с производительной архитектурой
(батчевая генерация вместо посточной, быстрая JSON-сериализация,
распараллеливание по модельным дням). Читаемость генератора для менти —
@@ -676,10 +705,10 @@ v2, этап 0).
наблюдаемости и того, чем платит каждый режим, — как задание. Третий способ,
типизированный чтец с `kafka_handle_error_mode`, на стенде не живёт: он
отвергнут обоими ADR, и остаётся материалом для рассказа;
- матвью как рабочий механизм, а не диковина: их видно на приёме и на сборке
ODS, а пакетная работа начинается выше. Отдельным заданием — как читать из ODS
последние версии, через `FINAL` или оконной функцией: что нагляднее, решаем на
месте;
- матвью как рабочий механизм, а не диковина: события проходят из STG в ODS на
лету, заказы — пакетным заданием. `ods.order_v` показывает границу между
физическими версиями и точным текущим состоянием; сравнение `FINAL` с
альтернативными способами чтения остаётся материалом задания;
- лекция про идентичность «как в бою»: `setUserID` и first-party id,
детерминированная против вероятностной склейки, identity graph,
кросс-девайс, CDP — с рамкой «мы склеили через транзакции, потому что трекер
+35 -4
View File
@@ -25,6 +25,32 @@
касается этапа 7 — расширение подтверждено владельцем в резолюции
«Производительности».
## Чем меряется генератор
Рамка владельца, по которой принимались решения ниже. Она шире этапа 2 и
живёт дальше него: спеки следующих этапов ссылаются сюда, а не переписывают.
- **Генератор меряется тем, что приезжает на выход, а не устройством
внутри.** Он играет два источника с разными темпами: поток Метрики и
ночную выгрузку бэкенда магазина. Любое внутреннее решение отвечает перед
одним вопросом — правдоподобны ли данные на выходе.
- **Менти внутрь не смотрит.** Он видит внешние эффекты и по ним составляет
представление, как работает магазин и обо что там спотыкаются. Значит,
конструкция внутри оправдана только наблюдаемым эффектом: без него это
сложность, за которую никто не заплатил.
- **Повторимы должны быть эффекты, а не конкретные числа.** Правка свойств
мира сдвигает почти все числа разом — генератор случайных чисел выдаёт
другую последовательность. Опираться на конкретные числа поэтому нечем.
Хеши описи обещают не постоянство чисел, а то, что их сдвиг не пройдёт
молча: пересобрал опись — увидел diff.
- **Цена ошибки здесь мала.** В худшем случае испорчен урок, и чинится он
правкой генератора после взгляда на получившиеся распределения. Оборона,
которая дороже ошибки, не заводится.
- **Сложность изолирована нарочно.** Генератор вынесен отдельной службой,
чтобы его сложность не расползалась по стенду. Плата за изоляцию —
дешевизна правки: выбирается устройство, которое дёшево менять, а не то,
которое всё предусмотрело.
## Целевая картина одним взглядом
- **Мир — функция, не состояние.** Состав мира — чистая функция зерна;
@@ -93,9 +119,10 @@
человеку предыстории, чьё окно активности таких дней не оставляет,
пара не назначается. День-функция обязана назначенные заказы
реализовать; остальные покупки — вольные, их решает генератор торговых
событий (#40). Опись считает пары, реализованные в горизонте
снимка. Условность в данных не видна: дни назначены той же
случайностью, просто брошенной планом один раз.
событий (#40). План считает пары, реализованные в запрошенном горизонте;
опись это число не хранит — оно растёт вместе с миром. Условность в данных
не видна: дни назначены той же случайностью, просто брошенной планом один
раз.
- **Своя ось модельного времени.** Ось событий начинается в
фиксированный день D0 (понедельник — см. раздел 5); реальный
календарь в модели не участвует. В `EventDate`/`UTCEventTime` дни оси
@@ -171,7 +198,11 @@
целое, поэтому у предыстории своя ветвь состава, отдельная от оси
(уточнение при исполнении #38, 2026-08-02);
(зерно, день) → подпоток дня → именованные подпотоки компонентов: трафик,
торговые события, расхождения, опоздания — в фиксированном порядке.
торговые события, заказная и событийная стороны (судьба заказов и порча
событийного потока) — в фиксированном порядке (уточнение [судьбой
заказа](../architecture/orders/fate.md): прежние «расхождения» и
«опоздания» были названы по классам витрины, а компонент называет часть
мира).
По построению: параллельный прогон равен последовательному; продление
истории днём N+1 не трогает дни 1…N; правка одного компонента меняет
только его часть снимка — в описи меняются хеши только затронутых
+3
View File
@@ -4,6 +4,9 @@ USER root
# Официальный провайдер Kafka использует тот же confluent-kafka. Пробникам
# не нужны его подключения и обёртки, поэтому в образ добавлен сам клиент.
#
# Провайдера docker, которым пульт мира зовёт генератор, ставить не нужно:
# база несёт его сама — 4.5.7 к Airflow 3.3.0, замер 13 августа 2026 года.
RUN uv pip install --python /home/airflow/.local/bin/python --no-cache \
"clickhouse-connect==1.6.0" \
"confluent-kafka==2.15.0" \