feat(ods): типизированное событие, строгий приём и таблица ошибок #62
@@ -48,7 +48,7 @@ smoke:
|
|||||||
COMPOSE_BIN="$(COMPOSE)" ./scripts/stand-smoke.sh
|
COMPOSE_BIN="$(COMPOSE)" ./scripts/stand-smoke.sh
|
||||||
|
|
||||||
check-clickhouse:
|
check-clickhouse:
|
||||||
./scripts/clickhouse-smoke.sh
|
./scripts/check-clickhouse.sh
|
||||||
|
|
||||||
check-services:
|
check-services:
|
||||||
COMPOSE_BIN="$(COMPOSE)" ./scripts/stand-services.sh
|
COMPOSE_BIN="$(COMPOSE)" ./scripts/stand-services.sh
|
||||||
|
|||||||
@@ -125,15 +125,20 @@ STG сырьё исключительно от брака: слой сырых
|
|||||||
присутствие, и потому сужен до пяти колонок.
|
присутствие, и потому сужен до пяти колонок.
|
||||||
|
|
||||||
Цена решения. Контракт получает ещё два места: типы сорока семи колонок в
|
Цена решения. Контракт получает ещё два места: типы сорока семи колонок в
|
||||||
выражениях матвью и список тех же имён для сверки ключей. Раздел 1.4 спеки
|
выражениях матвью и список тех же имён для сверки ключей — а список этот
|
||||||
предупреждает, что, повторяясь примерно в семи местах, они расходятся молча, а
|
повторён дважды, по разу на матвью. Раздел 1.4 спеки предупреждает, что,
|
||||||
contract-тест из #43 сюда не дотягивается — он сравнивает `system.columns`
|
повторяясь примерно в семи местах, они расходятся молча.
|
||||||
целевой таблицы со схемой генератора и о выражениях матвью ничего не знает.
|
|
||||||
Смягчение работает не везде: опечатка в имени скалярного поля уводит строки в
|
И сверка ключей эту цену покрывает не всю: она смотрит на ключи сообщения, а
|
||||||
таблицу ошибок пачкой и видна сразу, а опечатка в имени массива даёт пустой
|
не на выражения матвью. Опечатка в имени внутри `JSONExtract` даёт умолчание
|
||||||
массив тихо — сверка ключей проверяет ключи сообщения, а не выражения матвью.
|
типа — ноль, пустую строку, пустой массив, — и молчит она у сорока двух
|
||||||
Эти двенадцать колонок сторожит smoke: известное событие с товарами обязано
|
обычных колонок ровно так же, как у двенадцати массивов. Громко ломаются
|
||||||
доезжать с непустыми массивами.
|
только пять опорных: у них разбор `Nullable` стоит в предикате, и опечатка
|
||||||
|
уводит в таблицу ошибок все строки до единой. Остальные сорок две сторожит
|
||||||
|
сверка разобранного события против сырого текста — разовый опыт при
|
||||||
|
исполнении #43, а не постоянная проверка; сила его в том, что выражения
|
||||||
|
сверки собираются из контракта, а выражения матвью написаны руками по
|
||||||
|
описанию выгрузки, и одна опечатка в двух местах не повторяется.
|
||||||
|
|
||||||
Второе — разбор функциями дороже разбора форматом. На объёмах стенда это
|
Второе — разбор функциями дороже разбора форматом. На объёмах стенда это
|
||||||
несущественно; если станет заметно, сорок семь вызовов сворачиваются в один
|
несущественно; если станет заметно, сорок семь вызовов сворачиваются в один
|
||||||
@@ -177,15 +182,62 @@ contract-тест из #43 сюда не дотягивается — он ср
|
|||||||
Тогда же нашлась и граница обещания — запись с пустым значением и
|
Тогда же нашлась и граница обещания — запись с пустым значением и
|
||||||
запись-надгробие не дают строки вовсе; из-за неё в «Решении» и в «Почему»
|
запись-надгробие не дают строки вовсе; из-за неё в «Решении» и в «Почему»
|
||||||
приписано слово «непустое». Замер целиком — в [доке
|
приписано слово «непустое». Замер целиком — в [доке
|
||||||
хранилища](../architecture/storage.md), раздел «Что проверено». Остальные три
|
хранилища](../architecture/storage.md), раздел «Что проверено».
|
||||||
по-прежнему ждут живого стенда; это однострочные `SELECT`, их довольно прогнать
|
|
||||||
заодно:
|
|
||||||
|
|
||||||
- форма именованного кортежа в `JSONExtract` с `Nullable`-членами — нужна для
|
Остальные четыре закрыты при исполнении #43, на стенде 7 августа 2026 года,
|
||||||
свёртки сорока семи вызовов в один, если разбор окажется дорогим;
|
ClickHouse 26.3.17.56. Все четыре ответили так, как ждала постановка:
|
||||||
- `isValidJSON('123')` возвращает единицу, а `JSONExtractKeys` от скаляра —
|
|
||||||
пустой массив. На обоих стоят классы брака и их приоритет, а документация
|
- `isValidJSON('123')` возвращает единицу — скаляр законный JSON, и первый
|
||||||
поведение на не-объекте не описывает: два `SELECT` закрывают вопрос;
|
класс брака поэтому проверяет именно объект, а не валидность;
|
||||||
- `JSONAsString` действительно падает на некорректном JSON, а не пропускает
|
- `JSONExtractKeys('123')` возвращает пустой массив, и он же приходит от
|
||||||
строку. На этом стоит отказ от него в пользу `RawBLOB`; документация про
|
вовсе не-JSON. Значит скаляр проваливает и сверку ключей — отсюда
|
||||||
ошибочный ввод молчит.
|
обязательный порядок классов. Раз `isValidJSON` объект от скаляра не
|
||||||
|
отличает, первый класс держит `JSONType`: она возвращает `Object` у объекта,
|
||||||
|
`Int64` у скаляра `123` и `Null` у вовсе не-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.
|
||||||
|
Настройка `cast_string_to_date_time_mode = 'best_effort'` дела не меняет:
|
||||||
|
с ней суффикс берёт `CAST`, а `JSONExtract` по-прежнему отдаёт NULL — то есть
|
||||||
|
внутренний разбор `JSONExtract` её не слушает. `EventDate` уезжает как
|
||||||
|
`2026-06-01` и разбирается `JSONExtract` без оговорок.
|
||||||
|
|
||||||
|
Разбор метки времени идёт
|
||||||
|
`parseDateTimeOrNull(JSONExtractString(raw, 'UTCEventTime'), '%Y-%m-%dT%H:%i:%SZ')`
|
||||||
|
— по буквально названному формату, а не через `parseDateTimeBestEffort`.
|
||||||
|
Обе функции ISO-8601 понимают и обе в варианте `*OrNull` отдают NULL вместо
|
||||||
|
исключения, то есть годятся в предикат. Выбран точный формат потому, что
|
||||||
|
широта здесь работает против строгого приёма: `parseDateTimeBestEffort` на
|
||||||
|
непонятной строке не краснеет, а достраивает недостающее. Измерено 7 августа
|
||||||
|
2026 года — обрезанное `20:00:21` он превращает в `2026-01-01 20:00:21`,
|
||||||
|
подставив текущий год и первое января; голая дата `2026-05-31` становится
|
||||||
|
полуночью; строка цифр читается числом эпохи. Такое сообщение прошло бы
|
||||||
|
строгий приём с тихо неверным временем — ровно с той порчей, ради которой
|
||||||
|
класс `key_field_unparsed` и заведён. Разбор по названному формату отдаёт на
|
||||||
|
всех трёх NULL. Источник у топика один и шлёт одну запись, так что широта не
|
||||||
|
нужна вовсе, а платится за неё отключённой проверкой.
|
||||||
|
|
||||||
|
Цена выбора измерена на настоящих данных: по всем 101 252 строкам сырья
|
||||||
|
модельного дня (день залит дважды) точный формат разобрал метку у каждой, и
|
||||||
|
ни на одной не разошёлся с `parseDateTimeBestEffort`. Различаются они только
|
||||||
|
на порче.
|
||||||
|
|
||||||
|
Оговорка к обнуляемому разбору даты, измеренная там же: `Nullable(Date)`
|
||||||
|
даёт NULL на строке, которая датой не является вовсе («мусор»), и на числе,
|
||||||
|
но невозможную дату `2026-13-99` молча приводит к `1970-01-01`. То есть
|
||||||
|
класс `key_field_unparsed` ловит порчу типа, а не порчу значения внутри типа.
|
||||||
|
|||||||
@@ -6,10 +6,11 @@
|
|||||||
раздел «Что проверено» — чему в этом тексте верить и на каком основании.
|
раздел «Что проверено» — чему в этом тексте верить и на каком основании.
|
||||||
|
|
||||||
**Что здесь описано и чего ещё нет.** Собран этап 1: кластер из двух шардов,
|
**Что здесь описано и чего ещё нет.** Собран этап 1: кластер из двух шардов,
|
||||||
keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/` уже лежат базы слоёв
|
keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/` лежит вся цепочка
|
||||||
и объекты STG — чтец топика `hits`, таблицы сырья и матвью приёма. Объектов ODS
|
`Kafka → STG → ODS` — чтец топика `hits`, таблицы сырья, типизированное
|
||||||
в репозитории пока нет. Дальше по тексту устройство описано так, как оно
|
событие с таблицей ошибок и три матвью. Дальше по тексту устройство описано
|
||||||
проектируется; построенное от заложенного отличает карта таблиц в конце.
|
так, как оно проектируется; построенное от заложенного отличает карта таблиц в
|
||||||
|
конце.
|
||||||
|
|
||||||
Зона ответственности у документа одна — хранилище. Генератор описан отдельно:
|
Зона ответственности у документа одна — хранилище. Генератор описан отдельно:
|
||||||
его замысел — в [спеке генератора](../specs/2026-08-01-generator.md), формат
|
его замысел — в [спеке генератора](../specs/2026-08-01-generator.md), формат
|
||||||
@@ -218,9 +219,11 @@ Airflow она стоит — там это обычный `SETTINGS` у зап
|
|||||||
нахлёста.** Одна забирает годные строки в `ods.event_dist`, вторая — брак в
|
нахлёста.** Одна забирает годные строки в `ods.event_dist`, вторая — брак в
|
||||||
`ods.event_errors_dist`. Строка, подошедшая обеим, задвоится; не подошедшая ни
|
`ods.event_errors_dist`. Строка, подошедшая обеим, задвоится; не подошедшая ни
|
||||||
одной — исчезнет молча. Держится это формой: второе условие пишется буквальным
|
одной — исчезнет молча. Держится это формой: второе условие пишется буквальным
|
||||||
отрицанием первого, а сам предикат собирается только из функций, не возвращающих
|
отрицанием первого, а сам предикат NULL не возвращает ни в одной своей части, —
|
||||||
NULL, — иначе трёхзначная логика даст строку, которую не возьмёт ни `условие`,
|
иначе трёхзначная логика даст строку, которую не возьмёт ни `условие`, ни
|
||||||
ни `NOT условие`. Те же функции не должны и бросать исключений: упавшая матвью
|
`NOT условие`. Обнуляемый разбор в предикате поэтому есть, но заканчивается
|
||||||
|
`IS NOT NULL`, а сравнения дают 0 или 1. Функции предиката не должны и бросать
|
||||||
|
исключений: упавшая матвью
|
||||||
роняет вставку и останавливает потребление до починки
|
роняет вставку и останавливает потребление до починки
|
||||||
([ADR 0005](../adr/0005-event-ingestion.md)).
|
([ADR 0005](../adr/0005-event-ingestion.md)).
|
||||||
|
|
||||||
@@ -229,9 +232,12 @@ NULL, — иначе трёхзначная логика даст строку,
|
|||||||
разъедется само: сырьё истечёт раньше, а перезаливка модельного дня задвоит его,
|
разъедется само: сырьё истечёт раньше, а перезаливка модельного дня задвоит его,
|
||||||
тогда как в ODS тот же повтор схлопнется. Правило, красное в норме, учит не
|
тогда как в ODS тот же повтор схлопнется. Правило, красное в норме, учит не
|
||||||
смотреть на оповещения
|
смотреть на оповещения
|
||||||
([ADR 0002](../adr/0002-monitoring-scope.md)). Равенство проверяется разово в
|
([ADR 0002](../adr/0002-monitoring-scope.md)). Равенство проверено разовым
|
||||||
smoke на управляемой пачке: отправили N сообщений — получили N строк сырья и N в
|
опытом при исполнении #43, на управляемой пачке: отправили N сообщений —
|
||||||
сумме событий и ошибок. Счёт по ODS идёт через `FINAL`: голый `count()` по
|
получили N строк сырья и N в сумме событий и ошибок. Постоянной целью такой
|
||||||
|
опыт не становится, и почему — в [карте
|
||||||
|
проверок](testing.md), раздел «Интеграционная проверка постоянной целью не
|
||||||
|
становится». Счёт по ODS идёт через `FINAL`: голый `count()` по
|
||||||
`ReplacingMergeTree` зависит от того, сколько мержей успело пройти, и спека это
|
`ReplacingMergeTree` зависит от того, сколько мержей успело пройти, и спека это
|
||||||
прямо запрещает (раздел 6).
|
прямо запрещает (раздел 6).
|
||||||
|
|
||||||
@@ -279,7 +285,10 @@ D0 и к реальному календарю не привязана; паке
|
|||||||
брака нет разобранных полей: шардируется `cityHash64` сырой строки — `ClientID` у
|
брака нет разобранных полей: шардируется `cityHash64` сырой строки — `ClientID` у
|
||||||
строки, которая не разобралась, взять неоткуда; нарезается по дню загрузки, как и
|
строки, которая не разобралась, взять неоткуда; нарезается по дню загрузки, как и
|
||||||
сырьё; живёт месяц. Дольше сырья — намеренно: если брак истекает вместе с ним,
|
сырьё; живёт месяц. Дольше сырья — намеренно: если брак истекает вместе с ним,
|
||||||
разбираться к моменту разбирательства будет уже нечем.
|
разбираться к моменту разбирательства будет уже нечем. Снятие — целыми кусками
|
||||||
|
(`ttl_only_drop_parts`), как у сырья и по той же причине: строки в партиции дня
|
||||||
|
загрузки разного возраста не более чем на сутки, и куску незачем переживать
|
||||||
|
срок из-за самой свежей строки.
|
||||||
|
|
||||||
Класс брака лежит в колонке `error_class` типа `LowCardinality(String)`. Без неё
|
Класс брака лежит в колонке `error_class` типа `LowCardinality(String)`. Без неё
|
||||||
в таблице копятся строки «что-то не так» без ответа на «что именно», а витрине
|
в таблице копятся строки «что-то не так» без ответа на «что именно», а витрине
|
||||||
@@ -345,8 +354,7 @@ ODS. Второе: матвью приёма создаётся последне
|
|||||||
|
|
||||||
## Карта таблиц
|
## Карта таблиц
|
||||||
|
|
||||||
Ниже — то, что закладывает этап 2. DDL слоя STG уже лежит в `sql/ddl/`;
|
Ниже — то, что закладывает этап 2; всё перечисленное лежит в `sql/ddl/`.
|
||||||
объектов ODS в репозитории пока нет.
|
|
||||||
|
|
||||||
| Слой | Объект | Что это |
|
| Слой | Объект | Что это |
|
||||||
|---|---|---|
|
|---|---|---|
|
||||||
@@ -382,9 +390,42 @@ ODS. Второе: матвью приёма создаётся последне
|
|||||||
завязан на движок базы `Atomic`. `ON CLUSTER` ждёт все хосты и бросает по
|
завязан на движок базы `Atomic`. `ON CLUSTER` ждёт все хосты и бросает по
|
||||||
таймауту; `CREATE ... IF NOT EXISTS` на существующем объекте не бросает.
|
таймауту; `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` кладёт блок в спул и сразу возвращает управление, так что на
|
||||||
|
сносе локальной ошибка всплыла бы фоном и утверждение показалось бы
|
||||||
|
опровергнутым.
|
||||||
|
- `_load_ts` в `ods.event` — это метка исходной строки сырья, а не время
|
||||||
|
разбора. Сверено по `WatchID` на двух тысячах событий модельного дня,
|
||||||
|
залитого дважды: у каждого события метка совпала с меткой одной из двух его
|
||||||
|
доставок, а после `FINAL` — с меткой поздней. Случаев «метки нет среди
|
||||||
|
доставок» ноль, то есть `now64()` в матвью разбора нет.
|
||||||
|
- Пересозданная матвью пропускает ближайшие сообщения. Снятые и заново
|
||||||
|
созданные матвью разбора при живом чтеце: сообщение, отправленное сразу
|
||||||
|
после, легло в сырьё и не попало в ODS никуда — ни в событие, ни в ошибки;
|
||||||
|
то же сообщение через минуту разобралось штатно. Воспроизведено дважды
|
||||||
|
7 августа 2026 года. Это тот же зазор, о котором предупреждает нумерация
|
||||||
|
файлов DDL, только приходит он с другой стороны — не при первом создании, а
|
||||||
|
при замене матвью на работающем стенде. Практический вывод один: правишь
|
||||||
|
матвью — не верь ближайшей отправке, повтори её. Чем именно держится
|
||||||
|
задержка, не измерено; наблюдение записано как наблюдение.
|
||||||
|
- Форма ключа `ods.event` принимается такой, как её задумала спека: выражение
|
||||||
|
`intHash32(ClientID)` стоит в ключе сортировки `ReplacingMergeTree`, а
|
||||||
|
`SAMPLE BY` — по тому же выражению. Вопрос стоял открытым в разделе 11
|
||||||
|
мастер-спеки; ответ — DDL применяется и таблица работает.
|
||||||
|
|
||||||
- `RawBLOB` даёт ровно одну строку на каждое непустое сообщение. Три сообщения
|
- `RawBLOB` даёт ровно одну строку на каждое непустое сообщение. Три сообщения
|
||||||
с ключами, поставленные в очередь до одного сброса продюсера, стали тремя
|
с ключами, поставленные в очередь до одного сброса продюсера, стали тремя
|
||||||
@@ -433,12 +474,5 @@ Kafka с пустым значением (ноль байт) и запись-н
|
|||||||
таблице. Это вычитано, а не измерено. На устройство приёма оговорка не влияет:
|
таблице. Это вычитано, а не измерено. На устройство приёма оговорка не влияет:
|
||||||
служебные колонки мы заполняем выражением при любом ответе.
|
служебные колонки мы заполняем выражением при любом ответе.
|
||||||
|
|
||||||
**Сказано по памяти, проверки пока нет.** Осталось два утверждения, и оба ждут
|
**Сказано по памяти, проверки нет.** Группа пуста: оба утверждения, ждавшие
|
||||||
одного и того же — матвью разбора, а она приходит с #43.
|
матвью разбора, закрыты опытами при исполнении #43 и переехали выше.
|
||||||
|
|
||||||
- Матвью с источником-`Distributed` срабатывает на вставку именно в эту
|
|
||||||
распределённую таблицу, до раскладки по шардам. В документации случая нет
|
|
||||||
вовсе, утверждение держится на опыте владельца.
|
|
||||||
- Упавшая матвью роняет вставку и останавливает потребление до починки. На этой
|
|
||||||
фразе стоит правило «грязные записи не валят пайплайн», а сама она стоит пока
|
|
||||||
на одном рассуждении.
|
|
||||||
|
|||||||
@@ -22,10 +22,10 @@
|
|||||||
макросов `shard`: обе ноды здоровы и порты отвечают. `make check-clickhouse` не
|
макросов `shard`: обе ноды здоровы и порты отвечают. `make check-clickhouse` не
|
||||||
заметит потерянного подключения Superset: он про ClickHouse и только.
|
заметит потерянного подключения Superset: он про ClickHouse и только.
|
||||||
|
|
||||||
Отсюда правило для новой проверки: **спроси, кого она спрашивает.** Договор со
|
Отсюда правило для новой проверки: **спроси, кого она спрашивает.** Счётчики
|
||||||
схемой событий — вопрос к ClickHouse, значит дом ему в `check-clickhouse`, даже
|
против манифеста — вопрос к ClickHouse: строки в `ods.event` считает сам
|
||||||
если по цене он подошёл бы смоуку. Счётчики против манифеста — тоже вопрос к
|
сервер и отвечает сразу, значит дом им в `check-clickhouse`, даже если по цене
|
||||||
ClickHouse: строки в `ods.event` считает сам сервер и отвечает сразу.
|
они подошли бы смоуку.
|
||||||
|
|
||||||
## Карта целей
|
## Карта целей
|
||||||
|
|
||||||
@@ -128,11 +128,7 @@ ClickHouse отвечает сразу.
|
|||||||
| `make config-test` | `scripts/config-test.sh` |
|
| `make config-test` | `scripts/config-test.sh` |
|
||||||
| `make smoke` | `scripts/stand-smoke.sh` |
|
| `make smoke` | `scripts/stand-smoke.sh` |
|
||||||
| `make check-services` | `scripts/stand-services.sh` |
|
| `make check-services` | `scripts/stand-services.sh` |
|
||||||
| `make check-clickhouse` | `scripts/clickhouse-smoke.sh` |
|
| `make check-clickhouse` | `scripts/check-clickhouse.sh` |
|
||||||
|
|
||||||
Имя `clickhouse-smoke.sh` осталось от прежнего имени цели — `smoke-cluster`.
|
|
||||||
Файл переименуют при следующем касании: сейчас в него встраивается проверка
|
|
||||||
договора со схемой, и переименование устроило бы конфликт на ровном месте.
|
|
||||||
|
|
||||||
Общее у смоука и `check-services` — счёт проверок, обращение к Compose и две
|
Общее у смоука и `check-services` — счёт проверок, обращение к Compose и две
|
||||||
проверки — вынесено в `scripts/stand-common.sh`; сам он не запускается.
|
проверки — вынесено в `scripts/stand-common.sh`; сам он не запускается.
|
||||||
|
|||||||
@@ -172,8 +172,9 @@ Ecommerce (заполнены только у торговых событий):
|
|||||||
рендеренное «описание выгрузки» в доках — аналог документации Метрики.
|
рендеренное «описание выгрузки» в доках — аналог документации Метрики.
|
||||||
Сторона хранилища (DDL, SELECT матвью, трансформации, витрины) пишется по
|
Сторона хранилища (DDL, SELECT матвью, трансформации, витрины) пишется по
|
||||||
этой документации на своих этапах, как в бою хранилище адаптируется к
|
этой документации на своих этапах, как в бою хранилище адаптируется к
|
||||||
источнику; границу сторожат строгий приём (раздел 6) и contract-тест в
|
источнику; границу сторожит строгий приём (раздел 6). Вторым сторожем здесь
|
||||||
smoke — сравнение `system.columns` поднятого стенда со схемой генератора.
|
стояла сверка объявлений — `system.columns` поднятого стенда против схемы
|
||||||
|
генератора; она снята при исполнении #43 как ничего не добавляющая к соседу.
|
||||||
Без контракта 47 колонок, повторяясь примерно в семи местах, расходятся
|
Без контракта 47 колонок, повторяясь примерно в семи местах, расходятся
|
||||||
молча. Заодно это учебный артефакт: менти видит на живом примере, что
|
молча. Заодно это учебный артефакт: менти видит на живом примере, что
|
||||||
такое data contract.
|
такое data contract.
|
||||||
@@ -584,8 +585,10 @@ v2 стартует пустым, поэтому объём ниже — это
|
|||||||
|
|
||||||
Список убывает по мере постройки: проверенное уходит отсюда, а ответ с датой
|
Список убывает по мере постройки: проверенное уходит отсюда, а ответ с датой
|
||||||
остаётся там, где на него опираются. Формат чтеца и форма виртуальной метки
|
остаётся там, где на него опираются. Формат чтеца и форма виртуальной метки
|
||||||
времени закрыты при исполнении #37 — см. [доку
|
времени закрыты при исполнении #37; форма ключа ODS и поведение матвью над
|
||||||
хранилища](../architecture/storage.md), раздел «Что проверено».
|
`Distributed` — при исполнении #43, ответы в [доке
|
||||||
|
хранилища](../architecture/storage.md), раздел «Что проверено»; запасной
|
||||||
|
именованный кортеж — там же в [ADR 0005](../adr/0005-event-ingestion.md).
|
||||||
|
|
||||||
- Поведение соединения двух Distributed-таблиц и `distributed_product_mode` —
|
- Поведение соединения двух Distributed-таблиц и `distributed_product_mode` —
|
||||||
эмпирически на стенде (хвост #14).
|
эмпирически на стенде (хвост #14).
|
||||||
@@ -593,14 +596,6 @@ v2 стартует пустым, поэтому объём ниже — это
|
|||||||
прогонами, отсутствие дублей при штатной работе. Закрыто пока наполовину: что
|
прогонами, отсутствие дублей при штатной работе. Закрыто пока наполовину: что
|
||||||
обе ноды читают топик и обе партиции доезжают, показал #37; что дублей нет и
|
обе ноды читают топик и обе партиции доезжают, показал #37; что дублей нет и
|
||||||
как раскладка меняется между прогонами — нет.
|
как раскладка меняется между прогонами — нет.
|
||||||
- Точная форма `ORDER BY` ODS-таблиц (выражение `intHash32` в ключе
|
|
||||||
ReplacingMergeTree).
|
|
||||||
- Матвью с источником-`Distributed` срабатывает на вставку именно в эту
|
|
||||||
распределённую таблицу, до раскладки по шардам: на этом стоит цепочка
|
|
||||||
STG → ODS (ADR 0005). Проверено владельцем на рабочих проектах, в документации
|
|
||||||
ClickHouse этот случай не описан.
|
|
||||||
- Форма именованного кортежа в `JSONExtract` с `Nullable`-членами — ею
|
|
||||||
сворачиваются 47 вызовов в один, если разбор окажется дорогим (ADR 0005).
|
|
||||||
- Размер артефакта эталонного мира после пересборки.
|
- Размер артефакта эталонного мира после пересборки.
|
||||||
- Спорные API (Airflow Datasets/сенсоры, ClickHouse DDL) — перед кодом
|
- Спорные API (Airflow Datasets/сенсоры, ClickHouse DDL) — перед кодом
|
||||||
сверять через MCP Context7 (правило AGENTS.md).
|
сверять через MCP Context7 (правило AGENTS.md).
|
||||||
|
|||||||
@@ -33,8 +33,8 @@
|
|||||||
- **Детерминизм до байта.** Одно зерно — побайтово тот же снимок; сверка —
|
- **Детерминизм до байта.** Одно зерно — побайтово тот же снимок; сверка —
|
||||||
хешами манифеста. Транспорт (офсеты Kafka, темп) — вне обещания.
|
хешами манифеста. Транспорт (офсеты Kafka, темп) — вне обещания.
|
||||||
- **Схема — контракт генератора.** Python-модуль с чистыми данными;
|
- **Схема — контракт генератора.** Python-модуль с чистыми данными;
|
||||||
хранилище строится по рендеренной документации, границу сторожит
|
хранилище строится по рендеренной документации, границу сторожит строгий
|
||||||
contract-тест.
|
приём на стороне хранилища.
|
||||||
- **Один сериализатор, глупые приёмники.** День-функция выдаёт канонические
|
- **Один сериализатор, глупые приёмники.** День-функция выдаёт канонические
|
||||||
байты; приёмники — файл, Kafka пачкой, Kafka с темпом.
|
байты; приёмники — файл, Kafka пачкой, Kafka с темпом.
|
||||||
- **Числа.** Средний день ~50 тыс. событий; эталонный снимок — 14 дней;
|
- **Числа.** Средний день ~50 тыс. событий; эталонный снимок — 14 дней;
|
||||||
@@ -229,10 +229,13 @@
|
|||||||
- **Сторона хранилища пишется по документации, не генерируется.** DDL
|
- **Сторона хранилища пишется по документации, не генерируется.** DDL
|
||||||
`ods.event`, SELECT матвью, `dds.event_v`, трансформации — работа
|
`ods.event`, SELECT матвью, `dds.event_v`, трансформации — работа
|
||||||
следующих этапов по «описанию выгрузки», как в бою хранилище адаптируется
|
следующих этапов по «описанию выгрузки», как в бою хранилище адаптируется
|
||||||
к источнику. Границу сторожат два боевых механизма: строгий приём
|
к источнику. Границу сторожит боевой механизм — строгий приём
|
||||||
(`Nullable`-разбор со сверкой набора ключей, таблицы `*_errors` — раздел 6
|
(`Nullable`-разбор со сверкой набора ключей, таблицы `*_errors` — раздел 6
|
||||||
мастер-спеки) и contract-тест в smoke — сравнение `system.columns`
|
мастер-спеки). Сверка объявлений (`system.columns` поднятого стенда против
|
||||||
поднятого стенда со схемой генератора.
|
контракта) здесь стояла вторым механизмом и снята при исполнении #43:
|
||||||
|
сверх строгого приёма она ловила ровно одно — смену типа колонки, — а её
|
||||||
|
ловит и сверка разобранного события, причём на живых данных, а не на
|
||||||
|
объявлениях.
|
||||||
|
|
||||||
Отклонено с доводами:
|
Отклонено с доводами:
|
||||||
|
|
||||||
@@ -263,10 +266,20 @@
|
|||||||
- **Даты и время на проводе — ISO-8601.** `EventDate` уезжает как `2026-06-01`,
|
- **Даты и время на проводе — ISO-8601.** `EventDate` уезжает как `2026-06-01`,
|
||||||
`UTCEventTime` — как `2026-06-01T12:34:56Z`. Довод — читаемость сырья: весь
|
`UTCEventTime` — как `2026-06-01T12:34:56Z`. Довод — читаемость сырья: весь
|
||||||
смысл слоя STG в том, что менти открывает колонку `raw` в обычном клиенте и
|
смысл слоя STG в том, что менти открывает колонку `raw` в обычном клиенте и
|
||||||
разбирает событие глазами, а число эпохи этот урок убивает. Разбору это
|
разбирает событие глазами, а число эпохи этот урок убивает. Колонка
|
||||||
ничего не стоит: `JSONExtract(raw, 'UTCEventTime', 'Nullable(DateTime)')`
|
`ecommerce` — строка, внутри которой лежит экранированный JSON, как отдаёт
|
||||||
принимает ISO без плясок. Колонка `ecommerce` — строка, внутри которой лежит
|
Метрика.
|
||||||
экранированный JSON, как отдаёт Метрика.
|
|
||||||
|
Оговорка про цену разбора, вписанная сюда 6 августа и оказавшаяся неверной:
|
||||||
|
здесь стояло, что `JSONExtract(raw, 'UTCEventTime', 'Nullable(DateTime)')`
|
||||||
|
«принимает ISO без плясок». Не принимает — на строке с суффиксом `Z` он
|
||||||
|
отдаёт NULL, и при исполнении #43 это увело бы в брак все события до
|
||||||
|
единого. Измерено на стенде 7 августа 2026 года; форма на проводе от этого
|
||||||
|
не меняется, меняется выражение разбора на стороне хранилища — `parseDateTime`
|
||||||
|
по буквально названному формату вместо `JSONExtract`
|
||||||
|
([ADR 0005](../adr/0005-event-ingestion.md)). То, что форма на проводе одна и
|
||||||
|
каноническая, здесь работает на хранилище: раз запись ровно одна, разбирать
|
||||||
|
её можно строго, не принимая заодно десяток чужих записей.
|
||||||
|
|
||||||
Форму реализует сериализатор (#41), хранилище (#43) читает то, что он
|
Форму реализует сериализатор (#41), хранилище (#43) читает то, что он
|
||||||
положил: порядок тикетов развёрнут 6 августа 2026 года, и отправитель идёт
|
положил: порядок тикетов развёрнут 6 августа 2026 года, и отправитель идёт
|
||||||
@@ -283,7 +296,8 @@
|
|||||||
ReplacingMergeTree.
|
ReplacingMergeTree.
|
||||||
- **Эталонный снимок при старте стенда — через Kafka, пакетным режимом
|
- **Эталонный снимок при старте стенда — через Kafka, пакетным режимом
|
||||||
проигрывателя.** Отдельный механизм заливки не строится: каждый `make up`
|
проигрывателя.** Отдельный механизм заливки не строится: каждый `make up`
|
||||||
бесплатно прогоняет весь конвейер и contract-тест на настоящих данных.
|
бесплатно прогоняет весь конвейер на настоящих данных, и строгий приём
|
||||||
|
хранилища проверяет контракт тем же прогоном.
|
||||||
Оговорка «если заливка уйдёт в десятки минут — вернуться к прямой
|
Оговорка «если заливка уйдёт в десятки минут — вернуться к прямой
|
||||||
загрузке» проверена при фиксации чисел: 14 × 50 тыс. ≈ 700 тыс. событий —
|
загрузке» проверена при фиксации чисел: 14 × 50 тыс. ≈ 700 тыс. событий —
|
||||||
расчётно минута-две, запас есть.
|
расчётно минута-две, запас есть.
|
||||||
@@ -395,8 +409,8 @@ pytest-тест с маркером `perf` и таймаутом-обрубан
|
|||||||
Внесены в мастер-спеку тем же коммитом, что и эта спека:
|
Внесены в мастер-спеку тем же коммитом, что и эта спека:
|
||||||
|
|
||||||
- **Раздел 1.4**: «из контракта выводятся DDL и валидация» заменено на data
|
- **Раздел 1.4**: «из контракта выводятся DDL и валидация» заменено на data
|
||||||
contract — хранилище пишется по документации, границу сторожит
|
contract — хранилище пишется по документации, границу сторожит строгий
|
||||||
contract-тест (раздел 3 здесь).
|
приём (раздел 3 здесь).
|
||||||
- **Раздел 8**: артефакт `data/startup_history/` в git заменён манифестом;
|
- **Раздел 8**: артефакт `data/startup_history/` в git заменён манифестом;
|
||||||
снимок генерируется на месте (раздел 5 здесь). Туман «политика
|
снимок генерируется на месте (раздел 5 здесь). Туман «политика
|
||||||
версионирования артефакта» закрыт этим же ходом: версионируется манифест.
|
версионирования артефакта» закрыт этим же ходом: версионируется манифест.
|
||||||
|
|||||||
@@ -4,9 +4,9 @@
|
|||||||
1.1–1.2) и здесь не переоткрываются — модуль записывает их машинно-читаемо.
|
1.1–1.2) и здесь не переоткрываются — модуль записывает их машинно-читаемо.
|
||||||
Контракт принадлежит генератору и кормит трёх потребителей: сам генератор,
|
Контракт принадлежит генератору и кормит трёх потребителей: сам генератор,
|
||||||
его валидацию и «описание выгрузки» в доках (`schema_doc`). Хранилище
|
его валидацию и «описание выгрузки» в доках (`schema_doc`). Хранилище
|
||||||
строится по описанию, а не по модулю; границу будет сторожить contract-тест,
|
строится по описанию, а не по модулю; границу сторожит строгий приём на его
|
||||||
сверяющий `system.columns` поднятого стенда с этим контрактом, — он придёт
|
стороне — сверка набора ключей сообщения с контрактным списком, и
|
||||||
вместе с типизированным ODS (спека генератора, раздел 3).
|
разошедшееся уходит в таблицу ошибок (спека генератора, раздел 3).
|
||||||
|
|
||||||
Что несёт описатель колонки:
|
Что несёт описатель колонки:
|
||||||
|
|
||||||
@@ -16,7 +16,8 @@
|
|||||||
`system.columns`: параметры входят в имя типа целиком, без сокращений
|
`system.columns`: параметры входят в имя типа целиком, без сокращений
|
||||||
(`LowCardinality(String)`, `Array(Float64)`). Сверено 2026-08-01 —
|
(`LowCardinality(String)`, `Array(Float64)`). Сверено 2026-08-01 —
|
||||||
по документации ClickHouse через Context7 и запросом к узлу стенда
|
по документации ClickHouse через Context7 и запросом к узлу стенда
|
||||||
(26.3.17.56); от этой записи зависит будущий contract-тест.
|
(26.3.17.56). Запись важна потому, что по ней человек пишет DDL: тип,
|
||||||
|
сокращённый здесь, приедет в таблицу сокращённым же.
|
||||||
- `numpy_dtype` — чем колонка представлена внутри генератора; у массивов это
|
- `numpy_dtype` — чем колонка представлена внутри генератора; у массивов это
|
||||||
тип элемента. Строки живут в `object`-массивах: numpy-строки фиксированной
|
тип элемента. Строки живут в `object`-массивах: numpy-строки фиксированной
|
||||||
длины стенду ничего не дают.
|
длины стенду ничего не дают.
|
||||||
|
|||||||
@@ -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));
|
||||||
@@ -0,0 +1,191 @@
|
|||||||
|
-- ODS: разбор сырья в событие и в таблицу ошибок.
|
||||||
|
--
|
||||||
|
-- Матвью две, и вместе они обязаны делить поток без зазора и без нахлёста:
|
||||||
|
-- строка, подошедшая обеим, задвоится, а не подошедшая ни одной — исчезнет
|
||||||
|
-- молча. Держится это формой: годность считает один и тот же предикат из трёх
|
||||||
|
-- частей, а вторая матвью берёт его буквальное отрицание.
|
||||||
|
--
|
||||||
|
-- Отсюда два запрета на выражения предиката, и оба серьёзные.
|
||||||
|
--
|
||||||
|
-- Первый: ни одна его часть не возвращает NULL. Трёхзначная логика дала бы
|
||||||
|
-- строку, которую не берёт ни «условие», ни «NOT условие», — и потерялась бы
|
||||||
|
-- она бесшумно. Поэтому сравнения здесь дают 0 или 1, а разбор в Nullable
|
||||||
|
-- заканчивается IS NOT NULL.
|
||||||
|
--
|
||||||
|
-- Второй: ни одна его часть не бросает исключений. Упавшая матвью роняет
|
||||||
|
-- вставку, офсеты не коммитятся, и потребление топика встаёт до починки —
|
||||||
|
-- это единственное, на чём сейчас держится правило репозитория «грязные
|
||||||
|
-- записи не валят пайплайн». Поэтому приведение Nullable к необнуляемому типу
|
||||||
|
-- (assumeNotNull) стоит только в списке колонок, за предикатом, который NULL
|
||||||
|
-- уже отсёк, а в самом предикате живут только функции семейства
|
||||||
|
-- JSONExtract/*OrNull. Довод целиком — ADR 0005.
|
||||||
|
--
|
||||||
|
-- Источник у обеих — stg.hits_raw_dist, лицо слоя, а не локальная таблица:
|
||||||
|
-- потребитель цепляется к лицу, а матвью с источником-Distributed срабатывает
|
||||||
|
-- на вставку именно в распределённую таблицу, до раскладки по шардам.
|
||||||
|
--
|
||||||
|
-- Порядок колонок в SELECT совпадает с порядком в целевой таблице.
|
||||||
|
|
||||||
|
-- Годное событие: строка, прошедшая строгий приём.
|
||||||
|
--
|
||||||
|
-- Строгий приём — это три вопроса: объект ли это JSON, тот ли набор ключей,
|
||||||
|
-- разбираются ли в свой тип пять опорных колонок. Почему именно так — почему
|
||||||
|
-- объект, а не валидность; почему сверка ключей заменяет сорок семь проверок
|
||||||
|
-- на присутствие; почему опорных пять, а не сорок семь — ADR 0005, «Решение».
|
||||||
|
--
|
||||||
|
-- arraySort ниже стоит с обеих сторон, и убирать его нельзя: без обёртки сорок
|
||||||
|
-- семь CamelCase-имён пришлось бы держать ровно в байтовом порядке, а сбой
|
||||||
|
-- порядка увёл бы в брак вообще всё, и молча (ADR 0005).
|
||||||
|
--
|
||||||
|
-- Метка времени — единственная из пяти, кого разбирает не JSONExtract, и это
|
||||||
|
-- измеренная необходимость, а не вкус. На проводе UTCEventTime уезжает в
|
||||||
|
-- ISO-8601 с суффиксом зоны — «2026-06-01T12:34:56Z» (спека генератора,
|
||||||
|
-- раздел 4), а JSONExtract с типом DateTime такую строку не берёт и отдаёт
|
||||||
|
-- NULL. Оставь его здесь — и в брак уехали бы все события до единого.
|
||||||
|
--
|
||||||
|
-- Формат назван буквально, а не отдан parseDateTimeBestEffort, и вот почему.
|
||||||
|
-- Best-effort понимает десяток записей и на непонятной не краснеет, а
|
||||||
|
-- достраивает недостающее: обрезанное «20:00:21» он превращает в первое
|
||||||
|
-- января текущего года. Такая строка прошла бы строгий приём с тихо неверным
|
||||||
|
-- временем — ровно с той порчей, ради которой класс key_field_unparsed и
|
||||||
|
-- заведён. Источник у топика один и шлёт одну запись, так что широта здесь не
|
||||||
|
-- нужна вовсе, а стоит она отключённой проверкой. Замеры — 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 parseDateTimeOrNull(JSONExtractString(raw, 'UTCEventTime'),
|
||||||
|
'%Y-%m-%dT%H:%i:%SZ') 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(parseDateTimeOrNull(
|
||||||
|
JSONExtractString(raw, 'UTCEventTime'),
|
||||||
|
'%Y-%m-%dT%H:%i:%SZ')) 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 parseDateTimeOrNull(JSONExtractString(raw, 'UTCEventTime'),
|
||||||
|
'%Y-%m-%dT%H:%i:%SZ') 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);
|
||||||
Reference in New Issue
Block a user