From bdcbe48b595a0d561dd67befa21b91bac84cfe63 Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Tue, 18 Aug 2026 23:10:01 +0300 Subject: [PATCH] =?UTF-8?q?feat(orders):=20=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2?= =?UTF-8?q?=D0=BB=D0=B5=D0=BD=20=D0=B2=D0=B5=D1=80=D1=81=D0=B8=D0=BE=D0=BD?= =?UTF-8?q?=D0=BD=D1=8B=D0=B9=20ODS=20=D0=B7=D0=B0=D0=BA=D0=B0=D0=B7=D0=BE?= =?UTF-8?q?=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - менти должен различать версию заказа, наблюдение источника и запуск загрузки. - Что: - добавлены таблицы версий и брака заказов, а также представление ods.order_v. - даг orders_ingest дополнен строгим переходом одного среза STG в две цели ODS. - результаты опытов записаны в документации, дефект генератора вынесен в #102. - Проверка: - make lint config-test smoke check-clickhouse check-services. --- dags/orders_ingest.py | 177 ++++++++++++++++++++++++-- docs/architecture/orders/README.md | 3 - docs/architecture/orders/ingestion.md | 25 +++- docs/architecture/storage.md | 61 ++++++++- sql/ddl/20-ods-tables.sql | 86 +++++++++++++ sql/ddl/30-ods-views.sql | 16 +++ 6 files changed, 342 insertions(+), 26 deletions(-) diff --git a/dags/orders_ingest.py b/dags/orders_ingest.py index 49da20b..5f4324f 100644 --- a/dags/orders_ingest.py +++ b/dags/orders_ingest.py @@ -1,8 +1,9 @@ """Приём заказов: пакетный забор слепка из Kafka в хранилище. -Путь Kafka → STG → ODS у заказов один и живёт одним дагом; здесь его первый -шаг — забор. Расписания у дага нет: его дёргает тот, кто положил слепок в -топик, — работник пульта мира после проигрыша дня, — и ждёт конца. +Путь Kafka → STG → ODS у заказов один и живёт одним дагом: сначала забор +порции в сырьё, затем разбор того же среза в типизированный ODS и в таблицу +брака. Расписания у дага нет: его дёргает тот, кто положил слепок в топик, — +работник пульта мира после проигрыша дня, — и ждёт конца. Контраст с приёмом событий и есть урок. События тянет матвью, навсегда подписанная на чтеца: приём идёт сам, пока идёт поток. Слепок заказов @@ -10,8 +11,8 @@ и забирает её запрос по команде. Push против pull — два режима на одном стенде, каждый там, где ему место по природе источника. -Решения и доводы целиком — ADR 0008; форма перехода STG → ODS, который -прирастёт сюда вторым шагом, — docs/architecture/orders/ingestion.md. +Решения и доводы целиком — ADR 0008 (забор) и ADR 0010 (версии в ODS); форма +перехода — docs/architecture/orders/ingestion.md. """ from __future__ import annotations @@ -50,6 +51,133 @@ SETTINGS distributed_foreground_insert = 1 """ +# Строгий приём: проверяется форма провода и ничего сверх неё. +# +# Ключей одиннадцать, и проверяются все: десять скалярных значений прямо +# образуют типизированную строку заказа. Внутрь items приём не смотрит — +# содержимое позиций, переходы статуса и равенства сумм остаются ниже границы +# (docs/architecture/orders/ingestion.md, «Граница строгого приёма»). +# +# Блок общий у обеих вставок нарочно: годная ветвь берёт row_is_valid, брак — +# буквальное NOT. Разойдись условия хоть на символ — строка либо задвоится, +# либо исчезнет молча. +# +# Отсюда два запрета на выражения предиката, и оба серьёзные. Первый: ни одно +# не возвращает NULL — трёхзначная логика дала бы строку, которую не берёт ни +# условие, ни его отрицание. Поэтому сравнения дают 0 или 1, а обнуляемый +# разбор заканчивается IS NOT NULL. Второй: ни одно не бросает исключений на +# произвольном raw — иначе одна грязная строка роняет весь переход, ради +# отсутствия чего таблица брака и заведена. +WIRE_CONTRACT = r""" +WITH + -- Каноническое время провода: UTC, ровно три знака долей секунды. + 'yyyy-MM-dd\'T\'HH:mm:ss.SSS\'Z\'' AS ts_mask, + JSONType(raw) = 'Object' AS is_object, + arraySort(JSONExtractKeys(raw)) = arraySort([ + 'order_id', 'user_id', 'status', 'created_at', 'updated_at', + 'items_total', 'discount', 'delivery', 'total', 'items', + 'snapshot_date' + ]) AS keys_match, + JSONType(raw, 'order_id') = 'String' + AND JSONType(raw, 'status') = 'String' + -- Целое у JSONType зовётся двумя именами, и нужны оба: с одним + -- Int64 законный идентификатор за 2^63 уехал бы в брак. + AND JSONType(raw, 'user_id') IN ('Int64', 'UInt64') + AND JSONExtract(raw, 'user_id', 'Nullable(UInt64)') IS NOT NULL + AND JSONType(raw, 'items') = 'Array' + -- Форму держит регулярка, реальность — разбор, и порознь они дырявы: + -- маска берёт «2026-6-3T…» без ведущих нулей, а регулярка пропускает + -- 30 февраля. Обе половины измерены — docs/architecture/storage.md, + -- «Что проверено». + AND arrayAll(k -> + JSONType(raw, k) = 'String' + AND match(JSONExtractString(raw, k), + '^\\d{4}-\\d{2}-\\d{2}T\\d{2}:\\d{2}:\\d{2}\\.\\d{3}Z$') + AND parseDateTime64InJodaSyntaxOrNull( + JSONExtractString(raw, k), ts_mask, 'UTC') IS NOT NULL, + ['created_at', 'updated_at']) + -- Деньги — строка с ровно двумя знаками после точки. Регулярка держит + -- форму, разбор — вместимость: у строки из двадцати девяток форма + -- хороша, а в Decimal(18, 2) она не влезает. + AND arrayAll(k -> + JSONType(raw, k) = 'String' + AND match(JSONExtractString(raw, k), '^\\d+\\.\\d{2}$') + AND toDecimal64OrNull(JSONExtractString(raw, k), 2) IS NOT NULL, + ['items_total', 'discount', 'delivery', 'total']) + AND JSONType(raw, 'snapshot_date') = 'String' + AND match(JSONExtractString(raw, 'snapshot_date'), '^\\d{4}-\\d{2}-\\d{2}$') + AND parseDateTimeInJodaSyntaxOrNull( + JSONExtractString(raw, 'snapshot_date'), 'yyyy-MM-dd', 'UTC') IS NOT NULL + AS fields_valid, + is_object AND keys_match AND fields_valid AS row_is_valid +""" + +# Годные версии заказа. Срез читается по _load_id — тому же, что проставил +# забор: разбор идёт по неизменной порции, а не по «всему, что появилось». +# +# Разобранные значения берёт assumeNotNull: обнуляемый разбор стоит за +# предикатом, который NULL уже отсёк, и приведение здесь не может упасть. +PARSE_GOOD_ROWS = ( + "INSERT INTO ods.order_snapshot_dist" + + WIRE_CONTRACT + + r""" +SELECT + JSONExtractString(raw, 'order_id') AS order_id, + JSONExtract(raw, 'user_id', 'UInt64') AS user_id, + JSONExtractString(raw, 'status') AS status, + assumeNotNull(parseDateTime64InJodaSyntaxOrNull( + JSONExtractString(raw, 'created_at'), ts_mask, 'UTC')) AS created_at, + assumeNotNull(parseDateTime64InJodaSyntaxOrNull( + JSONExtractString(raw, 'updated_at'), ts_mask, 'UTC')) AS updated_at, + assumeNotNull(toDecimal64OrNull( + JSONExtractString(raw, 'items_total'), 2)) AS items_total, + assumeNotNull(toDecimal64OrNull(JSONExtractString(raw, 'discount'), 2)) AS discount, + assumeNotNull(toDecimal64OrNull(JSONExtractString(raw, 'delivery'), 2)) AS delivery, + assumeNotNull(toDecimal64OrNull(JSONExtractString(raw, 'total'), 2)) AS total, + -- items кладётся сырым фрагментом JSON, а не разобранной структурой. + JSONExtractRaw(raw, 'items') AS items, + toDate(assumeNotNull(parseDateTimeInJodaSyntaxOrNull( + JSONExtractString(raw, 'snapshot_date'), + 'yyyy-MM-dd', 'UTC'))) AS snapshot_date, + -- Метки запуска и прибытия переносятся как есть. Поставь здесь now64(3) — + -- и _load_ts молча ответила бы на другой вопрос: «когда разобрали». + _load_id, + _load_ts +FROM stg.orders_raw_dist +WHERE _load_id = {load_id:String} AND row_is_valid +SETTINGS distributed_foreground_insert = 1 +""" +) + +# Брак: тот же срез и буквальное отрицание того же предиката. +# +# Классы перекрываются, поэтому проверяются по порядку, а в error_class идёт +# первый совпавший: скаляр проваливает и проверку на объект, и сверку ключей — +# без объявленного порядка он попал бы то в один класс, то в другой. +PARSE_BAD_ROWS = ( + "INSERT INTO ods.order_snapshot_errors_dist" + + WIRE_CONTRACT + + r""" +SELECT + raw, + multiIf( + NOT is_object, 'not_an_object', + NOT keys_match, 'keyset_mismatch', + 'field_invalid' + ) AS error_class, + kafka_topic, + kafka_partition, + kafka_offset, + kafka_timestamp, + consumer_host, + _load_id, + _load_ts +FROM stg.orders_raw_dist +WHERE _load_id = {load_id:String} AND NOT row_is_valid +SETTINGS distributed_foreground_insert = 1 +""" +) + @dag( dag_id="orders_ingest", @@ -63,17 +191,16 @@ SETTINGS tags=["заказы"], ) def orders_ingest(): - """Забрать приехавший слепок заказов из топика в сырьё.""" + """Забрать приехавший слепок заказов и разложить его по слоям.""" - @task - def pull_batch() -> None: - # Импорт внутри задачи: обработчик DAG разбирает этот файл снова и - # снова, и импорт наверху оплачивался бы каждым разбором. + def clickhouse_client(): + # Импорт при выполнении задачи, а не при разборе файла: обработчик DAG + # разбирает его снова и снова, и импорт наверху оплачивался бы каждым + # разбором. import clickhouse_connect - load_id = get_current_context()["run_id"] connection = Connection.get("clickhouse_default") - client = clickhouse_connect.get_client( + return clickhouse_connect.get_client( host=connection.host, port=connection.port, username=connection.login, @@ -82,6 +209,11 @@ def orders_ingest(): connect_timeout=5, send_receive_timeout=30, ) + + @task + def pull_batch() -> None: + load_id = get_current_context()["run_id"] + client = clickhouse_client() try: summary = client.command(TAKE_ONE_BATCH, parameters={"load_id": load_id}) finally: @@ -94,7 +226,26 @@ def orders_ingest(): load_id, ) - pull_batch() + @task + def parse_batch() -> None: + """Разобрать срез сырья в версии заказов и в брак. + + Обе вставки в одном task_id: транзакции между ними ClickHouse не даёт, + а повтор задачи безопасен — срез читается по тому же неизменному + _load_id (docs/architecture/orders/ingestion.md, «Поток данных»). + """ + load_id = get_current_context()["run_id"] + client = clickhouse_client() + try: + client.command(PARSE_GOOD_ROWS, parameters={"load_id": load_id}) + client.command(PARSE_BAD_ROWS, parameters={"load_id": load_id}) + finally: + client.close() + # Счётчиков строк нет: у запроса с WHERE read_rows считает прочитанное + # с диска, а не подошедшее (storage.md, «Что проверено»). + logging.info("срез разобран: _load_id %s", load_id) + + pull_batch() >> parse_batch() orders_ingest() diff --git a/docs/architecture/orders/README.md b/docs/architecture/orders/README.md index 46bf515..53af499 100644 --- a/docs/architecture/orders/README.md +++ b/docs/architecture/orders/README.md @@ -57,9 +57,6 @@ доли классов, стоимость доставки — калибровка при реализации; финальная фиксация чисел — пересборка эталонного мира, этап 7. При пересборке правки потребуют только числа, не устройство. -- **Проверки приёма** — опыты из [«Рисков и проверки»](ingestion.md) про брак - и версии в ODS — тикет перехода STG → ODS (#94). Допущение «один запуск — - одно чтение» принято живым прогоном при исполнении #93. - **`_load_id` выше ODS** — вместе с устройством `dds.order` (#85). - **Контур проверок качества для расхождений** (даг DQ, `dm.dq_summary`) — остаётся в тумане карты #69; естественное место разговора — этап 4. diff --git a/docs/architecture/orders/ingestion.md b/docs/architecture/orders/ingestion.md index 787abcb..7ccbfbd 100644 --- a/docs/architecture/orders/ingestion.md +++ b/docs/architecture/orders/ingestion.md @@ -60,10 +60,10 @@ JSON-объектом с точным набором ключей: `order_id`, ` `items`, `snapshot_date`. Скалярные поля проверяются по типу JSON. Деньги дополнительно обязаны быть -строками с ровно двумя знаками после точки, времена — строками RFC 3339 в UTC с -обязательными миллисекундами, дата слепка — строкой `YYYY-MM-DD`. `items` -проверяется только как JSON-массив. Каноническая форма и основания выбора -зафиксированы в +строками неотрицательной суммы с ровно двумя знаками после точки — минус в эту +форму не входит; времена — строками RFC 3339 в UTC с обязательными +миллисекундами; дата слепка — строкой `YYYY-MM-DD`. `items` проверяется только +как JSON-массив. Каноническая форма и основания выбора зафиксированы в [исследовании формата](../../research/2026-08-16-order-snapshot-wire-format.md). Проверять все верхнеуровневые поля здесь уместно: их одиннадцать, и десять @@ -156,6 +156,23 @@ JSON-объектом с точным набором ключей: `order_id`, ` ## Что проверено +Переход STG → ODS снят на живом стенде 18 августа 2026 года при исполнении #94. +Опыты этого раздела прогнаны и подтвердили обещанное: разбиение сырья на годные +строки и брак полное и непересекающееся; две версии одного заказа легли в одну +партицию и на один шард, а `ods.order_v` вернуло позднюю независимо от фонового +слияния; повтор задачи с тем же `_load_id` строку в `ods.order_v` и её +`_load_ts` не изменил, а таблица ошибок записала тот же брак второй раз; строки +опытов убраны. Числа, механика опытов и поведение ClickHouse, на которое всё это +опирается, — [документ хранилища](../storage.md), «Что проверено». + +Одно расхождение с ожиданием осталось, и оно снаружи приёма: восемь настоящих +строк слепка получили класс `field_invalid` — все версии одного заказа с +отрицательным `total`. Класс заслужен, граница верна, дефект в генераторе и +заведён отдельным issue +[#102](https://git.dementev.space/ddmitry/clickstream-data-platform/issues/102). +Пока он не починен, критерий «честный прогон дня даёт пустой `_errors`» на +стартовом мире не выполняется. + Забор из Kafka в STG снят на живом стенде 18 августа 2026 года при исполнении #93: одно прямое чтение приносит весь слепок дня, метаданные доставки доступны, офсеты коммитятся, а сбой забора не двигает позицию мира. Числа — [ADR diff --git a/docs/architecture/storage.md b/docs/architecture/storage.md index 283fe47..4304e4a 100644 --- a/docs/architecture/storage.md +++ b/docs/architecture/storage.md @@ -9,10 +9,11 @@ keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/` лежит вся цепочка `Kafka → STG → ODS` — чтец топика `hits`, таблицы сырья, типизированное событие с таблицей ошибок, поверхность актуального состояния и три матвью. -Этап 3 добавил вход второго источника: топик `orders`, свой чтец и своё сырьё, -которое наполняет даг `orders_ingest`, а не матвью. Дальше по тексту устройство -описано так, как оно проектируется; построенное от заложенного отличает карта -таблиц в конце. +Этап 3 добавил вход второго источника и довёл его до ODS: топик `orders`, свой +чтец, своё сырьё, версии заказов с таблицей ошибок и поверхность текущего +состояния. Наполняет всю цепочку даг `orders_ingest` двумя шагами, а не матвью. +Дальше по тексту устройство описано так, как оно проектируется; построенное от +заложенного отличает карта таблиц в конце. Зона ответственности у документа одна — хранилище. Генератор описан отдельно: его замысел — в [спеке генератора](../specs/2026-08-01-generator.md), формат @@ -422,8 +423,8 @@ kafka_offset)`: смотрят такую таблицу от класса, а |---|---| | `00-databases.sql` | базы слоёв | | `10-stg-tables.sql` | чтецы топиков `hits` и `orders`, локальные и распределённые таблицы сырья обоих источников | -| `20-ods-tables.sql` | типизированное событие и таблица ошибок | -| `30-ods-views.sql` | актуальные события и матвью разбора в ODS | +| `20-ods-tables.sql` | типизированное событие, версии заказа и обе таблицы ошибок | +| `30-ods-views.sql` | актуальные события, текущие заказы и матвью разбора в ODS | | `40-stg-views.sql` | матвью приёма: чтец в сырьё | Порядок задают два правила. Первое: матвью принадлежит слою своей цели, а не @@ -482,6 +483,12 @@ ODS. Второе: матвью приёма создаётся последне | ODS | `ods.event_v` | актуальная версия события с полями источника | | ODS | `ods.event_errors_rep` / `_dist` | строки, не прошедшие строгий приём | | ODS | `ods.event_mv`, `ods.event_errors_mv` | разбор сырья в событие и в ошибки | +| ODS | `ods.order_snapshot_rep` / `_dist` | типизированные версии заказа | +| ODS | `ods.order_v` | текущая версия заказа на языке источника | +| ODS | `ods.order_snapshot_errors_rep` / `_dist` | строки слепка, не прошедшие строгий приём | + +Матвью разбора у заказов нет: срез сырья раскладывают по этим двум целям два +`INSERT SELECT` шага `parse_batch` в даге `orders_ingest`. Слои DDS и DM появляются на следующих этапах; их состав задан разделом 7 мастер-спеки и переносится сюда по мере постройки. @@ -502,6 +509,48 @@ MCP Context7 подтвердила обычное представление с `ods.event_v` остался равен `ods.event_dist FINAL`. Имена и типы всех колонок представления совпали с распределённой таблицей. +**Проверка версий заказов 18 августа 2026 года (#94).** MCP Context7 подтвердил, +что `JSONType` возвращает имя типа значения JSON, — на нём стоит проверка типов +в предикате приёма заказов. Остальное снято на закреплённом ClickHouse 26.3. + +Три находки касаются не заказов, а самого ClickHouse, и знать их стоит любому, +кто пишет здесь строгий разбор: + +- **`toDateOrNull` календарь не проверяет.** `'2026-02-30'` он молча превращает + в `2026-03-02`, `'2026-13-01'` — в `1970-01-01`. Для строгого приёма он + поэтому не годится: нужен `parseDateTimeInJodaSyntaxOrNull(…, 'yyyy-MM-dd', + 'UTC')`, который на обеих строках даёт `NULL`. +- **Маска разбора формы не держит.** `parseDateTime64InJodaSyntaxOrNull` по + маске `yyyy-MM-dd'T'HH:mm:ss.SSS'Z'` берёт и `2026-6-3T14:21:07.123Z` — без + ведущих нулей. Форму приходится сверять отдельно, регулярным выражением; + зато календарь маска проверяет честно (30 февраля, 13-й месяц, 25-й час + дают `NULL`), как и лишние или недостающие знаки долей секунды и смещение + `+00:00` вместо `Z`. +- **У временных типов есть потолок, и он ниже, чем кажется.** `DateTime64(3)` + заканчивается на `2299-12-31`, `DateTime` — на `2106-02-07`; строки за + потолком разбор отдаёт как `NULL`. Это ловит опыт, который берёт «заведомо + далёкий» год: 2999-й не разбирается вовсе. + +Про сам приём заказов снято следующее. Предикат формы провода разложил все +13 661 строку сырья, приехавшую при исполнении #93, без исключений и без +`NULL`: 13 653 годных и 8 брака, сумма сошлась с общим счётом по каждому +`_load_id`, пересечение ветвей — ноль. Все восемь строк брака — версии одного +заказа с отрицательным `total`; это дефект генератора, разобранный в +[спецификации приёма](orders/ingestion.md), «Что проверено». На управляемой +порции из пяти строк (две версии +одного заказа плюс по одной строке каждого класса брака) годные ушли в +`ods.order_snapshot`, брак — в `ods.order_snapshot_errors` с ожидаемыми +классами. Обе версии легли в одну партицию и на один шард; `FINAL` через +распределённую таблицу вернул одну строку — ту, у которой `updated_at` позже, — +тогда как физический счёт показывал две. Повтор задачи с тем же `_load_id` +строку в `ods.order_v` и её `_load_ts` не изменил, а таблица ошибок записала +тот же брак второй раз, как и обещано спецификацией приёма. + +Отдельно измерено, что **сводка вставки не годится в счётчики строк**: у +запроса с `WHERE` `read_rows` считает прочитанное с диска, а не подошедшее, — +на одном и том же срезе из пяти строк две вставки дали 4 и 5. Сколько чего +легло, спрашивают у самих таблиц по `_load_id`. + **Сверено с документацией.** Собственная колонка с именем виртуальной делает виртуальную недоступной. При вставке в `Distributed` шард выбирается по ключу шардирования; фоновый режим — умолчание, а `distributed_foreground_insert = 1` diff --git a/sql/ddl/20-ods-tables.sql b/sql/ddl/20-ods-tables.sql index b4854d3..825c4c9 100644 --- a/sql/ddl/20-ods-tables.sql +++ b/sql/ddl/20-ods-tables.sql @@ -150,3 +150,89 @@ SETTINGS ttl_only_drop_parts = 1; CREATE TABLE IF NOT EXISTS ods.event_errors_dist ON CLUSTER clickstream_cluster AS ods.event_errors_rep ENGINE = Distributed('clickstream_cluster', 'ods', 'event_errors_rep', cityHash64(raw)); + +-- Локальная таблица версий заказа. +-- +-- Слепок привозит состояние заказов окна изменяемости, а не поток изменений, +-- и одна и та же сущность приезжает в нём много дней подряд. Поэтому строка +-- здесь — версия заказа, а не запись слепка: ключ сущности order_id, колонка +-- версии updated_at (время последнего изменения строки в источнике). Порядок +-- версий решает источник, а не хранилище. +-- +-- Четыре координаты отвечают на разные вопросы, и путать их нельзя: +-- updated_at — какая бизнес-версия новее; snapshot_date — в слепке какого +-- модельного дня источник показал строку; _load_id — какой запуск Airflow её +-- принял; _load_ts — когда она приехала в хранилище. Ключом сущности не +-- становится ни одна из трёх последних: они про наблюдение и загрузку. +-- +-- PARTITION BY toDate(created_at) — по дню создания строки в источнике. Он +-- неизменен у всех версий заказа, поэтому версии лежат в одной партиции и +-- встречаются при мерже. Днём покупки этот день не является: бизнес-время +-- живёт в событии purchase (docs/research/2026-08-16-order-snapshot-wire-format.md). +-- +-- Замена партиции сюда не годится и заменена версиями: у прямого чтения Kafka +-- нет признака конца слепка, а дата наблюдения не ключ публикации (ADR 0010). +-- +-- items остаётся сырым фрагментом JSON: приём проверяет только, что это +-- массив. Что внутри позиций — забота DDS, а не границы провода. +CREATE TABLE IF NOT EXISTS ods.order_snapshot_rep ON CLUSTER clickstream_cluster +( + order_id String, + user_id UInt64, + status LowCardinality(String), + created_at DateTime64(3, 'UTC'), + updated_at DateTime64(3, 'UTC'), + items_total Decimal(18, 2), + discount Decimal(18, 2), + delivery Decimal(18, 2), + total Decimal(18, 2), + items String, + snapshot_date Date, + _load_id String, + _load_ts DateTime64(3, 'UTC') +) +ENGINE = ReplicatedReplacingMergeTree('/clickhouse/tables/{shard}/{database}/{table}', '{replica}', updated_at) +PARTITION BY toDate(created_at) +ORDER BY order_id; + +-- Лицо слоя: пишем и читаем через него. Ключ шардирования — cityHash64(order_id), +-- и выбор здесь не про перекос, а про корректность: только так все версии +-- одного заказа попадают на один шард, и FINAL через распределённую таблицу +-- выбирает одного победителя, а не по победителю на шард. +CREATE TABLE IF NOT EXISTS ods.order_snapshot_dist ON CLUSTER clickstream_cluster +AS ods.order_snapshot_rep +ENGINE = Distributed('clickstream_cluster', 'ods', 'order_snapshot_rep', cityHash64(order_id)); + +-- Локальная таблица брака слепка. +-- +-- Устроена как ods.event_errors_rep и по тем же доводам (см. выше): сырой +-- текст, метаданные доставки, класс брака, нарезка по дню загрузки, месяц +-- жизни, снятие кусками целиком. Своя колонка одна — _load_id: по нему видно, +-- какой запуск привёз брак, и повтор задачи узнаётся по совпадению _load_id +-- с координатами доставки. +-- +-- Классы у заказов свои и их три: not_an_object, keyset_mismatch, +-- field_invalid. Имя провалившегося поля в класс не входит — сырой текст лежит +-- рядом, и единичный случай разбирается по нему, без постоянной детализации +-- предиката (docs/architecture/orders/ingestion.md). +CREATE TABLE IF NOT EXISTS ods.order_snapshot_errors_rep ON CLUSTER clickstream_cluster +( + raw String, + error_class LowCardinality(String), + kafka_topic LowCardinality(String), + kafka_partition UInt64, + kafka_offset UInt64, + kafka_timestamp Nullable(DateTime64(3, 'UTC')), + consumer_host LowCardinality(String), + _load_id String, + _load_ts DateTime64(3, 'UTC') +) +ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/{database}/{table}', '{replica}') +PARTITION BY toDate(_load_ts) +ORDER BY (error_class, kafka_partition, kafka_offset) +TTL toDateTime(_load_ts) + INTERVAL 1 MONTH +SETTINGS ttl_only_drop_parts = 1; + +CREATE TABLE IF NOT EXISTS ods.order_snapshot_errors_dist ON CLUSTER clickstream_cluster +AS ods.order_snapshot_errors_rep +ENGINE = Distributed('clickstream_cluster', 'ods', 'order_snapshot_errors_rep', cityHash64(raw)); diff --git a/sql/ddl/30-ods-views.sql b/sql/ddl/30-ods-views.sql index f97087e..a5d7f29 100644 --- a/sql/ddl/30-ods-views.sql +++ b/sql/ddl/30-ods-views.sql @@ -8,6 +8,22 @@ AS SELECT * FROM ods.event_dist FINAL; +-- ODS: каноническое чтение текущего состояния заказов. +-- +-- То же, что у событий, но версию здесь задаёт источник: физическая таблица +-- хранит все приехавшие версии заказа, а FINAL оставляет ту, у которой +-- updated_at позже. Читать заказы голым SELECT по _dist значит зависеть от +-- того, сколько фоновых слияний успело пройти. +-- +-- Поля остаются языком источника, обогащения нет: DDS читает это +-- представление и отдельно решает, какой будет его модель (ADR 0010). +-- Способ выбора версии можно заменить, не трогая потребителей, — за именем +-- скрыт FINAL, а не устройство. +CREATE VIEW IF NOT EXISTS ods.order_v ON CLUSTER clickstream_cluster +AS +SELECT * +FROM ods.order_snapshot_dist FINAL; + -- ODS: разбор сырья в событие и в таблицу ошибок. -- -- Матвью две, и вместе они обязаны делить поток без зазора и без нахлёста: