# Приём заказов из 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.