From bdcbe48b595a0d561dd67befa21b91bac84cfe63 Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Tue, 18 Aug 2026 23:10:01 +0300 Subject: [PATCH 1/2] =?UTF-8?q?feat(orders):=20=D0=B4=D0=BE=D0=B1=D0=B0?= =?UTF-8?q?=D0=B2=D0=BB=D0=B5=D0=BD=20=D0=B2=D0=B5=D1=80=D1=81=D0=B8=D0=BE?= =?UTF-8?q?=D0=BD=D0=BD=D1=8B=D0=B9=20ODS=20=D0=B7=D0=B0=D0=BA=D0=B0=D0=B7?= =?UTF-8?q?=D0=BE=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: разбор сырья в событие и в таблицу ошибок. -- -- Матвью две, и вместе они обязаны делить поток без зазора и без нахлёста: -- 2.54.0 From a5508bd9ec6dc896f8e9f5e75d8463c6a773988c Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Tue, 18 Aug 2026 23:55:04 +0300 Subject: [PATCH 2/2] =?UTF-8?q?refactor(orders):=20SQL=20=D0=BF=D1=80?= =?UTF-8?q?=D0=B8=D1=91=D0=BC=D0=B0=20=D0=B2=D1=8B=D0=BD=D0=B5=D1=81=D0=B5?= =?UTF-8?q?=D0=BD=20=D0=B8=D0=B7=20DAG?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - преобразования хранилища должны читаться рядом с целевым слоем, а DAG должен показывать оркестрацию. - Что: - запросы забора и разбора заказов разложены по каталогам STG и ODS. - общий контракт провода подключён в обе ветви штатным шаблонизатором Airflow. - SQL смонтирован во все службы Airflow, а правило раскладки записано в архитектуре. - Проверка: - make lint config-test smoke check-services. - airflow tasks render для pull_batch и parse_batch; airflow tasks test для parse_batch. --- compose.yaml | 4 + dags/orders_ingest.py | 180 +++---------------------- docs/architecture/orders/ingestion.md | 10 +- docs/architecture/storage.md | 20 ++- sql/ods/_order_wire_contract.sql | 61 +++++++++ sql/ods/order_snapshot_errors_load.sql | 25 ++++ sql/ods/order_snapshot_load.sql | 33 +++++ sql/stg/orders_raw_load.sql | 27 ++++ 8 files changed, 195 insertions(+), 165 deletions(-) create mode 100644 sql/ods/_order_wire_contract.sql create mode 100644 sql/ods/order_snapshot_errors_load.sql create mode 100644 sql/ods/order_snapshot_load.sql create mode 100644 sql/stg/orders_raw_load.sql diff --git a/compose.yaml b/compose.yaml index 0e0687d..6e2944b 100644 --- a/compose.yaml +++ b/compose.yaml @@ -71,6 +71,9 @@ x-airflow-common: &airflow-common KAFKA_ORDERS_TOPIC: orders volumes: - ./dags:/opt/airflow/dags:ro + # Запросы дагов: даг называет файл, а текст читает Airflow при исполнении + # задачи. Каталог тот же, что применяет DDL, — вся SQL стенда живёт в sql/. + - ./sql:/opt/airflow/sql:ro - ./infra/airflow/init.sh:/opt/airflow/init.sh:ro - airflow_logs:/opt/airflow/logs - airflow_auth:/opt/airflow/auth @@ -391,6 +394,7 @@ services: volumes: - /var/run/docker.sock:/var/run/docker.sock - ./dags:/opt/airflow/dags:ro + - ./sql:/opt/airflow/sql:ro - ./infra/airflow/init.sh:/opt/airflow/init.sh:ro - airflow_logs:/opt/airflow/logs - airflow_auth:/opt/airflow/auth diff --git a/dags/orders_ingest.py b/dags/orders_ingest.py index 5f4324f..e655077 100644 --- a/dags/orders_ingest.py +++ b/dags/orders_ingest.py @@ -24,159 +24,9 @@ from airflow.sdk import Connection, dag, get_current_context, task START_DATE = datetime.datetime(2026, 1, 1, tzinfo=datetime.UTC) -# Забор — один прямой SELECT, без цикла до пустоты: одна порция ClickHouse -# берёт десятки тысяч сообщений, а слепок дня — порядка полутора тысяч строк. -# Короткая порция оставит хвост до следующего прогона, а отказ после чтения -# унесёт прочитанное с собой: офсеты коммитятся в момент чтения. Граница -# целиком — ADR 0008, «Следствия». -# -# stream_like_engine_allow_direct_select разрешает читать чтеца запросом; вторая -# половина пары объявлена на самой таблице (sql/ddl/10-stg-tables.sql). -# distributed_foreground_insert = 1 — конвенция ETL-вставок стенда: задача не -# должна зеленеть раньше, чем строки легли на шарды. -TAKE_ONE_BATCH = """ -INSERT INTO stg.orders_raw_dist -SELECT - raw, - _topic AS kafka_topic, - _partition AS kafka_partition, - _offset AS kafka_offset, - _timestamp_ms AS kafka_timestamp, - hostName() AS consumer_host, - {load_id:String} AS _load_id, - now64(3) AS _load_ts -FROM stg.orders_raw_kafka -SETTINGS - stream_like_engine_allow_direct_select = 1, - 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 -""" -) +# Запросы лежат файлами в sql/ репозитория, разложенные по слоям хранилища; +# сюда их монтирует compose — всем службам Airflow сразу. +SQL_ROOT = "/opt/airflow/sql" @dag( @@ -188,6 +38,7 @@ SETTINGS distributed_foreground_insert = 1 # дрались бы за неё, а слепок разъехался бы по двум _load_id; второй # прогон подождёт своей очереди. max_active_runs=1, + template_searchpath=SQL_ROOT, tags=["заказы"], ) def orders_ingest(): @@ -210,12 +61,16 @@ def orders_ingest(): send_receive_timeout=30, ) - @task - def pull_batch() -> None: + # Аргумент с расширением из templates_exts Airflow подменяет текстом файла: + # берёт его из template_searchpath и прогоняет через Jinja при исполнении + # задачи, а не при разборе дага. В задачу приезжает готовый запрос — вместе + # с тем, что файл подключил через include. + @task(templates_exts=(".sql",)) + def pull_batch(sql: str) -> None: load_id = get_current_context()["run_id"] client = clickhouse_client() try: - summary = client.command(TAKE_ONE_BATCH, parameters={"load_id": load_id}) + summary = client.command(sql, parameters={"load_id": load_id}) finally: client.close() # Размер порции — read_rows: written_rows у вставки в Distributed @@ -226,8 +81,8 @@ def orders_ingest(): load_id, ) - @task - def parse_batch() -> None: + @task(templates_exts=(".sql",)) + def parse_batch(good_rows_sql: str, bad_rows_sql: str) -> None: """Разобрать срез сырья в версии заказов и в брак. Обе вставки в одном task_id: транзакции между ними ClickHouse не даёт, @@ -237,15 +92,18 @@ def orders_ingest(): 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}) + client.command(good_rows_sql, parameters={"load_id": load_id}) + client.command(bad_rows_sql, parameters={"load_id": load_id}) finally: client.close() # Счётчиков строк нет: у запроса с WHERE read_rows считает прочитанное # с диска, а не подошедшее (storage.md, «Что проверено»). logging.info("срез разобран: _load_id %s", load_id) - pull_batch() >> parse_batch() + pull_batch("stg/orders_raw_load.sql") >> parse_batch( + "ods/order_snapshot_load.sql", + "ods/order_snapshot_errors_load.sql", + ) orders_ingest() diff --git a/docs/architecture/orders/ingestion.md b/docs/architecture/orders/ingestion.md index 7ccbfbd..7c3521e 100644 --- a/docs/architecture/orders/ingestion.md +++ b/docs/architecture/orders/ingestion.md @@ -38,19 +38,25 @@ дёргает его тот, кто положил слепок в топик, — работник пульта мира, — и ждёт конца прогона. Все строки получают `_load_id`, равный `run_id` Airflow. `_load_ts` вычисляется при этой записи и дальше переносится без пересчёта. +Запрос забора лежит в [`sql/stg/orders_raw_load.sql`](../../../sql/stg/orders_raw_load.sql): +даг задаёт порядок и параметры, а преобразование остаётся в SQL своего слоя. Один следующий `task_id` отвечает за весь переход STG → ODS. Внутри него два последовательных `INSERT SELECT` читают неизменный срез по `_load_id`: первый пишет годные строки в `ods.order_snapshot`, второй — брак в `ods.order_snapshot_errors`. Транзакции между запросами нет. При частичном сбое Airflow повторяет весь `task_id`; одинаковые исходные строки и служебные метки -не вычисляются заново. +не вычисляются заново. Запросы лежат рядом с целями: +[`order_snapshot_load.sql`](../../../sql/ods/order_snapshot_load.sql) и +[`order_snapshot_errors_load.sql`](../../../sql/ods/order_snapshot_errors_load.sql). Условия запросов взаимодополняющие: один общий предикат определяет брак, а годная ветвь использует его буквальное отрицание. Все функции предиката возвращают результат без исключения, а сам предикат всегда заканчивается в `true` или `false`, не в `NULL`. Постоянный классификатор между STG и ODS для -этого не нужен. +этого не нужен. Обе ветви включают один файл +[`_order_wire_contract.sql`](../../../sql/ods/_order_wire_contract.sql), поэтому +предикат нельзя случайно исправить только в одной из них. ## Граница строгого приёма diff --git a/docs/architecture/storage.md b/docs/architecture/storage.md index 4304e4a..7523bb2 100644 --- a/docs/architecture/storage.md +++ b/docs/architecture/storage.md @@ -414,7 +414,22 @@ kafka_offset)`: смотрят такую таблицу от класса, а разрастается до имени отдельного поля. Точная граница приёма — в [спецификации заказов](orders/ingestion.md). -## Раскладка DDL +## Раскладка SQL + +Исполняемые дагами запросы лежат по правилу +`sql/<слой>/<объект>_<роль>.sql`. Слой — слой цели запроса: забор заказов в +сырьё живёт в `sql/stg/orders_raw_load.sql`, две ветви разбора — в `sql/ods/`. +Даг называет файлы, передаёт параметры и задаёт порядок выполнения, но не +хранит текст преобразований. Общий фрагмент начинается с подчёркивания: файл +`sql/ods/_order_wire_contract.sql` сам не исполняется, его включают обе ветви +разбора. + +Airflow читает и собирает эти файлы штатным шаблонизатором при исполнении +задачи. Каталог `sql/` смонтирован во все его службы только для чтения. Поэтому +обработчик дагов не зависит от наличия SQL на машине при разборе Python, а +планировщик видит те же файлы при исполнении задачи. + +### DDL Файлы лежат в `sql/ddl/` и применяются по порядку имён. Сначала все статичные объекты, потом матвью — тогда к моменту создания матвью её цель уже существует. @@ -488,7 +503,8 @@ ODS. Второе: матвью приёма создаётся последне | ODS | `ods.order_snapshot_errors_rep` / `_dist` | строки слепка, не прошедшие строгий приём | Матвью разбора у заказов нет: срез сырья раскладывают по этим двум целям два -`INSERT SELECT` шага `parse_batch` в даге `orders_ingest`. +`INSERT SELECT` шага `parse_batch` из файлов `sql/ods/order_snapshot_load.sql` +и `sql/ods/order_snapshot_errors_load.sql`. Слои DDS и DM появляются на следующих этапах; их состав задан разделом 7 мастер-спеки и переносится сюда по мере постройки. diff --git a/sql/ods/_order_wire_contract.sql b/sql/ods/_order_wire_contract.sql new file mode 100644 index 0000000..1fafb6f --- /dev/null +++ b/sql/ods/_order_wire_contract.sql @@ -0,0 +1,61 @@ +-- Форма провода одной строки слепка: общая часть обеих вставок ODS. +-- +-- Строгий приём: проверяется форма провода и ничего сверх неё. +-- +-- Ключей одиннадцать, и проверяются все: десять скалярных значений прямо +-- образуют типизированную строку заказа. Внутрь items приём не смотрит — +-- содержимое позиций, переходы статуса и равенства сумм остаются ниже границы +-- (docs/architecture/orders/ingestion.md, «Граница строгого приёма»). +-- +-- Файл включают обе вставки нарочно: годная ветвь берёт row_is_valid, брак — +-- буквальное NOT. Разойдись условия хоть на символ — строка либо задвоится, +-- либо исчезнет молча. +-- +-- Отсюда два запрета на выражения предиката, и оба серьёзные. Первый: ни одно +-- не возвращает NULL — трёхзначная логика дала бы строку, которую не берёт ни +-- условие, ни его отрицание. Поэтому сравнения дают 0 или 1, а обнуляемый +-- разбор заканчивается IS NOT NULL. Второй: ни одно не бросает исключений на +-- произвольном raw — иначе одна грязная строка роняет весь переход, ради +-- отсутствия чего таблица брака и заведена. + +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 diff --git a/sql/ods/order_snapshot_errors_load.sql b/sql/ods/order_snapshot_errors_load.sql new file mode 100644 index 0000000..597e9ab --- /dev/null +++ b/sql/ods/order_snapshot_errors_load.sql @@ -0,0 +1,25 @@ +-- Брак: тот же срез и буквальное отрицание того же предиката. +-- +-- Классы перекрываются, поэтому проверяются по порядку, а в error_class идёт +-- первый совпавший: скаляр проваливает и проверку на объект, и сверку ключей — +-- без объявленного порядка он попал бы то в один класс, то в другой. + +INSERT INTO ods.order_snapshot_errors_dist +{% include "ods/_order_wire_contract.sql" %} +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 diff --git a/sql/ods/order_snapshot_load.sql b/sql/ods/order_snapshot_load.sql new file mode 100644 index 0000000..5c5f8a9 --- /dev/null +++ b/sql/ods/order_snapshot_load.sql @@ -0,0 +1,33 @@ +-- Годные версии заказа. Срез читается по _load_id — тому же, что проставил +-- забор: разбор идёт по неизменной порции, а не по «всему, что появилось». +-- +-- Разобранные значения берёт assumeNotNull: обнуляемый разбор стоит за +-- предикатом, который NULL уже отсёк, и приведение здесь не может упасть. + +INSERT INTO ods.order_snapshot_dist +{% include "ods/_order_wire_contract.sql" %} +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 diff --git a/sql/stg/orders_raw_load.sql b/sql/stg/orders_raw_load.sql new file mode 100644 index 0000000..1a458d4 --- /dev/null +++ b/sql/stg/orders_raw_load.sql @@ -0,0 +1,27 @@ +-- Забор порции слепка заказов из Kafka в сырьё STG. +-- +-- Забор — один прямой SELECT, без цикла до пустоты: одна порция ClickHouse +-- берёт десятки тысяч сообщений, а слепок дня — порядка полутора тысяч строк. +-- Короткая порция оставит хвост до следующего прогона, а отказ после чтения +-- унесёт прочитанное с собой: офсеты коммитятся в момент чтения. Граница +-- целиком — ADR 0008, «Следствия». +-- +-- stream_like_engine_allow_direct_select разрешает читать чтеца запросом; вторая +-- половина пары объявлена на самой таблице (sql/ddl/10-stg-tables.sql). +-- distributed_foreground_insert = 1 — конвенция ETL-вставок стенда: задача не +-- должна зеленеть раньше, чем строки легли на шарды. + +INSERT INTO stg.orders_raw_dist +SELECT + raw, + _topic AS kafka_topic, + _partition AS kafka_partition, + _offset AS kafka_offset, + _timestamp_ms AS kafka_timestamp, + hostName() AS consumer_host, + {load_id:String} AS _load_id, + now64(3) AS _load_ts +FROM stg.orders_raw_kafka +SETTINGS + stream_like_engine_allow_direct_select = 1, + distributed_foreground_insert = 1 -- 2.54.0