From 8d54a3cabac6ecd4ed3e28e0db3120ed959bbd50 Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Sun, 16 Aug 2026 17:40:23 +0300 Subject: [PATCH] =?UTF-8?q?docs(orders):=20=D1=83=D1=82=D0=BE=D1=87=D0=BD?= =?UTF-8?q?=D1=91=D0=BD=20=D1=84=D0=BE=D1=80=D0=BC=D0=B0=D1=82=20=D1=81?= =?UTF-8?q?=D0=BB=D0=B5=D0=BF=D0=BA=D0=B0=20=D0=B8=20=D0=B2=D1=80=D0=B5?= =?UTF-8?q?=D0=BC=D0=B5=D0=BD=D0=BD=D1=8B=D0=B5=20=D0=BF=D0=BE=D0=BB=D1=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - аудит источника нельзя смешивать с бизнес-временем и `_load_ts` ClickHouse. - Что: - зафиксированы JSON-контракт слепка, порядок ключей и намеренное различие `Array(Float64)` и `Decimal`. - разведены бизнес-время, аудит источника и загрузка; уточнены слой ODS, `snapshot_date` и технический ключ партиции. - добавлены термины словаря и исследование точного миллисекундного формата с проверками ClickHouse. - Проверка: - выполнен `git diff --cached --check`. --- CONTEXT.md | 14 +- docs/architecture/storage.md | 39 ++-- .../2026-08-16-order-snapshot-wire-format.md | 210 ++++++++++++++++++ docs/specs/2026-07-30-stand-v2-realism.md | 32 ++- 4 files changed, 271 insertions(+), 24 deletions(-) create mode 100644 docs/research/2026-08-16-order-snapshot-wire-format.md diff --git a/CONTEXT.md b/CONTEXT.md index ecc0094..f6230ba 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -162,8 +162,8 @@ Python-модуль с описателями колонок события — _Избегать_: описание схемы, документация контракта **Канонический сериализатор**: -Единственное место, где событие превращается в байты. Фиксированный порядок -ключей и строк — основа побайтовой воспроизводимости. +Единственная граница, где запись генератора превращается в байты. Фиксированный +порядок ключей и строк — основа побайтовой воспроизводимости. **Проигрыватель**: Компонент доставки готового потока дня в приёмник. Два режима: пакетный @@ -175,6 +175,16 @@ _Избегать_: описание схемы, документация кон может отсутствовать в ранних слепках; после первого появления ездит до конца своего окна. +**Дата слепка**: +Дата завершившегося модельного дня, состояние которого снято на исходящей +границе суток. Одна для всех записей выгрузки; не дата заказа и не время +отправки сообщения. + +**Время аудита источника**: +Время создания или последнего изменения строки заказа по часам базы источника. +В записи слепка это `created_at` и `updated_at`; бизнес-время покупки живёт в +событии `purchase` отдельно. + **Окно изменяемости**: Сколько модельных дней заказ ещё может измениться и потому продолжает ездить в слепках. Константа мира: семь дней. За окном заказ замёрз, выручка дня diff --git a/docs/architecture/storage.md b/docs/architecture/storage.md index caf8f1f..74325f9 100644 --- a/docs/architecture/storage.md +++ b/docs/architecture/storage.md @@ -116,8 +116,9 @@ 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)). Движок таблицы сырья — обычный `ReplicatedMergeTree`, `ORDER BY (kafka_partition, kafka_offset)`: разбор полётов идёт от «какое сообщение», другого ключа у сырья и @@ -125,15 +126,23 @@ kafka_offset)`: разбор полётов идёт от «какое сооб которого он заведён: повтор доставки в сырье обязан быть виден. Метка времени загрузки зовётся `_load_ts`, тип `DateTime64(3, 'UTC')`. Ставится -она один раз, в матвью приёма, и дальше переносится из STG в ODS как есть: -колонка отвечает на вопрос «когда строка приехала в хранилище», а не «когда её -разобрали». В ODS она же служит колонкой версии `ReplacingMergeTree`, и работа у -этой версии ровно одна — схлопнуть повтор доставки. Содержимое у повтора то же -самое, отличается только метка, поэтому какая из двух строк переживёт мерж, -безразлично. Пакетной переобработки у ODS нет: слой наполняет матвью, а не -задание Airflow, и работа с партициями начинается выше. Переделать разобранное -руками можно — вставкой из сырья с фильтром по `_load_ts`, в пределах -трёхсуточного окна; ничья по версии разрешается в пользу вставленного позже. +она один раз при записи в STG: для событий — матвью приёма, для дневного слепка +заказов — пакетным шагом. Дальше метка переносится в ODS как есть и отвечает на +вопрос «когда строка приехала в хранилище», а не «когда её разобрали». В +`ods.event` она же служит колонкой версии `ReplacingMergeTree` и схлопывает +повтор доставки. `ods.order_snapshot` повтор не схлопывает: пакетный шаг +заменяет целиком партицию `snapshot_date`. + +`created_at` и `updated_at` заказа к служебным колонкам хранилища не относятся. +Они приезжают в сообщении как аудит строки в БД источника и в ODS разбираются в +`DateTime64(3, 'UTC')`; ClickHouse их не создаёт и добавляет рядом собственную +`_load_ts`. Совпадение слов «техническое время» не делает эти часы одной осью. + +Пакетной переобработки у событий в ODS нет: слой наполняют матвью, а не задание +Airflow. Переделать разобранное руками можно вставкой из сырья с фильтром по +`_load_ts` в пределах трёхсуточного окна; ничья по версии разрешается в пользу +вставленного позже. Заказы, напротив, по построению перерабатываются дневными +партициями пакетного шага. Имя согласовано с каноном служебных полей соседнего учебного стенда на Greenplum, чтобы словарь был общим у двух хранилищ; ведущее подчёркивание у @@ -141,10 +150,10 @@ Greenplum, чтобы словарь был общим у двух хранил её же используют Fivetran, Airbyte и Stitch. С правилом выше это не спорит: запрещено совпадать с именами виртуальных колонок, а не носить подчёркивание. -Идентификатора пачки загрузки (`_load_id`) пока нет. В STG и ODS данные приезжают -потоком через матвью, у которого нет ни батча, ни `run_id`, и колонка была бы -пустой формальностью. В слоях, которые наполняет Airflow, `run_id` появится -по-настоящему — тогда и заведём, тем же стилем имени. +Общего идентификатора пачки загрузки (`_load_id`) нет. У потока событий нет ни +пачки, ни `run_id`, и колонка была бы пустой формальностью. Пакетный шаг заказов +не делает её общей конвенцией: конкретный слой заведёт `run_id`, только когда у +него появится названный читатель этой координаты. ## Часовые пояса diff --git a/docs/research/2026-08-16-order-snapshot-wire-format.md b/docs/research/2026-08-16-order-snapshot-wire-format.md new file mode 100644 index 0000000..0167770 --- /dev/null +++ b/docs/research/2026-08-16-order-snapshot-wire-format.md @@ -0,0 +1,210 @@ +# Формат дневного слепка заказов на проводе + +Дата исследования: 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)` в +[мастер-спеке](../specs/2026-07-30-stand-v2-realism.md) используется только как +стабильный технический ключ партиции `dds.order`: это день создания строки +источника, а не доказательство дня бизнес-события. + +Минимальное решение — не добавлять `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` уже является языком спеки, партиции ODS и 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` дали одно значение. Проверка денежного формата +не нужна сериализатору этого слепка: он сам выпускает ровно два знака. Считать +ли остальные формы браком при приёме, решает задача #80; это решение отдельного +лексического валидатора не требует. diff --git a/docs/specs/2026-07-30-stand-v2-realism.md b/docs/specs/2026-07-30-stand-v2-realism.md index cf79609..a4a6115 100644 --- a/docs/specs/2026-07-30-stand-v2-realism.md +++ b/docs/specs/2026-07-30-stand-v2-realism.md @@ -190,17 +190,34 @@ 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 | | +| `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 | завершившийся модельный день, состояние которого снято на исходящей границе суток | + + На проводе один заказ — один документ 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` — партиция по `snapshot_date`, **без дедупа**, @@ -208,7 +225,8 @@ Ecommerce (заполнены только у торговых событий): партиции дня слепка, а не ReplacingMergeTree. - Дедуп до последней версии — **argMax** в трансформации при сборке `dds.order`. `dds.order` — единственная дедуплицированная таблица: - партиция по дню заказа (`toDate(created_at)`), + партиция по дню создания строки источника (`toDate(created_at)`) — это + стабильный технический ключ, а не бизнес-день покупки, ReplacingMergeTree(`updated_at`), `ORDER BY order_id` — заказ всегда лежит в одной партиции, дедуп работает. @@ -394,7 +412,7 @@ README. | ODS | `ods.order_snapshot` (+`_errors`) | слепки заказов как приехали, партиция по `snapshot_date`, без дедупа | | 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.order` | единственная дедуплицированная таблица заказа: партиция по техническому дню создания строки источника (`toDate(created_at)`), ReplacingMergeTree(`updated_at`), `ORDER BY order_id`, дедуп до последней версии — argMax в трансформации при сборке | | DDS | `dds.identity_map` | карта кука↔пользователь | | DDS | словарь `products` | каталог из CSV | | DM | витрины `dm.*_v`, `dm.dq_summary` | см. ниже |