Files
clickstream-data-platform/docs/specs/2026-08-16-order-ingestion.md
T
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

171 lines
13 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# Приём заказов из 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.