From ede1df765c4c8360c768285d5bf26a6790fcbe4e Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Sun, 16 Aug 2026 21:04:29 +0300 Subject: [PATCH] =?UTF-8?q?docs(orders):=20=D0=B7=D0=B0=D1=84=D0=B8=D0=BA?= =?UTF-8?q?=D1=81=D0=B8=D1=80=D0=BE=D0=B2=D0=B0=D0=BD=20=D0=B2=D0=B5=D1=80?= =?UTF-8?q?=D1=81=D0=B8=D0=BE=D0=BD=D0=BD=D1=8B=D0=B9=20=D0=BF=D1=80=D0=B8?= =?UTF-8?q?=D1=91=D0=BC=20=D0=B7=D0=B0=D0=BA=D0=B0=D0=B7=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - отменённая подмена партиции snapshot_date противоречила порционному чтению Kafka и могла обучать потере ранее принятых версий. - Что: - добавлены спецификация приёма заказов и ADR о версионном ODS с ods.order_v. - согласованы мастер-спека, дока хранилища, ADR 0008 и исследование формата. - зафиксированы граница приёма, координаты загрузки, диагностические повторы и отложенное проектирование DDS. - Проверка: - git diff --cached --check. - горячее ревью по правилам репозитория и принятому решению. - два прохода холодного ревью. --- docs/adr/0008-order-ingestion.md | 18 +- docs/adr/0010-order-versions-in-ods.md | 37 ++++ docs/architecture/storage.md | 58 ++++-- .../2026-08-16-order-snapshot-wire-format.md | 16 +- docs/specs/2026-07-30-stand-v2-realism.md | 77 ++++---- docs/specs/2026-08-16-order-ingestion.md | 170 ++++++++++++++++++ 6 files changed, 311 insertions(+), 65 deletions(-) create mode 100644 docs/adr/0010-order-versions-in-ods.md create mode 100644 docs/specs/2026-08-16-order-ingestion.md diff --git a/docs/adr/0008-order-ingestion.md b/docs/adr/0008-order-ingestion.md index 1c93f50..d6f6698 100644 --- a/docs/adr/0008-order-ingestion.md +++ b/docs/adr/0008-order-ingestion.md @@ -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 описана в +[спецификации приёма заказов](../specs/2026-08-16-order-ingestion.md). ## Решение @@ -128,7 +135,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 +147,9 @@ - Прямое чтение офсеты по умолчанию **не** коммитит; коммит включается настройкой `kafka_commit_on_select` на самой таблице. Это тот подводный камень, который молчит на первом прогоне и вылезает на втором. -- Сколько строк отдаёт одно чтение, задаёт `kafka_max_block_size`. При слепке - порядка полутора тысяч строк это один блок с запасом. +- Прямое чтение возвращает одну порцию, полученную одним опросом Kafka. При + настройках стенда её предел — 65 409 сообщений, поэтому слепок порядка + полутора тысяч строк помещается с запасом. - Про `Distributed` поверх Kafka документация не говорит ничего — ни поддержки, ни запрета. diff --git a/docs/adr/0010-order-versions-in-ods.md b/docs/adr/0010-order-versions-in-ods.md new file mode 100644 index 0000000..6668cae --- /dev/null +++ b/docs/adr/0010-order-versions-in-ods.md @@ -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 собрана полная партиция, и её подмена могла бы удалить уже принятые +версии прошлого дня. Маркер конца, опись или фиксация конечных офсетов сделали +бы границу настоящей, но добавили бы новый протокол без нужного стенду урока. + +Малый объём позволяет оставить одно чтение на запуск как проверяемое +эксплуатационное допущение, а не границу полноты. Полный контракт разбора, +граница брака и поведение повторов заданы в +[спецификации приёма заказов](../specs/2026-08-16-order-ingestion.md). diff --git a/docs/architecture/storage.md b/docs/architecture/storage.md index 74325f9..8e4bd03 100644 --- a/docs/architecture/storage.md +++ b/docs/architecture/storage.md @@ -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 заранее. ## Служебные колонки @@ -118,7 +119,8 @@ keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/` Само сообщение лежит в колонке `raw` тем, чем пришло: чтец читает байты и ничего не проверяет, поэтому там оказываются и целые сообщения, и мусор. События ниже разбирают матвью ODS ([ADR 0005](../adr/0005-event-ingestion.md)), заказы — -пакетный шаг ([ADR 0008](../adr/0008-order-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)`: разбор полётов идёт от «какое сообщение», другого ключа у сырья и @@ -130,8 +132,8 @@ kafka_offset)`: разбор полётов идёт от «какое сооб заказов — пакетным шагом. Дальше метка переносится в ODS как есть и отвечает на вопрос «когда строка приехала в хранилище», а не «когда её разобрали». В `ods.event` она же служит колонкой версии `ReplacingMergeTree` и схлопывает -повтор доставки. `ods.order_snapshot` повтор не схлопывает: пакетный шаг -заменяет целиком партицию `snapshot_date`. +повтор доставки. У `ods.order_snapshot` версию задаёт `updated_at` источника; +`_load_ts` только показывает, когда конкретная строка приехала. `created_at` и `updated_at` заказа к служебным колонкам хранилища не относятся. Они приезжают в сообщении как аудит строки в БД источника и в ODS разбираются в @@ -141,8 +143,8 @@ kafka_offset)`: разбор полётов идёт от «какое сооб Пакетной переобработки у событий в ODS нет: слой наполняют матвью, а не задание Airflow. Переделать разобранное руками можно вставкой из сырья с фильтром по `_load_ts` в пределах трёхсуточного окна; ничья по версии разрешается в пользу -вставленного позже. Заказы, напротив, по построению перерабатываются дневными -партициями пакетного шага. +вставленного позже. Заказы разбирает пакетный шаг из неизменного среза STG по +`_load_id`; дневные партиции ODS он не заменяет. Имя согласовано с каноном служебных полей соседнего учебного стенда на Greenplum, чтобы словарь был общим у двух хранилищ; ведущее подчёркивание у @@ -151,9 +153,27 @@ Greenplum, чтобы словарь был общим у двух хранил запрещено совпадать с именами виртуальных колонок, а не носить подчёркивание. Общего идентификатора пачки загрузки (`_load_id`) нет. У потока событий нет ни -пачки, ни `run_id`, и колонка была бы пустой формальностью. Пакетный шаг заказов -не делает её общей конвенцией: конкретный слой заведёт `run_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) и +[спецификации приёма заказов](../specs/2026-08-16-order-ingestion.md). ## Часовые пояса @@ -367,6 +387,12 @@ D0 и к реальному календарю не привязана; паке kafka_offset)`: смотрят такую таблицу от класса, а внутри класса — по координатам доставки. +`ods.order_snapshot_errors` держит тот же диагностический минимум и `_load_id` +запуска. У заказов три класса по приоритету: `not_an_object`, +`keyset_mismatch`, `field_invalid`. Сырой текст остаётся рядом, поэтому класс не +разрастается до имени отдельного поля. Точная граница приёма — в +[спецификации заказов](../specs/2026-08-16-order-ingestion.md). + ## Раскладка DDL Файлы лежат в `sql/ddl/` и применяются по порядку имён. Сначала все статичные diff --git a/docs/research/2026-08-16-order-snapshot-wire-format.md b/docs/research/2026-08-16-order-snapshot-wire-format.md index 0167770..1b3cfd4 100644 --- a/docs/research/2026-08-16-order-snapshot-wire-format.md +++ b/docs/research/2026-08-16-order-snapshot-wire-format.md @@ -63,8 +63,9 @@ Debezium проводит ту же границу внутри одного с `created_at` и `updated_at` — аудит строки источника. Поэтому `toDate(created_at)` в -[мастер-спеке](../specs/2026-07-30-stand-v2-realism.md) используется только как -стабильный технический ключ партиции `dds.order`: это день создания строки +[спецификации приёма заказов](../specs/2026-08-16-order-ingestion.md) +используется как стабильный технический ключ партиции `ods.order_snapshot`: +это день создания строки источника, а не доказательство дня бизнес-события. Минимальное решение — не добавлять `ordered_at` на всякий случай. Пока модель @@ -155,8 +156,8 @@ parseDateTime64InJodaSyntaxOrNull( слепок с `snapshot_date = ORIGIN` уезжает прогоном дня 1. `as_of_date` тоже могло бы означать дату, по состоянию на которую показаны -данные. Но `snapshot_date` уже является языком спеки, партиции ODS и ADR о -приёме заказов. Переименование не добавляет урока и может спутать дату +данные. Но `snapshot_date` уже является языком спеки и ADR о приёме заказов. +Переименование не добавляет урока и может спутать дату выгрузки с периодом бизнес-действия записи. Для этого стенда оставляем `snapshot_date`. @@ -205,6 +206,7 @@ KISS-вариант для [существующего модуля сериал Граница этого решения: `toDecimal64OrNull(..., 2)` проверяет числовую преобразуемость, но не лексическое правило «ровно два знака» — локально строки `1299.90`, `1299.9` и `1299.900` дали одно значение. Проверка денежного формата -не нужна сериализатору этого слепка: он сам выпускает ровно два знака. Считать -ли остальные формы браком при приёме, решает задача #80; это решение отдельного -лексического валидатора не требует. +не нужна сериализатору этого слепка: он сам выпускает ровно два знака. Приём ODS +проверяет ту же каноническую форму и считает остальные формы браком; граница +строгого приёма зафиксирована в +[спецификации заказов](../specs/2026-08-16-order-ingestion.md). diff --git a/docs/specs/2026-07-30-stand-v2-realism.md b/docs/specs/2026-07-30-stand-v2-realism.md index a4a6115..8ee94c4 100644 --- a/docs/specs/2026-07-30-stand-v2-realism.md +++ b/docs/specs/2026-07-30-stand-v2-realism.md @@ -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 — с рамкой «мы склеили через транзакции, потому что трекер diff --git a/docs/specs/2026-08-16-order-ingestion.md b/docs/specs/2026-08-16-order-ingestion.md new file mode 100644 index 0000000..265a43a --- /dev/null +++ b/docs/specs/2026-08-16-order-ingestion.md @@ -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.