docs(orders): зафиксирован версионный приём заказов

- Зачем:
  - отменённая подмена партиции snapshot_date противоречила порционному чтению Kafka и могла обучать потере ранее принятых версий.
- Что:
  - добавлены спецификация приёма заказов и ADR о версионном ODS с ods.order_v.
  - согласованы мастер-спека, дока хранилища, ADR 0008 и исследование формата.
  - зафиксированы граница приёма, координаты загрузки, диагностические повторы и отложенное проектирование DDS.
- Проверка:
  - git diff --cached --check.
  - горячее ревью по правилам репозитория и принятому решению.
  - два прохода холодного ревью.
This commit is contained in:
2026-08-16 21:42:11 +03:00
parent 8d54a3caba
commit ede1df765c
6 changed files with 311 additions and 65 deletions
+39 -38
View File
@@ -219,23 +219,19 @@ Ecommerce (заполнены только у торговых событий):
Обоснование и проверка разбора — в
[исследовании формата](../research/2026-08-16-order-snapshot-wire-format.md).
- Приём идемпотентный, но дедуп расщеплён на два слоя:
- `ods.order_snapshot` — партиция по `snapshot_date`, **без дедупа**,
хранит «как приехало»; идемпотентность повторного прогона — заменой
партиции дня слепка, а не ReplacingMergeTree.
- Дедуп до последней версии — **argMax** в трансформации при сборке
`dds.order`. `dds.order` — единственная дедуплицированная таблица:
партиция по дню создания строки источника (`toDate(created_at)`) — это
стабильный технический ключ, а не бизнес-день покупки,
ReplacingMergeTree(`updated_at`), `ORDER BY order_id` — заказ всегда
лежит в одной партиции, дедуп работает.
- `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.
@@ -360,18 +356,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` (раньше называлась
@@ -409,10 +405,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` | см. ниже |
@@ -422,19 +418,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` — свойство содержимого, а содержимое
@@ -636,9 +635,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-сериализация,
распараллеливание по модельным дням). Читаемость генератора для менти —
@@ -694,10 +695,10 @@ v2, этап 0).
наблюдаемости и того, чем платит каждый режим, — как задание. Третий способ,
типизированный чтец с `kafka_handle_error_mode`, на стенде не живёт: он
отвергнут обоими ADR, и остаётся материалом для рассказа;
- матвью как рабочий механизм, а не диковина: их видно на приёме и на сборке
ODS, а пакетная работа начинается выше. Отдельным заданием — как читать из ODS
последние версии, через `FINAL` или оконной функцией: что нагляднее, решаем на
месте;
- матвью как рабочий механизм, а не диковина: события проходят из STG в ODS на
лету, заказы — пакетным заданием. `ods.order_v` показывает границу между
физическими версиями и точным текущим состоянием; сравнение `FINAL` с
альтернативными способами чтения остаётся материалом задания;
- лекция про идентичность «как в бою»: `setUserID` и first-party id,
детерминированная против вероятностной склейки, identity graph,
кросс-девайс, CDP — с рамкой «мы склеили через транзакции, потому что трекер
+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.