From 6910440400dafdbb379c6814e1223ed5ea0cef40 Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Fri, 7 Aug 2026 16:06:46 +0300 Subject: [PATCH] =?UTF-8?q?feat(ods):=20=D1=82=D0=B8=D0=BF=D0=B8=D0=B7?= =?UTF-8?q?=D0=B8=D1=80=D0=BE=D0=B2=D0=B0=D0=BD=D0=BD=D0=BE=D0=B5=20=D1=81?= =?UTF-8?q?=D0=BE=D0=B1=D1=8B=D1=82=D0=B8=D0=B5,=20=D1=81=D1=82=D1=80?= =?UTF-8?q?=D0=BE=D0=B3=D0=B8=D0=B9=20=D0=BF=D1=80=D0=B8=D1=91=D0=BC=20?= =?UTF-8?q?=D0=B8=20=D1=82=D0=B0=D0=B1=D0=BB=D0=B8=D1=86=D0=B0=20=D0=BE?= =?UTF-8?q?=D1=88=D0=B8=D0=B1=D0=BE=D0=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Зачем: цепочка Kafka → STG → ODS достраивается последним этажом. Сырьё уже доезжает (#37), настоящие события в топике есть (#41), а типизированного слоя не было — событие негде было прочитать колонками, а брак негде увидеть. Что: - sql/ddl/20-ods-tables.sql — ods.event_rep/_dist на ReplacingMergeTree с версией _load_ts, партиция по EventDate, ключ по разделу 1.3 спеки, шардирование cityHash64(ClientID); ods.event_errors_rep/_dist с классом брака, своими ключами и сроком жизни в месяц. - sql/ddl/30-ods-views.sql — две матвью над stg.hits_raw_dist. Годность считает предикат из трёх частей, вторая матвью берёт его дословное отрицание, класс брака пишется первым совпавшим из трёх. - Метку времени разбирает parseDateTimeBestEffortOrNull, а не JSONExtract: ISO-8601 с суффиксом Z JSONExtract не берёт вовсе. Спека генератора обещала обратное — обещание поправлено, форма на проводе не менялась. - Сверка объявлений (contract-тест) снята из документов и из докстрингов schema.py: сверх строгого приёма она ловила только смену типа. - Документация приведена в соответствие: ADR 0005, дока хранилища и обе спеки; группа «сказано по памяти» в доке хранилища опустела. Проверка: make up && make check-clickhouse (8 проверок, 7,5 с); make lint, make typecheck, make test (406), make docs без диффа. Разовые опыты при исполнении — в теле PR. Co-Authored-By: Claude Opus 5 --- docs/adr/0005-event-ingestion.md | 72 +++++-- docs/architecture/storage.md | 53 +++-- docs/architecture/testing.md | 8 +- docs/specs/2026-07-30-stand-v2-realism.md | 19 +- docs/specs/2026-08-01-generator.md | 36 ++-- generator/src/clickstream_generator/schema.py | 9 +- sql/ddl/20-ods-tables.sql | 152 +++++++++++++ sql/ddl/30-ods-views.sql | 199 ++++++++++++++++++ 8 files changed, 475 insertions(+), 73 deletions(-) create mode 100644 sql/ddl/20-ods-tables.sql create mode 100644 sql/ddl/30-ods-views.sql diff --git a/docs/adr/0005-event-ingestion.md b/docs/adr/0005-event-ingestion.md index 00ada90..5b32f36 100644 --- a/docs/adr/0005-event-ingestion.md +++ b/docs/adr/0005-event-ingestion.md @@ -125,15 +125,20 @@ STG сырьё исключительно от брака: слой сырых присутствие, и потому сужен до пяти колонок. Цена решения. Контракт получает ещё два места: типы сорока семи колонок в -выражениях матвью и список тех же имён для сверки ключей. Раздел 1.4 спеки -предупреждает, что, повторяясь примерно в семи местах, они расходятся молча, а -contract-тест из #43 сюда не дотягивается — он сравнивает `system.columns` -целевой таблицы со схемой генератора и о выражениях матвью ничего не знает. -Смягчение работает не везде: опечатка в имени скалярного поля уводит строки в -таблицу ошибок пачкой и видна сразу, а опечатка в имени массива даёт пустой -массив тихо — сверка ключей проверяет ключи сообщения, а не выражения матвью. -Эти двенадцать колонок сторожит smoke: известное событие с товарами обязано -доезжать с непустыми массивами. +выражениях матвью и список тех же имён для сверки ключей — а список этот +повторён дважды, по разу на матвью. Раздел 1.4 спеки предупреждает, что, +повторяясь примерно в семи местах, они расходятся молча. + +И сверка ключей эту цену покрывает не всю: она смотрит на ключи сообщения, а +не на выражения матвью. Опечатка в имени внутри `JSONExtract` даёт умолчание +типа — ноль, пустую строку, пустой массив, — и молчит она у сорока двух +обычных колонок ровно так же, как у двенадцати массивов. Громко ломаются +только пять опорных: у них разбор `Nullable` стоит в предикате, и опечатка +уводит в таблицу ошибок все строки до единой. Остальные сорок две сторожит +сверка разобранного события против сырого текста — разовый опыт при +исполнении #43, а не постоянная проверка; сила его в том, что выражения +сверки собираются из контракта, а выражения матвью написаны руками по +описанию выгрузки, и одна опечатка в двух местах не повторяется. Второе — разбор функциями дороже разбора форматом. На объёмах стенда это несущественно; если станет заметно, сорок семь вызовов сворачиваются в один @@ -177,15 +182,42 @@ contract-тест из #43 сюда не дотягивается — он ср Тогда же нашлась и граница обещания — запись с пустым значением и запись-надгробие не дают строки вовсе; из-за неё в «Решении» и в «Почему» приписано слово «непустое». Замер целиком — в [доке -хранилища](../architecture/storage.md), раздел «Что проверено». Остальные три -по-прежнему ждут живого стенда; это однострочные `SELECT`, их довольно прогнать -заодно: +хранилища](../architecture/storage.md), раздел «Что проверено». -- форма именованного кортежа в `JSONExtract` с `Nullable`-членами — нужна для - свёртки сорока семи вызовов в один, если разбор окажется дорогим; -- `isValidJSON('123')` возвращает единицу, а `JSONExtractKeys` от скаляра — - пустой массив. На обоих стоят классы брака и их приоритет, а документация - поведение на не-объекте не описывает: два `SELECT` закрывают вопрос; -- `JSONAsString` действительно падает на некорректном JSON, а не пропускает - строку. На этом стоит отказ от него в пользу `RawBLOB`; документация про - ошибочный ввод молчит. +Остальные четыре закрыты при исполнении #43, на стенде 7 августа 2026 года, +ClickHouse 26.3.17.56. Все четыре ответили так, как ждала постановка: + +- `isValidJSON('123')` возвращает единицу — скаляр законный JSON, и первый + класс брака поэтому проверяет именно объект, а не валидность; +- `JSONExtractKeys('123')` возвращает пустой массив, и он же приходит от + вовсе не-JSON. Значит скаляр проваливает и сверку ключей — отсюда + обязательный порядок классов; +- `JSONAsString` на некорректном вводе падает, а не пропускает строку: код 117 + `INCORRECT_DATA`, «JSON object must begin with '{'». Падает и на скаляре + `123`. На этом стоит отказ от него в пользу `RawBLOB`; +- форма именованного кортежа с `Nullable`-членами работает и годится в + запасной вариант: `JSONExtract(raw, 'Tuple(WatchID Nullable(UInt64), …)')` + даёт NULL в тех членах, что не разобрались, обращение по имени члена + доступно, а на не-объекте кортеж выходит целиком из NULL и исключения нет. + +Тем же заходом нашлось то, о чём никто не спрашивал, и оно оказалось +блокирующим. **`JSONExtract` с типом `DateTime` не разбирает ISO-8601 с +суффиксом зоны.** На проводе `UTCEventTime` уезжает как +`2026-06-01T12:34:56Z` (спека генератора, раздел 4), а +`JSONExtract(raw, 'UTCEventTime', 'Nullable(DateTime)')` отдаёт на такой +строке NULL — то есть все события до единого уходили бы в брак с классом +`key_field_unparsed`. Ни `DateTime64`, ни `DateTime('UTC')` суффикс тоже не +берут; без `Z` та же строка разбирается. Спека генератора обещала обратное +(«принимает ISO без плясок») — обещание было ошибочным и исправлено тем же +PR. Разбор метки времени поэтому идёт +`parseDateTimeBestEffortOrNull(JSONExtractString(raw, 'UTCEventTime'))`: +документация ClickHouse прямо относит ISO-8601 к форматам +`parseDateTimeBestEffort` (сверено через Context7 7 августа 2026 года), а +вариант `*OrNull` возвращает NULL вместо исключения и потому годится в +предикат. `EventDate` уезжает как `2026-06-01` и разбирается `JSONExtract` +без оговорок. + +Оговорка к обнуляемому разбору даты, измеренная там же: `Nullable(Date)` +даёт NULL на строке, которая датой не является вовсе («мусор»), и на числе, +но невозможную дату `2026-13-99` молча приводит к `1970-01-01`. То есть +класс `key_field_unparsed` ловит порчу типа, а не порчу значения внутри типа. diff --git a/docs/architecture/storage.md b/docs/architecture/storage.md index 929ebd7..3ddb734 100644 --- a/docs/architecture/storage.md +++ b/docs/architecture/storage.md @@ -6,10 +6,11 @@ раздел «Что проверено» — чему в этом тексте верить и на каком основании. **Что здесь описано и чего ещё нет.** Собран этап 1: кластер из двух шардов, -keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/` уже лежат базы слоёв -и объекты STG — чтец топика `hits`, таблицы сырья и матвью приёма. Объектов ODS -в репозитории пока нет. Дальше по тексту устройство описано так, как оно -проектируется; построенное от заложенного отличает карта таблиц в конце. +keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/` лежит вся цепочка +`Kafka → STG → ODS` — чтец топика `hits`, таблицы сырья, типизированное +событие с таблицей ошибок и три матвью. Дальше по тексту устройство описано +так, как оно проектируется; построенное от заложенного отличает карта таблиц в +конце. Зона ответственности у документа одна — хранилище. Генератор описан отдельно: его замысел — в [спеке генератора](../specs/2026-08-01-generator.md), формат @@ -229,9 +230,12 @@ NULL, — иначе трёхзначная логика даст строку, разъедется само: сырьё истечёт раньше, а перезаливка модельного дня задвоит его, тогда как в ODS тот же повтор схлопнется. Правило, красное в норме, учит не смотреть на оповещения -([ADR 0002](../adr/0002-monitoring-scope.md)). Равенство проверяется разово в -smoke на управляемой пачке: отправили N сообщений — получили N строк сырья и N в -сумме событий и ошибок. Счёт по ODS идёт через `FINAL`: голый `count()` по +([ADR 0002](../adr/0002-monitoring-scope.md)). Равенство проверено разовым +опытом при исполнении #43, на управляемой пачке: отправили N сообщений — +получили N строк сырья и N в сумме событий и ошибок. Постоянной целью такой +опыт не становится, и почему — в [карте +проверок](testing.md), раздел «Интеграционная проверка постоянной целью не +становится». Счёт по ODS идёт через `FINAL`: голый `count()` по `ReplacingMergeTree` зависит от того, сколько мержей успело пройти, и спека это прямо запрещает (раздел 6). @@ -345,8 +349,7 @@ ODS. Второе: матвью приёма создаётся последне ## Карта таблиц -Ниже — то, что закладывает этап 2. DDL слоя STG уже лежит в `sql/ddl/`; -объектов ODS в репозитории пока нет. +Ниже — то, что закладывает этап 2; всё перечисленное лежит в `sql/ddl/`. | Слой | Объект | Что это | |---|---|---| @@ -382,9 +385,24 @@ ODS. Второе: матвью приёма создаётся последне завязан на движок базы `Atomic`. `ON CLUSTER` ждёт все хосты и бросает по таймауту; `CREATE ... IF NOT EXISTS` на существующем объекте не бросает. -**Проверено на стенде.** Опыты прогнаны на живом кластере при исполнении #37: -четыре — 5 августа 2026 года, пятый — 6 августа. Все подтвердили то, что здесь -написано. +**Проверено на стенде.** Опыты прогнаны на живом кластере: пять при исполнении +#37 (четыре 5 августа 2026 года, пятый 6 августа) и два при исполнении #43 +(7 августа). Все подтвердили то, что здесь написано. + +- Матвью с источником-`Distributed` срабатывает на вставку именно в эту + распределённую таблицу, до раскладки по шардам. Обе матвью разбора стоят над + `stg.hits_raw_dist`, а пишет в неё матвью приёма — и события доезжают до + `ods.event`; значит блок она видит. В документации ClickHouse случая нет + вовсе, до 7 августа утверждение держалось на опыте владельца. +- Упавшая матвью роняет вставку и останавливает потребление до починки. + Проверено сносом цели — распределённой `ods.event_dist` — при живом чтеце: за + двадцать секунд (сброс блока идёт за 7,5) в сырьё не приехало ничего, а + офсет группы застыл с отставанием в одно сообщение. Цель вернули — сообщение + доехало само, без повторной отправки, отставание ушло в ноль, событие + разобралось. Сносить надо именно распределённую таблицу: вставка в + `Distributed` кладёт блок в спул и сразу возвращает управление, так что на + сносе локальной ошибка всплыла бы фоном и утверждение показалось бы + опровергнутым. - `RawBLOB` даёт ровно одну строку на каждое непустое сообщение. Три сообщения с ключами, поставленные в очередь до одного сброса продюсера, стали тремя @@ -433,12 +451,5 @@ Kafka с пустым значением (ноль байт) и запись-н таблице. Это вычитано, а не измерено. На устройство приёма оговорка не влияет: служебные колонки мы заполняем выражением при любом ответе. -**Сказано по памяти, проверки пока нет.** Осталось два утверждения, и оба ждут -одного и того же — матвью разбора, а она приходит с #43. - -- Матвью с источником-`Distributed` срабатывает на вставку именно в эту - распределённую таблицу, до раскладки по шардам. В документации случая нет - вовсе, утверждение держится на опыте владельца. -- Упавшая матвью роняет вставку и останавливает потребление до починки. На этой - фразе стоит правило «грязные записи не валят пайплайн», а сама она стоит пока - на одном рассуждении. +**Сказано по памяти, проверки нет.** Группа пуста: оба утверждения, ждавшие +матвью разбора, закрыты опытами при исполнении #43 и переехали выше. diff --git a/docs/architecture/testing.md b/docs/architecture/testing.md index 4f09fd7..6caa598 100644 --- a/docs/architecture/testing.md +++ b/docs/architecture/testing.md @@ -22,10 +22,10 @@ макросов `shard`: обе ноды здоровы и порты отвечают. `make check-clickhouse` не заметит потерянного подключения Superset: он про ClickHouse и только. -Отсюда правило для новой проверки: **спроси, кого она спрашивает.** Договор со -схемой событий — вопрос к ClickHouse, значит дом ему в `check-clickhouse`, даже -если по цене он подошёл бы смоуку. Счётчики против манифеста — тоже вопрос к -ClickHouse: строки в `ods.event` считает сам сервер и отвечает сразу. +Отсюда правило для новой проверки: **спроси, кого она спрашивает.** Счётчики +против манифеста — вопрос к ClickHouse: строки в `ods.event` считает сам +сервер и отвечает сразу, значит дом им в `check-clickhouse`, даже если по цене +они подошли бы смоуку. ## Карта целей diff --git a/docs/specs/2026-07-30-stand-v2-realism.md b/docs/specs/2026-07-30-stand-v2-realism.md index 48c3413..0aa9420 100644 --- a/docs/specs/2026-07-30-stand-v2-realism.md +++ b/docs/specs/2026-07-30-stand-v2-realism.md @@ -172,8 +172,9 @@ Ecommerce (заполнены только у торговых событий): рендеренное «описание выгрузки» в доках — аналог документации Метрики. Сторона хранилища (DDL, SELECT матвью, трансформации, витрины) пишется по этой документации на своих этапах, как в бою хранилище адаптируется к -источнику; границу сторожат строгий приём (раздел 6) и contract-тест в -smoke — сравнение `system.columns` поднятого стенда со схемой генератора. +источнику; границу сторожит строгий приём (раздел 6). Вторым сторожем здесь +стояла сверка объявлений — `system.columns` поднятого стенда против схемы +генератора; она снята при исполнении #43 как ничего не добавляющая к соседу. Без контракта 47 колонок, повторяясь примерно в семи местах, расходятся молча. Заодно это учебный артефакт: менти видит на живом примере, что такое data contract. @@ -584,8 +585,10 @@ v2 стартует пустым, поэтому объём ниже — это Список убывает по мере постройки: проверенное уходит отсюда, а ответ с датой остаётся там, где на него опираются. Формат чтеца и форма виртуальной метки -времени закрыты при исполнении #37 — см. [доку -хранилища](../architecture/storage.md), раздел «Что проверено». +времени закрыты при исполнении #37, форма ключа ODS, поведение матвью над +`Distributed` и запасной именованный кортеж — при исполнении #43; ответы — в +[доке хранилища](../architecture/storage.md) и +[ADR 0005](../adr/0005-event-ingestion.md), разделы «Что проверено». - Поведение соединения двух Distributed-таблиц и `distributed_product_mode` — эмпирически на стенде (хвост #14). @@ -593,14 +596,6 @@ v2 стартует пустым, поэтому объём ниже — это прогонами, отсутствие дублей при штатной работе. Закрыто пока наполовину: что обе ноды читают топик и обе партиции доезжают, показал #37; что дублей нет и как раскладка меняется между прогонами — нет. -- Точная форма `ORDER BY` ODS-таблиц (выражение `intHash32` в ключе - ReplacingMergeTree). -- Матвью с источником-`Distributed` срабатывает на вставку именно в эту - распределённую таблицу, до раскладки по шардам: на этом стоит цепочка - STG → ODS (ADR 0005). Проверено владельцем на рабочих проектах, в документации - ClickHouse этот случай не описан. -- Форма именованного кортежа в `JSONExtract` с `Nullable`-членами — ею - сворачиваются 47 вызовов в один, если разбор окажется дорогим (ADR 0005). - Размер артефакта эталонного мира после пересборки. - Спорные API (Airflow Datasets/сенсоры, ClickHouse DDL) — перед кодом сверять через MCP Context7 (правило AGENTS.md). diff --git a/docs/specs/2026-08-01-generator.md b/docs/specs/2026-08-01-generator.md index c33cb21..bbfcdf7 100644 --- a/docs/specs/2026-08-01-generator.md +++ b/docs/specs/2026-08-01-generator.md @@ -33,8 +33,8 @@ - **Детерминизм до байта.** Одно зерно — побайтово тот же снимок; сверка — хешами манифеста. Транспорт (офсеты Kafka, темп) — вне обещания. - **Схема — контракт генератора.** Python-модуль с чистыми данными; - хранилище строится по рендеренной документации, границу сторожит - contract-тест. + хранилище строится по рендеренной документации, границу сторожит строгий + приём на стороне хранилища. - **Один сериализатор, глупые приёмники.** День-функция выдаёт канонические байты; приёмники — файл, Kafka пачкой, Kafka с темпом. - **Числа.** Средний день ~50 тыс. событий; эталонный снимок — 14 дней; @@ -229,10 +229,13 @@ - **Сторона хранилища пишется по документации, не генерируется.** DDL `ods.event`, SELECT матвью, `dds.event_v`, трансформации — работа следующих этапов по «описанию выгрузки», как в бою хранилище адаптируется - к источнику. Границу сторожат два боевых механизма: строгий приём + к источнику. Границу сторожит боевой механизм — строгий приём (`Nullable`-разбор со сверкой набора ключей, таблицы `*_errors` — раздел 6 - мастер-спеки) и contract-тест в smoke — сравнение `system.columns` - поднятого стенда со схемой генератора. + мастер-спеки). Сверка объявлений (`system.columns` поднятого стенда против + контракта) здесь стояла вторым механизмом и снята при исполнении #43: + сверх строгого приёма она ловила ровно одно — смену типа колонки, — а её + ловит и сверка разобранного события, причём на живых данных, а не на + объявлениях. Отклонено с доводами: @@ -263,10 +266,18 @@ - **Даты и время на проводе — ISO-8601.** `EventDate` уезжает как `2026-06-01`, `UTCEventTime` — как `2026-06-01T12:34:56Z`. Довод — читаемость сырья: весь смысл слоя STG в том, что менти открывает колонку `raw` в обычном клиенте и - разбирает событие глазами, а число эпохи этот урок убивает. Разбору это - ничего не стоит: `JSONExtract(raw, 'UTCEventTime', 'Nullable(DateTime)')` - принимает ISO без плясок. Колонка `ecommerce` — строка, внутри которой лежит - экранированный JSON, как отдаёт Метрика. + разбирает событие глазами, а число эпохи этот урок убивает. Колонка + `ecommerce` — строка, внутри которой лежит экранированный JSON, как отдаёт + Метрика. + + Оговорка про цену разбора, вписанная сюда 6 августа и оказавшаяся неверной: + здесь стояло, что `JSONExtract(raw, 'UTCEventTime', 'Nullable(DateTime)')` + «принимает ISO без плясок». Не принимает — на строке с суффиксом `Z` он + отдаёт NULL, и при исполнении #43 это увело бы в брак все события до + единого. Измерено на стенде 7 августа 2026 года; форма на проводе от этого + не меняется, меняется выражение разбора на стороне хранилища — + `parseDateTimeBestEffortOrNull` вместо `JSONExtract` + ([ADR 0005](../adr/0005-event-ingestion.md)). Форму реализует сериализатор (#41), хранилище (#43) читает то, что он положил: порядок тикетов развёрнут 6 августа 2026 года, и отправитель идёт @@ -283,7 +294,8 @@ ReplacingMergeTree. - **Эталонный снимок при старте стенда — через Kafka, пакетным режимом проигрывателя.** Отдельный механизм заливки не строится: каждый `make up` - бесплатно прогоняет весь конвейер и contract-тест на настоящих данных. + бесплатно прогоняет весь конвейер на настоящих данных, и строгий приём + хранилища проверяет контракт тем же прогоном. Оговорка «если заливка уйдёт в десятки минут — вернуться к прямой загрузке» проверена при фиксации чисел: 14 × 50 тыс. ≈ 700 тыс. событий — расчётно минута-две, запас есть. @@ -395,8 +407,8 @@ pytest-тест с маркером `perf` и таймаутом-обрубан Внесены в мастер-спеку тем же коммитом, что и эта спека: - **Раздел 1.4**: «из контракта выводятся DDL и валидация» заменено на data - contract — хранилище пишется по документации, границу сторожит - contract-тест (раздел 3 здесь). + contract — хранилище пишется по документации, границу сторожит строгий + приём (раздел 3 здесь). - **Раздел 8**: артефакт `data/startup_history/` в git заменён манифестом; снимок генерируется на месте (раздел 5 здесь). Туман «политика версионирования артефакта» закрыт этим же ходом: версионируется манифест. diff --git a/generator/src/clickstream_generator/schema.py b/generator/src/clickstream_generator/schema.py index d035bb0..db942da 100644 --- a/generator/src/clickstream_generator/schema.py +++ b/generator/src/clickstream_generator/schema.py @@ -4,9 +4,9 @@ 1.1–1.2) и здесь не переоткрываются — модуль записывает их машинно-читаемо. Контракт принадлежит генератору и кормит трёх потребителей: сам генератор, его валидацию и «описание выгрузки» в доках (`schema_doc`). Хранилище -строится по описанию, а не по модулю; границу будет сторожить contract-тест, -сверяющий `system.columns` поднятого стенда с этим контрактом, — он придёт -вместе с типизированным ODS (спека генератора, раздел 3). +строится по описанию, а не по модулю; границу сторожит строгий приём на его +стороне — сверка набора ключей сообщения с контрактным списком, и +разошедшееся уходит в `ods.event_errors` (спека генератора, раздел 3). Что несёт описатель колонки: @@ -16,7 +16,8 @@ `system.columns`: параметры входят в имя типа целиком, без сокращений (`LowCardinality(String)`, `Array(Float64)`). Сверено 2026-08-01 — по документации ClickHouse через Context7 и запросом к узлу стенда - (26.3.17.56); от этой записи зависит будущий contract-тест. + (26.3.17.56). Запись важна потому, что по ней человек пишет DDL: тип, + сокращённый здесь, приедет в таблицу сокращённым же. - `numpy_dtype` — чем колонка представлена внутри генератора; у массивов это тип элемента. Строки живут в `object`-массивах: numpy-строки фиксированной длины стенду ничего не дают. diff --git a/sql/ddl/20-ods-tables.sql b/sql/ddl/20-ods-tables.sql new file mode 100644 index 0000000..7c2cb40 --- /dev/null +++ b/sql/ddl/20-ods-tables.sql @@ -0,0 +1,152 @@ +-- ODS: типизированное событие и таблица ошибок разбора. +-- +-- Состав, имена и типы колонок списаны с «описания выгрузки» +-- (docs/formats/clickstream-event.md) — так же, как в бою хранилище пишут по +-- документации источника. Модуль контракта генератора здесь не читается: +-- граница «трекер | хранилище» проходит по документу. +-- Механика приёма и три класса брака — ADR 0005, конвенции имён, служебных +-- колонок и сроков — docs/architecture/storage.md. + +-- Локальная таблица события. +-- +-- Движок ReplacingMergeTree, колонка версии — _load_ts. Работа у версии ровно +-- одна: схлопнуть повтор доставки. Метка не ставится здесь заново, а +-- переносится из stg.hits_raw как есть — колонка отвечает на вопрос «когда +-- строка приехала в хранилище», а не «когда её разобрали». У повтора +-- содержимое то же самое, отличается только метка, поэтому какая из двух +-- строк переживёт мерж, безразлично. +-- +-- Колонка версии не входит ни в ключ партиции, ни в ключ сортировки, и это +-- не случайность: попади она туда — версии одной строки лягут в разные куски +-- или в разные места ключа и не встретятся при мерже, то есть дедупликация +-- перестанет работать молча. +-- +-- PARTITION BY EventDate — по дню события, а не по месяцу, как у Метрики. +-- Отступление осознанное: дневная партиция здесь единица переобработки, и +-- переиграть день X значит заменить одну партицию. Месячная партиция тянула +-- бы за собой тридцать чужих дней. +-- +-- ORDER BY вырожден, и это учебный факт, а не недосмотр: CounterID на стенде +-- константа (сайт один), EventDate константа внутри своей партиции — обе +-- головные колонки ключа не различают ни одной строки. Реальная сортировка +-- идёт по посетителю и событию: intHash32(ClientID), WatchID. Ключ написан в +-- боевой форме «по сайту за период по посетителю» — на стенде она +-- вырождается, в бою нет. Хвост WatchID работает ещё и на дедупликацию: без +-- него ReplacingMergeTree схлопнул бы в одну строку все события посетителя за +-- день. +-- +-- SAMPLE BY intHash32(ClientID) — по тому же выражению, что стоит в ключе +-- сортировки (иначе семплирование не разрешено). Семплирование берёт целиком +-- посетителей, а не события вразнобой, поэтому SAMPLE 0.1 не портит +-- uniq-метрики. +-- +-- Sign всегда равен 1: это колонка формата, взятая без механики. В бою +-- исправление записи шлют парой −1/+1 и считают через sum(Sign) поверх +-- CollapsingMergeTree; наш генератор исправлений не шлёт, поэтому колонка +-- есть, а механики за ней нет. +CREATE TABLE IF NOT EXISTS ods.event_rep ON CLUSTER clickstream_cluster +( + WatchID UInt64, + VisitID UInt64, + ClientID UInt64, + CounterID UInt32, + EventDate Date, + UTCEventTime DateTime, + ClientTimeZone Int16, + EventType LowCardinality(String), + Sign Int8, + URL String, + Referer String, + Title String, + UTMSource String, + UTMMedium String, + UTMCampaign String, + UTMContent String, + UTMTerm String, + LastTrafficSource String, + HasGCLID UInt8, + YCLID UInt64, + Browser String, + BrowserMajorVersion UInt16, + BrowserLanguage String, + OperatingSystem String, + OperatingSystemRoot String, + DeviceCategory UInt8, + MobilePhoneModel String, + ScreenWidth UInt16, + ScreenHeight UInt16, + IPAddress String, + RegionCountry String, + RegionCity String, + RegionCountryID UInt32, + RegionCityID UInt32, + GoalsReached Array(UInt32), + ParsedParamsKey1 Array(String), + purchaseID Array(String), + purchaseRevenue Array(Float64), + purchaseCurrency Array(String), + purchaseCoupon Array(String), + productID Array(String), + productName Array(String), + productCategory Array(String), + productPrice Array(Int64), + productQuantity Array(UInt64), + productEventType Array(String), + ecommerce String, + _load_ts DateTime64(3) +) +ENGINE = ReplicatedReplacingMergeTree('/clickhouse/tables/{shard}/{database}/{table}', '{replica}', _load_ts) +PARTITION BY EventDate +ORDER BY (CounterID, EventDate, intHash32(ClientID), WatchID) +SAMPLE BY intHash32(ClientID); + +-- Лицо слоя: пишем и читаем через него. Ключ шардирования — +-- cityHash64(ClientID), а не сырой ClientID: структурированный числовой +-- идентификатор перекашивает остаток по модулю числа шардов, хеш — нет. +-- События одной куки при этом остаются на одном шарде, и сессионизация со +-- склейкой идентичностей ниже по течению живут локально. +-- +-- У сырья ключ другой (хеш строки), поэтому сырая строка и разобранное из неё +-- событие почти всегда лежат на разных шардах. Пошардовые счётчики слоёв +-- сходиться не должны — сверять слои можно только через _dist. +CREATE TABLE IF NOT EXISTS ods.event_dist ON CLUSTER clickstream_cluster +AS ods.event_rep +ENGINE = Distributed('clickstream_cluster', 'ods', 'event_rep', cityHash64(ClientID)); + +-- Локальная таблица ошибок разбора. +-- +-- Сырой текст сообщения, метаданные доставки и класс брака. Разобранных полей +-- у брака нет по определению: строка сюда попала как раз потому, что не +-- разобралась. +-- +-- Движок — обычный ReplicatedMergeTree, без замены версий: у брака нет ключа +-- сущности, схлопывать его не по чему. +-- +-- Ключи собственные. Шардирование — cityHash64(raw): ClientID у неразобранной +-- строки взять неоткуда. ORDER BY (error_class, kafka_partition, +-- kafka_offset): такую таблицу смотрят от класса брака, а внутри класса — по +-- координатам доставки. Нарезка по дню загрузки, как у сырья: модельного дня +-- у брака тоже нет. +-- +-- Срок жизни — месяц, вдесятеро дольше сырья. Истеки брак вместе с сырьём — +-- к моменту разбирательства не осталось бы ни того, ни другого. +CREATE TABLE IF NOT EXISTS ods.event_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)), + consumer_host LowCardinality(String), + _load_ts DateTime64(3) +) +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.event_errors_dist ON CLUSTER clickstream_cluster +AS ods.event_errors_rep +ENGINE = Distributed('clickstream_cluster', 'ods', 'event_errors_rep', cityHash64(raw)); diff --git a/sql/ddl/30-ods-views.sql b/sql/ddl/30-ods-views.sql new file mode 100644 index 0000000..454c56e --- /dev/null +++ b/sql/ddl/30-ods-views.sql @@ -0,0 +1,199 @@ +-- ODS: разбор сырья в событие и в таблицу ошибок. +-- +-- Матвью две, и вместе они обязаны делить поток без зазора и без нахлёста: +-- строка, подошедшая обеим, задвоится, а не подошедшая ни одной — исчезнет +-- молча. Держится это формой: годность считает один и тот же предикат из трёх +-- частей, а вторая матвью берёт его буквальное отрицание. +-- +-- Отсюда два запрета на выражения предиката, и оба серьёзные. +-- +-- Первый: ни одна его часть не возвращает NULL. Трёхзначная логика дала бы +-- строку, которую не берёт ни «условие», ни «NOT условие», — и потерялась бы +-- она бесшумно. Поэтому сравнения здесь дают 0 или 1, а разбор в Nullable +-- заканчивается IS NOT NULL. +-- +-- Второй: ни одна его часть не бросает исключений. Упавшая матвью роняет +-- вставку, офсеты не коммитятся, и потребление топика встаёт до починки — +-- это единственное, на чём сейчас держится правило репозитория «грязные +-- записи не валят пайплайн». Поэтому приведение Nullable к необнуляемому типу +-- (assumeNotNull) стоит только в списке колонок, за предикатом, который NULL +-- уже отсёк, а в самом предикате живут только функции семейства +-- JSONExtract/*OrNull. Довод целиком — ADR 0005. +-- +-- Источник у обеих — stg.hits_raw_dist, лицо слоя, а не локальная таблица: +-- потребитель цепляется к лицу, а матвью с источником-Distributed срабатывает +-- на вставку именно в распределённую таблицу, до раскладки по шардам. +-- +-- Порядок колонок в SELECT совпадает с порядком в целевой таблице. + +-- Годное событие: строка, прошедшая строгий приём. +-- +-- Строгий приём — это три вопроса, и первые два держат весь контракт схемы. +-- +-- 1. Это вообще объект JSON? Проверяется именно объект, а не валидность: +-- isValidJSON('123') возвращает единицу — скаляр тоже законный JSON +-- (измерено на стенде 7 августа 2026 года). +-- +-- 2. Совпадает ли набор ключей с контрактным — все сорок семь имён, ни одного +-- лишнего. Одно это сравнение заменяет сорок семь проверок на присутствие +-- и ловит то, чего иначе не поймать вовсе: опечатку в имени поля (для +-- хранилища это одновременно пропавшее ожидаемое и появившееся лишнее), +-- молчаливое расширение контракта источником и любую подмену имени в +-- колонке-массиве. Обязательны все сорок семь: генератор шлёт их в каждом +-- событии, а «пусто» по контракту — пустое значение, а не отсутствие +-- ключа. +-- +-- arraySort стоит с обеих сторон, и это не украшение. Без него сорок семь +-- CamelCase-имён пришлось бы выписать руками ровно в байтовом порядке — +-- ошибка, которая увела бы в брак вообще всё, и притом молча. +-- +-- 3. Разбираются ли пять опорных колонок в свой тип. Не сорок семь, а пять: +-- идентификаторы события, визита и посетителя, дата партиции и метка +-- времени. Порча любой из них отравляет всё ниже по течению, тогда как +-- единственный производитель топика — свой генератор, сериализующий по +-- объявленным типам, и неверный тип может прийти только из руки. Сорок +-- семь проверок на NULL превратили бы матвью в простыню, не добавив +-- защиты. Присутствие остальных сорока двух держит вопрос 2. +-- +-- Метку времени разбирает не JSONExtract, а parseDateTimeBestEffortOrNull, и +-- это измеренная необходимость, а не вкус. На проводе UTCEventTime уезжает в +-- ISO-8601 с суффиксом зоны — «2026-06-01T12:34:56Z» (спека генератора, +-- раздел 4), а JSONExtract с типом DateTime такую строку не берёт и отдаёт +-- NULL. Оставь его здесь — и в брак уехали бы все события до единого. Замер и +-- его подробности — ADR 0005, «Что проверено». EventDate в такой подпорке не +-- нуждается: дата уезжает как «2026-06-01», и JSONExtract её берёт. + +CREATE MATERIALIZED VIEW IF NOT EXISTS ods.event_mv ON CLUSTER clickstream_cluster +TO ods.event_dist +AS +WITH + JSONType(raw) = 'Object' AS is_object, + arraySort(JSONExtractKeys(raw)) = arraySort([ + 'WatchID', 'VisitID', 'ClientID', 'CounterID', 'EventDate', + 'UTCEventTime', 'ClientTimeZone', 'EventType', 'Sign', + 'URL', 'Referer', 'Title', 'UTMSource', 'UTMMedium', 'UTMCampaign', + 'UTMContent', 'UTMTerm', 'LastTrafficSource', 'HasGCLID', 'YCLID', + 'Browser', 'BrowserMajorVersion', 'BrowserLanguage', 'OperatingSystem', + 'OperatingSystemRoot', 'DeviceCategory', 'MobilePhoneModel', + 'ScreenWidth', 'ScreenHeight', 'IPAddress', 'RegionCountry', + 'RegionCity', 'RegionCountryID', 'RegionCityID', + 'GoalsReached', 'ParsedParamsKey1', + 'purchaseID', 'purchaseRevenue', 'purchaseCurrency', 'purchaseCoupon', + 'productID', 'productName', 'productCategory', 'productPrice', + 'productQuantity', 'productEventType', 'ecommerce' + ]) AS keys_match, + JSONExtract(raw, 'WatchID', 'Nullable(UInt64)') IS NOT NULL + AND JSONExtract(raw, 'VisitID', 'Nullable(UInt64)') IS NOT NULL + AND JSONExtract(raw, 'ClientID', 'Nullable(UInt64)') IS NOT NULL + AND JSONExtract(raw, 'EventDate', 'Nullable(Date)') IS NOT NULL + AND parseDateTimeBestEffortOrNull(JSONExtractString(raw, 'UTCEventTime')) + IS NOT NULL AS key_fields_parsed +SELECT + JSONExtract(raw, 'WatchID', 'UInt64') AS WatchID, + JSONExtract(raw, 'VisitID', 'UInt64') AS VisitID, + JSONExtract(raw, 'ClientID', 'UInt64') AS ClientID, + JSONExtract(raw, 'CounterID', 'UInt32') AS CounterID, + JSONExtract(raw, 'EventDate', 'Date') AS EventDate, + assumeNotNull(parseDateTimeBestEffortOrNull( + JSONExtractString(raw, 'UTCEventTime'))) AS UTCEventTime, + JSONExtract(raw, 'ClientTimeZone', 'Int16') AS ClientTimeZone, + JSONExtract(raw, 'EventType', 'String') AS EventType, + JSONExtract(raw, 'Sign', 'Int8') AS Sign, + JSONExtract(raw, 'URL', 'String') AS URL, + JSONExtract(raw, 'Referer', 'String') AS Referer, + JSONExtract(raw, 'Title', 'String') AS Title, + JSONExtract(raw, 'UTMSource', 'String') AS UTMSource, + JSONExtract(raw, 'UTMMedium', 'String') AS UTMMedium, + JSONExtract(raw, 'UTMCampaign', 'String') AS UTMCampaign, + JSONExtract(raw, 'UTMContent', 'String') AS UTMContent, + JSONExtract(raw, 'UTMTerm', 'String') AS UTMTerm, + JSONExtract(raw, 'LastTrafficSource', 'String') AS LastTrafficSource, + JSONExtract(raw, 'HasGCLID', 'UInt8') AS HasGCLID, + JSONExtract(raw, 'YCLID', 'UInt64') AS YCLID, + JSONExtract(raw, 'Browser', 'String') AS Browser, + JSONExtract(raw, 'BrowserMajorVersion', 'UInt16') AS BrowserMajorVersion, + JSONExtract(raw, 'BrowserLanguage', 'String') AS BrowserLanguage, + JSONExtract(raw, 'OperatingSystem', 'String') AS OperatingSystem, + JSONExtract(raw, 'OperatingSystemRoot', 'String') AS OperatingSystemRoot, + JSONExtract(raw, 'DeviceCategory', 'UInt8') AS DeviceCategory, + JSONExtract(raw, 'MobilePhoneModel', 'String') AS MobilePhoneModel, + JSONExtract(raw, 'ScreenWidth', 'UInt16') AS ScreenWidth, + JSONExtract(raw, 'ScreenHeight', 'UInt16') AS ScreenHeight, + JSONExtract(raw, 'IPAddress', 'String') AS IPAddress, + JSONExtract(raw, 'RegionCountry', 'String') AS RegionCountry, + JSONExtract(raw, 'RegionCity', 'String') AS RegionCity, + JSONExtract(raw, 'RegionCountryID', 'UInt32') AS RegionCountryID, + JSONExtract(raw, 'RegionCityID', 'UInt32') AS RegionCityID, + JSONExtract(raw, 'GoalsReached', 'Array(UInt32)') AS GoalsReached, + JSONExtract(raw, 'ParsedParamsKey1', 'Array(String)') AS ParsedParamsKey1, + JSONExtract(raw, 'purchaseID', 'Array(String)') AS purchaseID, + JSONExtract(raw, 'purchaseRevenue', 'Array(Float64)') AS purchaseRevenue, + JSONExtract(raw, 'purchaseCurrency', 'Array(String)') AS purchaseCurrency, + JSONExtract(raw, 'purchaseCoupon', 'Array(String)') AS purchaseCoupon, + JSONExtract(raw, 'productID', 'Array(String)') AS productID, + JSONExtract(raw, 'productName', 'Array(String)') AS productName, + JSONExtract(raw, 'productCategory', 'Array(String)') AS productCategory, + JSONExtract(raw, 'productPrice', 'Array(Int64)') AS productPrice, + JSONExtract(raw, 'productQuantity', 'Array(UInt64)') AS productQuantity, + JSONExtract(raw, 'productEventType', 'Array(String)') AS productEventType, + JSONExtract(raw, 'ecommerce', 'String') AS ecommerce, + -- Метка загрузки переносится из сырья как есть, а не ставится заново: + -- она отвечает на вопрос «когда строка приехала в хранилище». Поставь + -- здесь now64(3) — и колонка версии молча ответила бы на другой вопрос. + _load_ts +FROM stg.hits_raw_dist +WHERE is_object AND keys_match AND key_fields_parsed; + +-- Брак: буквальное отрицание того же предиката. +-- +-- Классы брака пересекаются, поэтому проверяются по порядку, а в error_class +-- пишется первый совпавший. Скаляр проваливает и проверку на объект, и сверку +-- ключей — JSONExtractKeys от него даёт пустой массив (измерено на стенде +-- 7 августа 2026 года), — и без объявленного порядка попал бы то в один +-- класс, то в другой. +-- +-- Предикат повторён здесь дословно, и это выбор, а не безвыходность: назвать +-- его один раз на две матвью позволил бы CREATE FUNCTION. Отвергнуто — условие +-- разбора ушло бы за имя, в отдельный объект со своей жизнью, и слой перестал +-- бы читаться по своему же DDL. Расхождения двух копий сторожит соседство: +-- обе живут в одном файле, в полусотне строк друг от друга. +CREATE MATERIALIZED VIEW IF NOT EXISTS ods.event_errors_mv ON CLUSTER clickstream_cluster +TO ods.event_errors_dist +AS +WITH + JSONType(raw) = 'Object' AS is_object, + arraySort(JSONExtractKeys(raw)) = arraySort([ + 'WatchID', 'VisitID', 'ClientID', 'CounterID', 'EventDate', + 'UTCEventTime', 'ClientTimeZone', 'EventType', 'Sign', + 'URL', 'Referer', 'Title', 'UTMSource', 'UTMMedium', 'UTMCampaign', + 'UTMContent', 'UTMTerm', 'LastTrafficSource', 'HasGCLID', 'YCLID', + 'Browser', 'BrowserMajorVersion', 'BrowserLanguage', 'OperatingSystem', + 'OperatingSystemRoot', 'DeviceCategory', 'MobilePhoneModel', + 'ScreenWidth', 'ScreenHeight', 'IPAddress', 'RegionCountry', + 'RegionCity', 'RegionCountryID', 'RegionCityID', + 'GoalsReached', 'ParsedParamsKey1', + 'purchaseID', 'purchaseRevenue', 'purchaseCurrency', 'purchaseCoupon', + 'productID', 'productName', 'productCategory', 'productPrice', + 'productQuantity', 'productEventType', 'ecommerce' + ]) AS keys_match, + JSONExtract(raw, 'WatchID', 'Nullable(UInt64)') IS NOT NULL + AND JSONExtract(raw, 'VisitID', 'Nullable(UInt64)') IS NOT NULL + AND JSONExtract(raw, 'ClientID', 'Nullable(UInt64)') IS NOT NULL + AND JSONExtract(raw, 'EventDate', 'Nullable(Date)') IS NOT NULL + AND parseDateTimeBestEffortOrNull(JSONExtractString(raw, 'UTCEventTime')) + IS NOT NULL AS key_fields_parsed +SELECT + raw, + multiIf( + NOT is_object, 'not_an_object', + NOT keys_match, 'keyset_mismatch', + 'key_field_unparsed' + ) AS error_class, + kafka_topic, + kafka_partition, + kafka_offset, + kafka_timestamp, + consumer_host, + _load_ts +FROM stg.hits_raw_dist +WHERE NOT (is_object AND keys_match AND key_fields_parsed);