docs(storage): конвенции и приём событий выправлены после ревью

- Зачем:
  - три холодных ревью и сверка с документацией ClickHouse нашли противоречия
    между докой, ADR и спекой: исполнитель #37 получал два разных ответа на
    один вопрос, а два утверждения о движке оказались неверными.
- Что:
  - раскладка файлов DDL перестроена — сначала таблицы, матвью приёма
    последней: иначе часть событий тихо минует ODS.
  - синхронная вставка снята с пути приёма: настройка недостижима для потока
    Kafka-движка и связывает шарды; на ETL-вставках осталась.
  - у таблицы ошибок появился класс брака с порядком проверки, у сырья и
    ошибок названы движки и ключи сортировки.
  - в доку добавлен раздел «Что проверено»: сверенное с документацией,
    проверяемое на стенде и сказанное по памяти разведены.
  - в спеке выправлены источник матвью разбора, пять опорных колонок, имена
    четырёх витрин и ссылка на несуществующую цель make.
- Проверка:
  - make config-test
  - grep по устаревшим именам файлов DDL и витрин — пусто

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
2026-08-05 21:25:35 +03:00
co-authored by Claude Opus 5
parent dedb6c6739
commit 31b274175a
6 changed files with 284 additions and 65 deletions
+4
View File
@@ -125,3 +125,7 @@ Airflow) и названия из кода. Если для понятия ес
- Имена файлов в `docs/adr/``NNNN-краткое-имя.md`: сквозной номер из четырёх - Имена файлов в `docs/adr/``NNNN-краткое-имя.md`: сквозной номер из четырёх
цифр и слаг (`0001-stand-services.md`). Решения нумеруются подряд, дата цифр и слаг (`0001-stand-services.md`). Решения нумеруются подряд, дата
в имени не нужна. в имени не нужна.
- Имена файлов в `docs/architecture/` — слаг строчными латинскими буквами через
дефис (`storage.md`). Здесь живут рабочие справочники по зонам
ответственности: не событие истории и не решение, а текущее устройство —
один файл на зону, правится по мере постройки.
+5
View File
@@ -204,6 +204,11 @@ Superset закреплён на 6.1.0; драйвер `clickhouse-connect`, ф
## Документация ## Документация
- [docs/specs/](docs/specs/) — спеки: источник истины о задуманном. - [docs/specs/](docs/specs/) — спеки: источник истины о задуманном.
- [docs/adr/](docs/adr/) — принятые решения с доводами и отвергнутыми
вариантами: почему сделано так, а не иначе.
- [docs/architecture/](docs/architecture/) — рабочие справочники по зонам
ответственности; сейчас это [хранилище](docs/architecture/storage.md):
конвенции имён, раскладка по шардам, приём событий, карта таблиц.
- [docs/research/](docs/research/) — исследования; сейчас это формат - [docs/research/](docs/research/) — исследования; сейчас это формат
кликстрима Яндекса, по которому строится модель события. кликстрима Яндекса, по которому строится модель события.
- [docs/formats/](docs/formats/) — описания форматов источников: по ним - [docs/formats/](docs/formats/) — описания форматов источников: по ним
+51 -12
View File
@@ -5,8 +5,12 @@
## Решение ## Решение
Топик `hits` читает одна Kafka-таблица формата `RawBLOB`: сообщение ложится в Топик `hits` читает одна Kafka-таблица формата `RawBLOB`: сообщение ложится в
`stg.hits_raw_rep` строкой, как пришло, рядом с метаданными доставки. Ни `stg.hits_raw_dist` строкой, как пришло, рядом с метаданными доставки. Ни
типизации, ни проверки на этом шаге нет — слой сырья ничего не интерпретирует. типизации, ни проверки на этом шаге нет — слой сырья ничего не интерпретирует.
У формата есть следствие для DDL: он читает вход в одно значение и рассчитан на
таблицу с единственной колонкой, поэтому у чтеца она ровно одна — `raw`, а
метаданные доставки берутся только из виртуальных колонок и добавить к чтецу
что-либо своё нельзя.
Типизированный слой наполняют две матвью, привязанные к `stg.hits_raw_dist`. Типизированный слой наполняют две матвью, привязанные к `stg.hits_raw_dist`.
Поля достаются `JSONExtract`. В таблицу ошибок уходят три класса брака: Поля достаются `JSONExtract`. В таблицу ошибок уходят три класса брака:
@@ -16,8 +20,17 @@
на валидность: `isValidJSON('123')` возвращает единицу, скаляр — тоже законный на валидность: `isValidJSON('123')` возвращает единицу, скаляр — тоже законный
JSON. JSON.
Классы пересекаются: скаляр проваливает заодно и сверку ключей, потому что
`JSONExtractKeys` от него даёт пустой массив. Поэтому они проверяются по порядку,
а в колонку `error_class` пишется первый совпавший — `not_an_object`,
`keyset_mismatch`, `key_field_unparsed`. Приём тот же, что у `mismatch_class` в
витрине сверки: пересекающиеся классы плюс объявленный приоритет.
Присутствие полей целиком держит сверка набора ключей — одно сравнение Присутствие полей целиком держит сверка набора ключей — одно сравнение
отсортированного `JSONExtractKeys` с контрактным списком. Обязательны все сорок `arraySort(JSONExtractKeys(raw))` с контрактным списком, завёрнутым в тот же
`arraySort`. Обёртка с обеих сторон стоит ноль и снимает ошибку, которая иначе
увела бы в брак вообще всё: сорок семь CamelCase-имён, выписанных руками ровно в
байтовом порядке. Обязательны все сорок
семь полей: генератор шлёт их все в каждом событии, а «пусто» по контракту — семь полей: генератор шлёт их все в каждом событии, а «пусто» по контракту —
пустое значение, а не отсутствие ключа. Этим же закрыт критерий #43 про опечатку пустое значение, а не отсутствие ключа. Этим же закрыт критерий #43 про опечатку
в имени. в имени.
@@ -40,6 +53,16 @@ JSON.
'stream'` не используется тоже: при чтении байтами на входе нечему ломаться, и 'stream'` не используется тоже: при чтении байтами на входе нечему ломаться, и
ошибке разбора взяться неоткуда. ошибке разбора взяться неоткуда.
Отсюда ограничение на форму выражений разбора: они собираются только из функций,
которые не бросают исключений. `JSONExtract` и родственные возвращают значение по
умолчанию или NULL, но не падают, — и правило репозитория «грязные записи не
валят пайплайн» держится теперь именно на этом. Исключение в матвью не ошибка
формата, его не перехватит никакой режим Kafka-движка: вставка упадёт, офсеты не
закоммитятся, блок пойдёт читаться снова. Ломается это громко и чинится без
потерь — поправил матвью, потребление продолжилось с некоммиченного офсета, — но
пока не починено, топик стоит. Поэтому приведение `Nullable` к необнуляемому типу
и любая арифметика в этих выражениях живут за предикатом, который NULL уже отсёк.
Этим решение снимает ограничение, записанное в постановке #37: «ошибки разбора Этим решение снимает ограничение, записанное в постановке #37: «ошибки разбора
должны рождаться на шаге Kafka-движка». Оно ставилось как условие достижимости должны рождаться на шаге Kafka-движка». Оно ставилось как условие достижимости
критерия #43 про громкую ошибку на опечатку в имени поля — критерий достижим и критерия #43 про громкую ошибку на опечатку в имени поля — критерий достижим и
@@ -116,26 +139,42 @@ contract-тест из #43 сюда не дотягивается — он ср
сообщений, которые не разобрались, и оставляет их пустыми для разобранных. сообщений, которые не разобрались, и оставляет их пустыми для разобранных.
Отсюда весь довод о том, что типизированный чтец не может наполнить слой сырья. Отсюда весь довод о том, что типизированный чтец не может наполнить слой сырья.
Составные типы `Nullable` не оборачивает: `Nullable(Array)`, `Nullable(Map)` и Составные типы `Nullable` не оборачивает: `Nullable(Array)` и `Nullable(Map)` не
`Nullable(Tuple)` не поддерживаются, а `Nullable` внутри них — да. Отсюда поддерживаются, а `Nullable` внутри них — да. Отсюда слепота `Nullable`-разбора
слепота `Nullable`-разбора на двенадцати колонках и необходимость сверки ключей. на двенадцати колонках и необходимость сверки ключей. Оговорка про кортеж:
`Nullable` внутри кортежа разрешён, поэтому запасной вариант со свёрткой в `Nullable(Tuple)` в ClickHouse всё-таки есть, но за настройкой
именованный кортеж жив. `enable_nullable_tuple_type` и в статусе беты; на довод это не влияет — кортежей
в контракте нет. `Nullable` внутри кортежа разрешён без всяких флагов, поэтому
запасной вариант со свёрткой в именованный кортеж жив.
Оттуда же: у Kafka-движка есть третий режим обработки ошибок — Оттуда же: у Kafka-движка есть третий режим обработки ошибок —
`dead_letter_queue` с записью в системную таблицу. Он отвергнут независимо от `dead_letter_queue` с записью в системную таблицу. Он отвергнут независимо от
остального: спека требует свои `*_errors`, системная таблица их не заменяет. остального: спека требует свои `*_errors`, системная таблица их не заменяет.
Форматы `RawBLOB` и `LineAsString` в ClickHouse есть, и Kafka-движок Форматы `RawBLOB` и `LineAsString` в ClickHouse есть, и Kafka-движок
поддерживает все форматы. поддерживает все форматы. `RawBLOB` читает вход в одно значение и рассчитан на
таблицу с единственным полем `String` — отсюда ограничение на форму чтеца.
Матвью с источником-`Distributed` срабатывает на вставку именно в эту Матвью с источником-`Distributed` срабатывает на вставку именно в эту
распределённую таблицу — блок она видит до разрезания по шардам. Проверено распределённую таблицу — блок она видит до разрезания по шардам. Проверено
владельцем на рабочих проектах; на стенде подтверждается заодно с приёмкой #37. владельцем на рабочих проектах; на стенде подтверждается заодно с приёмкой #37.
На живом стенде проверяется при исполнении #37, и то же внесено в раздел 11 На живом стенде проверяется при исполнении #37. Первые два пункта внесены в
спеки: раздел 11 спеки как несущие; остальные — однострочные `SELECT`, их довольно
прогнать заодно:
- `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение; - `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение. Проверять это
нужно первым и до написания DDL: формулировка «читает вход в одно значение»
про файл понятна, а про пачку сообщений из топика — нет, и если сообщения
склеятся, переделывать придётся решение целиком, а не DDL. Опыт стоит трёх
сообщений и одного `count()`. Запасной вариант — `LineAsString`: он режет по
переводу строки, а события у нас однострочные; цена запасного — сообщение с
переводом строки внутри даст две строки вместо одной;
- форма именованного кортежа в `JSONExtract` с `Nullable`-членами — нужна для - форма именованного кортежа в `JSONExtract` с `Nullable`-членами — нужна для
свёртки сорока семи вызовов в один, если разбор окажется дорогим. свёртки сорока семи вызовов в один, если разбор окажется дорогим;
- `isValidJSON('123')` возвращает единицу, а `JSONExtractKeys` от скаляра —
пустой массив. На обоих стоят классы брака и их приоритет, а документация
поведение на не-объекте не описывает: два `SELECT` закрывают вопрос;
- `JSONAsString` действительно падает на некорректном JSON, а не пропускает
строку. На этом стоит отказ от него в пользу `RawBLOB`; документация про
ошибочный ввод молчит.
+10 -3
View File
@@ -29,12 +29,14 @@
половине кластера, выглядит как честное число. При суффиксе вида голого имени половине кластера, выглядит как честное число. При суффиксе вида голого имени
нет вовсе, и та же ошибка становится громкой. нет вовсе, и та же ошибка становится громкой.
Второй довод — что группируется в `SHOW TABLES`. Суффикс группирует объекты по Второй довод — что группируется в списке по алфавиту. Суффикс группирует объекты
сущности: все четыре объекта топика `hits` стоят рядом, потому что различаются по сущности: все четыре объекта топика `hits` стоят рядом, потому что различаются
хвостом. Префиксный стиль, которым пользовался стенд-предшественник (`kafka_*`, хвостом. Префиксный стиль, которым пользовался стенд-предшественник (`kafka_*`,
`mv_*`, `v_*`), группирует по технологии, и объекты одной сущности `mv_*`, `v_*`), группирует по технологии, и объекты одной сущности
расползаются по алфавиту. В хранилище, где у одной сущности живёт по три-четыре расползаются по алфавиту. В хранилище, где у одной сущности живёт по три-четыре
воплощения, полезнее первое. воплощения, полезнее первое. Довод про дерево в клиенте и про `ORDER BY name`:
порядок выдачи `SHOW TABLES` документация не оговаривает, так что на него здесь
опираться нельзя.
Третий — преемственность: `_rep` и `_dist` уже используются владельцем в других Третий — преемственность: `_rep` и `_dist` уже используются владельцем в других
хранилищах на ClickHouse, и общий словарь между стендами стоит больше, чем хранилищах на ClickHouse, и общий словарь между стендами стоит больше, чем
@@ -60,3 +62,8 @@
объектов слоёв STG, ODS, DDS и DM, включая пары локальная/распределённая, объектов слоёв STG, ODS, DDS и DM, включая пары локальная/распределённая,
представления и матвью. Единственным объектом без выводимого суффикса оказался представления и матвью. Единственным объектом без выводимого суффикса оказался
словарь `products` — отсюда исключение в решении. словарь `products` — отсюда исключение в решении.
Первый прогон был неполным: четыре витрины из восьми остались с префиксом, и
заметило это холодное ревью, а не автор. Имена приведены в порядок 5 августа
2026 года. Урок не про имена: «прогнал по документу» — такое же утверждение,
как утверждение о поведении системы, и проверять его надо так же.
+179 -35
View File
@@ -40,9 +40,17 @@ keeper, Kafka, каркас сервисов. Объекты хранилища
возвращает данные одного шарда. Правило спеки «проверки и контрольные суммы — возвращает данные одного шарда. Правило спеки «проверки и контрольные суммы —
только по `Distributed`» такую тишину переживает плохо. только по `Distributed`» такую тишину переживает плохо.
Второе свойство — `SHOW TABLES` группирует объекты по сущности, а не по Второе свойство — в списке по алфавиту объекты группируются по сущности, а не по
технологии: все четыре объекта топика `hits` стоят рядом, потому что различаются технологии: все четыре объекта топика `hits` стоят рядом, потому что различаются
хвостом, а не началом имени. хвостом, а не началом имени. Речь про дерево в клиенте и про `ORDER BY name`:
порядок выдачи `SHOW TABLES` документация не оговаривает.
Начало имени — сущность, и берётся она в разных слоях из разных мест. В STG имя
приходит от транспорта: слой хранит то, что доехало по топику, и зовётся именем
топика — `hits_raw` от `hits`, `orders_raw` от `orders`. В типизированных слоях
имя приходит от предметной области и стоит в единственном числе: `event`,
`order_snapshot`, `session`. Граница между «как привезли» и «что это такое»
проходит по STG, и имена её показывают.
## Раскладка по шардам ## Раскладка по шардам
@@ -65,6 +73,18 @@ keeper, Kafka, каркас сервисов. Объекты хранилища
Пошардовые счётчики STG и ODS поэтому не сходятся и сходиться не должны — Пошардовые счётчики STG и ODS поэтому не сходятся и сходиться не должны —
сверять слои можно только через `_dist`. сверять слои можно только через `_dist`.
Ключи ко-локации названы заранее, потому что на них стоит политика соединений из
раздела 6 спеки: обычное соединение разрешено только по ключу ко-локации, всё
прочее — через `GLOBAL`. Значит `dds.session` и `dds.identity_map` шардируются по
`cityHash64(ClientID)`, а `dds.order` и производные от заказа — по
`cityHash64(order_id)`. Ключи витрин появятся вместе с самими витринами.
Известное ограничение правила «пишем только в `_dist`»: пакетные слои собираются
заменой дневных партиций, а `DROP/REPLACE PARTITION` работает только по локальным
таблицам. Чем и как раскладывать партицию-донор по шардам до замены, здесь не
решено — вопрос встаёт вместе со сборкой DDS, и решать его нужно тогда, а не
задним числом.
## Служебные колонки ## Служебные колонки
Собственные колонки не повторяют имён виртуальных. Виртуальные даёт движок: Собственные колонки не повторяют имён виртуальных. Виртуальные даёт движок:
@@ -76,14 +96,40 @@ keeper, Kafka, каркас сервисов. Объекты хранилища
`kafka_timestamp`, рядом — `consumer_host`, имя читавшей ноды: виртуальные `kafka_timestamp`, рядом — `consumer_host`, имя читавшей ноды: виртуальные
колонки его не несут, а после записи в `Distributed` он уже невосстановим. колонки его не несут, а после записи в `Distributed` он уже невосстановим.
Типы у них такие: `kafka_topic` и `consumer_host``LowCardinality(String)`,
значений мало и они повторяются; `kafka_partition` и `kafka_offset``UInt64`.
С `kafka_timestamp` сложнее: виртуальная колонка `_timestamp` заполнена не
всегда, а разрядность у секундной и миллисекундной версий разная. Поэтому колонка
объявляется `Nullable(DateTime)`, а точная форма проверяется на стенде
(раздел 11 спеки) — записать её в необнуляемый тип значит либо уронить приём на
первом сообщении, либо получить тихие нули за 1970 год.
Заполняются обе группы колонок выражением в `SELECT` матвью приёма, а не
`DEFAULT` в таблице. Для `consumer_host` это обязательно: `DEFAULT hostName()`
сработал бы на шарде-получателе и назвал бы не ту ноду, которая читала топик, —
то есть колонка молча отвечала бы на другой вопрос.
Само сообщение лежит в колонке `raw` тем, чем пришло: чтец читает байты и ничего Само сообщение лежит в колонке `raw` тем, чем пришло: чтец читает байты и ничего
не проверяет, поэтому там оказываются и целые события, и мусор. Разбирается всё не проверяет, поэтому там оказываются и целые события, и мусор. Разбирается всё
это ниже, в матвью ODS — см. [ADR 0005](../adr/0005-event-ingestion.md). это ниже, в матвью ODS — см. [ADR 0005](../adr/0005-event-ingestion.md).
Метка времени загрузки зовётся `_load_ts`, тип `DateTime64(3)`. В ODS она же Движок таблицы сырья — обычный `ReplicatedMergeTree`, `ORDER BY (kafka_partition,
служит колонкой версии `ReplacingMergeTree`. Имя согласовано с каноном служебных kafka_offset)`: разбор полётов идёт от «какое сообщение», другого ключа у сырья и
полей соседнего учебного стенда на Greenplum, чтобы словарь был общим у двух нет. Замену версий сюда ставить нельзя — она отменила бы свойство слоя, ради
хранилищ; ведущее подчёркивание у технических колонок — распространённая запись, которого он заведён: повтор доставки в сырье обязан быть виден.
Метка времени загрузки зовётся `_load_ts`, тип `DateTime64(3)`. Ставится она
один раз, в матвью приёма, и дальше переносится из STG в ODS как есть: колонка
отвечает на вопрос «когда строка приехала в хранилище», а не «когда её
разобрали». В ODS она же служит колонкой версии `ReplacingMergeTree`, и работа у
этой версии ровно одна — схлопнуть повтор доставки. Содержимое у повтора то же
самое, отличается только метка, поэтому какая из двух строк переживёт мерж,
безразлично. Переобработки как стадии у ODS нет вовсе: слой наполняет матвью, а
не пакетное задание, и пакетная работа с партициями начинается выше.
Имя согласовано с каноном служебных полей соседнего учебного стенда на
Greenplum, чтобы словарь был общим у двух хранилищ; ведущее подчёркивание у
технических колонок — распространённая запись,
её же используют Fivetran, Airbyte и Stitch. С правилом выше это не спорит: её же используют Fivetran, Airbyte и Stitch. С правилом выше это не спорит:
запрещено совпадать с именами виртуальных колонок, а не носить подчёркивание. запрещено совпадать с именами виртуальных колонок, а не носить подчёркивание.
@@ -106,6 +152,13 @@ keeper, Kafka, каркас сервисов. Объекты хранилища
Цепочка одна: чтец топика → матвью → сырьё STG → матвью разбора → событие и Цепочка одна: чтец топика → матвью → сырьё STG → матвью разбора → событие и
таблица ошибок ODS. таблица ошибок ODS.
Чтец стоит на обеих нодах и читает одной группой потребителей — имя группы
`clickstream_hits`, и оно одинаково на обеих нодах по построению: DDL идёт
`ON CLUSTER` и макросов в имени не содержит. Разные группы дали бы каждой ноде
полную копию топика, и это отдельная сцена для лабы, а не рабочий режим. Имя
кластера в `ON CLUSTER` и в движке `Distributed``clickstream_cluster`, оно
задано в `infra/clickhouse/config.d/cluster.xml`.
**Источник матвью разбора — `stg.hits_raw_dist`, а не локальная таблица.** **Источник матвью разбора — `stg.hits_raw_dist`, а не локальная таблица.**
У матвью две привязки: источник, на вставку в который она срабатывает, и цель, У матвью две привязки: источник, на вставку в который она срабатывает, и цель,
куда пишет. Распределённая таблица — лицо слоя, локальная — его хранилище; куда пишет. Распределённая таблица — лицо слоя, локальная — его хранилище;
@@ -113,14 +166,24 @@ keeper, Kafka, каркас сервисов. Объекты хранилища
той же ноде, что читала Kafka, в момент вставки первой матвью — до раскладки по той же ноде, что читала Kafka, в момент вставки первой матвью — до раскладки по
шардам. шардам.
**Вставка синхронная.** На пути приёма стоит **Вставка фоновая, и окно потери мы принимаем.** Вставка в распределённую
`distributed_foreground_insert = 1`. По умолчанию вставка в распределённую
таблицу кладёт блок в локальный спул и сразу возвращает управление, а Kafka таблицу кладёт блок в локальный спул и сразу возвращает управление, а Kafka
коммитит офсеты по факту работы матвью — то есть по факту записи в спул. Топик коммитит офсеты по факту работы матвью — то есть по факту записи в спул. Топик
уже считает сообщение прочитанным, хотя на шарде его нет. Синхронный режим не уже считает сообщение прочитанным, хотя на шарде его ещё нет: умри нода в этом
добавляет работы, он её не прячет: кусок на шарде будет записан всё равно, промежутке — сообщения не перечитаются.
вопрос лишь в том, ждём ли мы этого внутри вставки. Платится ожидание один раз
на блок Kafka в десятки тысяч строк, а не на событие. Закрывает окно настройка `distributed_foreground_insert = 1`, и на ETL-вставках
Airflow она стоит — там это обычный `SETTINGS` у запроса. На пути приёма её нет,
и по трём причинам. Вставку выполняет фоновый поток Kafka-движка, своего запроса
у него не бывает, так что настройка уровня запроса доехала бы только профилем
пользователя в конфигурации ноды. Синхронный режим связывает шарды: пока второй
недоступен, вставка падает, офсеты не коммитятся, и приём встаёт целиком — тогда
как при фоновом первая нода продолжает принимать и копит спул для соседа.
Платится при этом не одно ожидание на блок, а три распределённые вставки — сырьё,
событие, ошибки, — и все внутри потока-потребителя, что само по себе повод для
ребаланса по таймауту сессии. Против всего этого — окно в сотню миллисекунд на
ноутбуке, где мир пересобирается одной командой. Размен не в пользу настройки, а
компромисс полезнее показать, чем спрятать за галочкой.
**Гарантия — «хотя бы один раз», не транзакция.** Падение после записи на шард, **Гарантия — «хотя бы один раз», не транзакция.** Падение после записи на шард,
но до коммита офсетов даёт повтор при перечитывании. В ODS повтор схлопнет но до коммита офсетов даёт повтор при перечитывании. В ODS повтор схлопнет
@@ -129,38 +192,52 @@ keeper, Kafka, каркас сервисов. Объекты хранилища
Это свойство слоя, а не поломка, — но обещание идемпотентности конвейера к Это свойство слоя, а не поломка, — но обещание идемпотентности конвейера к
сырому слою не относится. сырому слою не относится.
Оговорка к последнему: у семейства `Replicated*` есть своя дедупликация — блок с
тем же хешем, вставленный повторно, отбрасывается (`insert_deduplicate`).
Удвоение сырья проходит мимо неё только потому, что при повторном чтении
`_load_ts` новый и хеш блока другой. Свойство слоя держится на этом, а не на
отсутствии механизма.
**Матвью разбора две, и их условия обязаны делить поток без зазора и без **Матвью разбора две, и их условия обязаны делить поток без зазора и без
нахлёста.** Одна забирает годные строки в `ods.event_dist`, вторая — брак в нахлёста.** Одна забирает годные строки в `ods.event_dist`, вторая — брак в
`ods.event_errors_dist`. Строка, подошедшая обеим, задвоится; не подошедшая ни `ods.event_errors_dist`. Строка, подошедшая обеим, задвоится; не подошедшая ни
одной — исчезнет молча. Держится это формой: второе условие пишется буквальным одной — исчезнет молча. Держится это формой: второе условие пишется буквальным
отрицанием первого, а сам предикат собирается только из функций, не возвращающих отрицанием первого, а сам предикат собирается только из функций, не возвращающих
NULL, — иначе трёхзначная логика даст строку, которую не возьмёт ни `условие`, NULL, — иначе трёхзначная логика даст строку, которую не возьмёт ни `условие`,
ни `NOT условие`. ни `NOT условие`. Те же функции не должны и бросать исключений: упавшая матвью
роняет вставку и останавливает потребление до починки
([ADR 0005](../adr/0005-event-ingestion.md)).
Постоянной сверки счётчиков при этом нет и не должно быть. У сырья срок жизни Постоянной сверки счётчиков при этом нет и не должно быть. У сырья срок жизни
трое суток, а ODS хранит всё, поэтому равенство «сырьё = события + ошибки» трое суток, а ODS хранит всё, поэтому равенство «сырьё = события + ошибки»
разъедется само; переобработка добавит строк в ODS, перезаливка задвоит сырьё, а разъедется само: сырьё истечёт раньше, а перезаливка модельного дня задвоит его,
ODS её схлопнет. Правило, красное в норме, учит не смотреть на оповещения тогда как в ODS тот же повтор схлопнется. Правило, красное в норме, учит не
смотреть на оповещения
([ADR 0002](../adr/0002-monitoring-scope.md)). Равенство проверяется разово в ([ADR 0002](../adr/0002-monitoring-scope.md)). Равенство проверяется разово в
smoke на управляемой пачке: отправили N сообщений — получили N строк сырья и N в smoke на управляемой пачке: отправили N сообщений — получили N строк сырья и N в
сумме событий и ошибок. сумме событий и ошибок. Счёт по ODS идёт через `FINAL`: голый `count()` по
`ReplacingMergeTree` зависит от того, сколько мержей успело пройти, и спека это
прямо запрещает (раздел 6).
## Срок жизни сырья ## Срок жизни сырья
Сырьё в STG живёт трое суток реального времени и уходит само. Трое — это окно Сырьё в STG живёт трое суток реального времени и уходит само. Трое — это окно
отладки: столько сырьё лежит на ноутбуке, чтобы менти успел разобрать полёты, отладки: столько сырьё лежит на ноутбуке, чтобы менти успел разобрать полёты,
после чего перестаёт занимать место. Нарезка — по дню загрузки, срок — по той же после чего перестаёт занимать место. Нарезка — по дню загрузки, срок — по той же
колонке `_load_ts`, снятие — целыми кусками (`ttl_only_drop_parts`). Значение колонке `_load_ts`, снятие — целыми кусками (`ttl_only_drop_parts`). В DDL
этой настройки между версиями ClickHouse менялось, поэтому в DDL оно проставлено значение проставлено явно, чтобы поведение не зависело от умолчания версии.
явно.
Нарезать сырьё по модельному дню события было бы соблазнительно — он единица Нарезать сырьё по модельному дню события было бы соблазнительно — он единица
переобработки и он же ключ партиции в ODS, — но чистку это ломает. Кусок переобработки и он же ключ партиции в ODS, — но чистку это ломает. Кусок
снимается целиком, только когда в нём истекли все строки, а обычный мерж внутри снимается целиком, только когда в нём истекли все строки, а в партиции модельного
партиции модельного дня склеит куски разного возраста, и данные переживут срок. дня лежит приехавшее в разное реальное время: мерж склеит куски разного возраста,
В партиции дня загрузки склеивать нечего: все строки в ней одного возраста, и самая свежая строка удержит весь кусок, и данные переживут срок неограниченно.
куски уходят по мере того, как истекает самый свежий из них — то есть сырьё
живёт трое суток с небольшим хвостом, а не ровно трое. В партиции дня загрузки склейка идёт точно так же, и строки в ней тоже разного
возраста — но не более чем на сутки, потому что партицию закрывает календарный
день. Отсюда и оценка: сырьё живёт трое суток плюс хвост до суток, а не ровно
трое. Оговорка про сроки: TTL исполняется на мержах, а не по будильнику, так что
«уходит само» здесь обещано, а «уходит вовремя» — нет.
Модельного дня среди колонок сырья нет вовсе. Он свойство содержимого, а Модельного дня среди колонок сырья нет вовсе. Он свойство содержимого, а
содержимое разбирает ODS — там `EventDate` и живёт, ключом партиции. Сырьё содержимое разбирает ODS — там `EventDate` и живёт, ключом партиции. Сырьё
@@ -182,32 +259,60 @@ D0 и к реальному календарю не привязана; паке
## Таблица ошибок ## Таблица ошибок
`ods.event_errors` держит строки, не прошедшие строгий приём, вместе с их сырым `ods.event_errors` держит строки, не прошедшие строгий приём, вместе с их сырым
текстом и метаданными доставки. Ключи её собственные, потому что у брака нет текстом, метаданными доставки и классом брака. Ключи её собственные, потому что у
разобранных полей: шардируется `cityHash64` сырой строки — `ClientID` у строки, брака нет разобранных полей: шардируется `cityHash64` сырой строки — `ClientID` у
которая не разобралась, взять неоткуда; нарезается по дню загрузки, как и сырьё; строки, которая не разобралась, взять неоткуда; нарезается по дню загрузки, как и
живёт месяц. Дольше сырья — намеренно: если брак истекает вместе с ним, сырьё; живёт месяц. Дольше сырья — намеренно: если брак истекает вместе с ним,
разбираться к моменту разбирательства будет уже нечем. разбираться к моменту разбирательства будет уже нечем.
Класс брака лежит в колонке `error_class` типа `LowCardinality(String)`. Без неё
в таблице копятся строки «что-то не так» без ответа на «что именно», а витрине
качества не на что опереться. Сами классы перечислены в
[ADR 0005](../adr/0005-event-ingestion.md) и проверяются по порядку, потому что
пересекаются: сообщение, не являющееся объектом JSON, проваливает заодно и сверку
набора ключей — `JSONExtractKeys` от скаляра даёт пустой массив. Побеждает первый
совпавший класс, тем же приёмом, что `mismatch_class` в витрине сверки.
Движок — обычный `ReplicatedMergeTree`, без замены версий: схлопывать брак не по
чему, у него нет ключа сущности. `ORDER BY` — `(error_class, kafka_partition,
kafka_offset)`: смотрят такую таблицу от класса, а внутри класса — по координатам
доставки.
## Раскладка DDL ## Раскладка DDL
Файлы лежат в `sql/ddl/` и применяются по порядку имён. Один файл — это слой и Файлы лежат в `sql/ddl/` и применяются по порядку имён. Сначала все статичные
роль: статичные объекты отдельно от матвью. объекты, потом матвью — тогда к моменту создания матвью её цель уже существует.
| Файл | Что в нём | | Файл | Что в нём |
|---|---| |---|---|
| `00-databases.sql` | базы слоёв | | `00-databases.sql` | базы слоёв |
| `10-stg-tables.sql` | Kafka-таблица, локальная и распределённая таблицы сырья | | `10-stg-tables.sql` | Kafka-таблица, локальная и распределённая таблицы сырья |
| `11-stg-views.sql` | матвью, наполняющая сырьё из Kafka |
| `20-ods-tables.sql` | типизированное событие и таблица ошибок | | `20-ods-tables.sql` | типизированное событие и таблица ошибок |
| `21-ods-views.sql` | матвью разбора: сырьё в событие и в ошибки | | `30-ods-views.sql` | матвью разбора: сырьё в событие и в ошибки |
| `40-stg-views.sql` | матвью приёма: чтец в сырьё |
Матвью принадлежит слою своей цели, а не источника: разбор из STG в ODS лежит Порядок задают два правила. Первое: матвью принадлежит слою своей цели, а не
среди файлов ODS, потому что наполняет ODS. источника, — разбор из STG в ODS лежит среди файлов ODS, потому что наполняет
ODS. Второе: матвью приёма создаётся последней из всех, и потому нарушает
нумерацию слоёв. Kafka-движок начинает читать топик ровно тогда, когда к нему
привязывают первую матвью; создай её раньше разбора — и всё, что доедет в
зазоре, ляжет в сырьё и не попадёт в ODS никуда, ни в событие, ни в ошибки. На
пустом топике зазор безвреден, поэтому первый прогон о нём не скажет. Проснётся
он, когда тома ClickHouse снесены, а данные Kafka целы, — то есть на обычной
отладке.
Применение — двумя одноразовыми сервисами при `make up`, по образцу уже Применение — двумя одноразовыми сервисами при `make up`, по образцу уже
работающих `airflow-init` и `superset-init`. Сначала `kafka-init` создаёт топик работающих `airflow-init` и `superset-init`. Сначала `kafka-init` создаёт топик
`hits` с двумя партициями, затем `clickhouse-init` дожидается его завершения и `hits` с двумя партициями, затем `clickhouse-init` дожидается его завершения и
применяет файлы с ноды 1, `ON CLUSTER`. Порядок страхует от автосоздания топика применяет файлы с ноды 1, `ON CLUSTER`.
Образцы копируются не целиком, и в двух местах. `clickhouse-init` обязан ждать
готовности **обеих** нод: `ON CLUSTER` ждёт исполнения на всех хостах и по
таймауту бросает, а `airflow-init` ждёт только первую ноду, `superset-init`
только вторую. И второе: оба образца переживают `make up --wait` лишь потому, что
от них зависят долгоживущие сервисы; у пары `kafka-init` / `clickhouse-init`
таких зависимых нет, и как поведёт себя `--wait` с одноразовым сервисом без них —
проверяется при исполнении #37. Порядок страхует от автосоздания топика
брокером с одной партицией: у потребителя librdkafka разрешение на автосоздание брокером с одной партицией: у потребителя librdkafka разрешение на автосоздание
по умолчанию выключено, так что случиться это не обязано, но урок «обе ноды по умолчанию выключено, так что случиться это не обязано, но урок «обе ноды
читают топик» умирает тихо, и полагаться на умолчание клиента здесь не стоит. читают топик» умирает тихо, и полагаться на умолчание клиента здесь не стоит.
@@ -234,3 +339,42 @@ D0 и к реальному календарю не привязана; паке
Слои DDS и DM появляются на следующих этапах; их состав задан разделом 7 Слои DDS и DM появляются на следующих этапах; их состав задан разделом 7
мастер-спеки и переносится сюда по мере постройки. мастер-спеки и переносится сюда по мере постройки.
## Что проверено
Документ описывает устройство, которого в репозитории ещё нет, и на каждом шагу
опирается на поведение ClickHouse. Поэтому утверждения о движке разведены на три
группы: насколько фразе можно верить, должно быть видно из текста, а не зависеть
от того, хорошо ли автор помнит документацию. Сверка — через MCP Context7,
5 августа 2026 года; то же разведение для механики приёма — в
[ADR 0005](../adr/0005-event-ingestion.md).
**Сверено с документацией.** Собственная колонка с именем виртуальной делает
виртуальную недоступной. При вставке в `Distributed` шард выбирается по ключу
шардирования; фоновый режим — умолчание, а `distributed_foreground_insert = 1`
завершает вставку только после записи на все шарды. У семейства `Replicated*`
есть дедупликация одинаковых блоков. Голый `count()` по `ReplacingMergeTree`
зависит от того, сколько мержей прошло. `ttl_only_drop_parts` снимает кусок
целиком и только когда истекли все строки в нём, а сам TTL исполняется на
фоновых мержах. Kafka-движок начинает читать топик, когда к нему привязывают
матвью, и одна группа потребителей на кластер спасает от дублей. Таблицы с
одинаковым путём в keeper становятся репликами друг друга, а макрос `{uuid}`
завязан на движок базы `Atomic`. `ON CLUSTER` ждёт все хосты и бросает по
таймауту; `CREATE ... IF NOT EXISTS` на существующем объекте не бросает.
**Записано как проверка на стенде** — раздел 11 спеки и ADR 0005. Срабатывание
матвью с источником-`Distributed` на вставку именно в неё: этого случая в
документации нет вовсе, утверждение держится на опыте владельца. Одна строка на
сообщение у `RawBLOB`. Обнуляемость и разрядность виртуальной колонки
`_timestamp`.
**Сказано по памяти, проверки пока нет.** Что `DROP/REPLACE PARTITION` не
работает по `Distributed` — прямого запрета в документации нет, все примеры даны
для семейства MergeTree. Что `DEFAULT hostName()` вычислился бы на
шарде-получателе, а не на вставляющей ноде, и что имя читавшей ноды после записи
в `Distributed` уже невосстановимо. Что у потребителя librdkafka автосоздание
топиков по умолчанию выключено. Что упавшая матвью роняет вставку и
останавливает потребление до починки — на этой фразе держится правило «грязные
записи не валят пайплайн», и стоит она пока на одном рассуждении. Проверяются
все пятеро дёшево и заодно с приёмкой #37; до тех пор это предположения, а не
знание.
+35 -15
View File
@@ -312,9 +312,13 @@ CSV в репозитории (`data/catalog/products.csv`: `sku`, `name`, `cate
`cityHash64(order_id)`, сырьё STG — `cityHash64(сырой строки)`; полный `cityHash64(order_id)`, сырьё STG — `cityHash64(сырой строки)`; полный
список и доводы — в [доке хранилища](../architecture/storage.md). Урок: список и доводы — в [доке хранилища](../architecture/storage.md). Урок:
«какая нода читала топик — меняется между прогонами, куда легли данные — нет». «какая нода читала топик — меняется между прогонами, куда легли данные — нет».
- **Приём строгий**: обязательные поля разбираются как `Nullable`, а набор - **Приём строгий**: пять опорных колонок — `WatchID`, `VisitID`, `ClientID`,
ключей сообщения сверяется с контрактным; строка с NULL среди обязательных `EventDate`, `UTCEventTime` — разбираются как `Nullable`, а набор ключей
полей или с разошедшимся набором ключей уходит в `*_errors`. Сверка ключей — сообщения сверяется с контрактным; строка с NULL среди опорных колонок
или с разошедшимся набором ключей уходит в `*_errors`. Опорными выбраны те, на
которых стоят ключ сортировки, партиция и дедупликация; остальные сорок две
достаются обычными типами — сорок семь проверок на NULL превратили бы матвью в
простыню, а присутствие и так целиком закрыто сверкой ключей. Сверка ключей —
не добавка: у массивов NULL не бывает, и пропавшее поле-массив иначе не добавка: у массивов NULL не бывает, и пропавшее поле-массив иначе
неотличимо от пустого по смыслу. На входе разбора нет вовсе — Kafka-таблица неотличимо от пустого по смыслу. На входе разбора нет вовсе — Kafka-таблица
читает сообщение байтами, строгость целиком в матвью ODS читает сообщение байтами, строгость целиком в матвью ODS
@@ -331,8 +335,9 @@ CSV в репозитории (`data/catalog/products.csv`: `sku`, `name`, `cate
легитимная GLOBAL-витрина (заказы малы). легитимная GLOBAL-витрина (заказы малы).
- **Конвейер без TRUNCATE**: поток — append-only в ReplacingMergeTree (дедуп - **Конвейер без TRUNCATE**: поток — append-only в ReplacingMergeTree (дедуп
через argMax); батчевая переобработка — по дневным партициям через argMax); батчевая переобработка — по дневным партициям
(`DROP/REPLACE PARTITION ON CLUSTER`); `TRUNCATE ... ON CLUSTER` остаётся (`DROP/REPLACE PARTITION ON CLUSTER`); `TRUNCATE ... ON CLUSTER` в конвейере не
только в `make reset`. `DROP/REPLACE PARTITION` работает только по применяется вовсе — полный сброс стенда делается `make clean && make up`, то
есть вместе с томами. `DROP/REPLACE PARTITION` работает только по
**локальным** таблицам ON CLUSTER, не по Distributed; замена через **локальным** таблицам ON CLUSTER, не по Distributed; замена через
DROP+INSERT неатомарна — дашборд в середине прогона честно моргает (это DROP+INSERT неатомарна — дашборд в середине прогона честно моргает (это
осознанная цена, не баг). осознанная цена, не баг).
@@ -374,7 +379,7 @@ README.
| Слой | Объект | Что это | | Слой | Объект | Что это |
|---|---|---| |---|---|---|
| Kafka | `hits`, `orders` | два топика, по 2 партиции | | Kafka | `hits`, `orders` | два топика, по 2 партиции |
| STG | `stg.hits_raw_kafka`, `stg.hits_raw` + MV; то же для orders | сырые строки, Kafka Engine на обеих нодах | | STG | `stg.hits_raw_kafka`, `stg.hits_raw` + MV; для orders — развилка этапа 3, см. раздел 12 | сырые строки, Kafka Engine на обеих нодах |
| ODS | `ods.event` (+`_errors`) | типизированное широкое событие, ReplacingMergeTree | | ODS | `ods.event` (+`_errors`) | типизированное широкое событие, ReplacingMergeTree |
| ODS | `ods.order_snapshot` (+`_errors`) | слепки заказов как приехали, партиция по `snapshot_date`, без дедупа | | ODS | `ods.order_snapshot` (+`_errors`) | слепки заказов как приехали, партиция по `snapshot_date`, без дедупа |
| DDS | `dds.session` | сборка сессий из событий (наследник `dds.click`) | | DDS | `dds.session` | сборка сессий из событий (наследник `dds.click`) |
@@ -409,11 +414,11 @@ Kafka день переигрывается генератором заново:
### Витрины DM ### Витрины DM
- **`v_revenue_daily`** (выручка, только от заказов): `report_date`, - **`revenue_daily_v`** (выручка, только от заказов): `report_date`,
`product_category` (через `dictGet` каталога + ARRAY JOIN позиций), `product_category` (через `dictGet` каталога + ARRAY JOIN позиций),
`orders`, `units`, `revenue`, `aov`. Считается по заказам в статусе `orders`, `units`, `revenue`, `aov`. Считается по заказам в статусе
`paid`; внутри окна K число дня «дышит». `paid`; внутри окна K число дня «дышит».
- **`v_purchase_vs_orders`** (сверка): FULL OUTER GLOBAL JOIN по - **`purchase_vs_orders_v`** (сверка): FULL OUTER GLOBAL JOIN по
`purchaseID = order_id`; колонки: `order_day`, `order_id`, `purchaseID = order_id`; колонки: `order_day`, `order_id`,
`declared_revenue` (клиент), `items_total` (бэкенд, сравнимая база — не `declared_revenue` (клиент), `items_total` (бэкенд, сравнимая база — не
`total`: промокод и доставка клиенту не видны), `status`, `mismatch_class` `total`: промокод и доставка клиенту не видны), `status`, `mismatch_class`
@@ -427,10 +432,10 @@ Kafka день переигрывается генератором заново:
(шестое значение `mismatch_class`, вне приоритетов расхождений). После (шестое значение `mismatch_class`, вне приоритетов расхождений). После
закрытия окна K таких строк не остаётся — сироты исключены построением закрытия окна K таких строк не остаётся — сироты исключены построением
(раздел 4). (раздел 4).
- **`v_utm_effectiveness`** — остаётся клиентской (атрибуция по трекеру); - **`utm_effectiveness_v`** — остаётся клиентской (атрибуция по трекеру);
счётчики `purchases`/`add_to_carts` оживают из таксономии, добавляется счётчики `purchases`/`add_to_carts` оживают из таксономии, добавляется
`declared_revenue` по UTM. `declared_revenue` по UTM.
- **`v_daily_traffic`** — расширяется парой «посетители» / «известные - **`daily_traffic_v`** — расширяется парой «посетители» / «известные
пользователи» (обогащение через `dds.identity_map`, локальное соединение пользователи» (обогащение через `dds.identity_map`, локальное соединение
по ключу ко-локации). по ключу ко-локации).
- `events_enriched_v`, `top_pages_daily_v`, `session_overview_v`, - `events_enriched_v`, `top_pages_daily_v`, `session_overview_v`,
@@ -498,7 +503,7 @@ v2 стартует пустым, поэтому объём ниже — это
`cancelled`) — на статусе `paid` стоит «дыхание» выручки; сужение до пары `cancelled`) — на статусе `paid` стоит «дыхание» выручки; сужение до пары
`created`/`cancelled` — запасной ход, если генератор заказов окажется дороже `created`/`cancelled` — запасной ход, если генератор заказов окажется дороже
ожиданий. Связка: если срез 1 сработает, определение выручки в ожиданий. Связка: если срез 1 сработает, определение выручки в
`v_revenue_daily` придётся сменить с «заказы в статусе `paid`» на «все `revenue_daily_v` придётся сменить с «заказы в статусе `paid`» на «все
неотменённые заказы». неотменённые заказы».
- **Отступление от порядка #15**: резолюция предписывала резать в порядке - **Отступление от порядка #15**: резолюция предписывала резать в порядке
1 → 2 → 3, спека применяет 3 и 2, а 1 держит в резерве. Довод: срезы 3 и 2 1 → 2 → 3, спека применяет 3 и 2, а 1 держит в резерве. Довод: срезы 3 и 2
@@ -569,8 +574,15 @@ smoke-проверки, а не «дашборд зелёный». Это мин
- Точная форма `ORDER BY` ODS-таблиц (выражение `intHash32` в ключе - Точная форма `ORDER BY` ODS-таблиц (выражение `intHash32` в ключе
ReplacingMergeTree). ReplacingMergeTree).
- `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение (ADR 0005). - `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение (ADR 0005).
- Матвью, привязанная к локальной таблице, срабатывает, когда строки приходят Проверять первым, до написания DDL: если сообщения склеятся, переделывать
вставкой через `Distributed`: на этом стоит цепочка STG → ODS (ADR 0005). придётся решение, а не запрос. Запасной формат — `LineAsString`.
- Тип виртуальной колонки `_timestamp` у Kafka-движка: обнуляемость и
разрядность (секунды против миллисекунд) — от этого зависит объявление
`kafka_timestamp` в таблице сырья.
- Матвью с источником-`Distributed` срабатывает на вставку именно в эту
распределённую таблицу, до раскладки по шардам: на этом стоит цепочка
STG → ODS (ADR 0005). Проверено владельцем на рабочих проектах, в документации
ClickHouse этот случай не описан.
- Форма именованного кортежа в `JSONExtract` с `Nullable`-членами — ею - Форма именованного кортежа в `JSONExtract` с `Nullable`-членами — ею
сворачиваются 47 вызовов в один, если разбор окажется дорогим (ADR 0005). сворачиваются 47 вызовов в один, если разбор окажется дорогим (ADR 0005).
- Размер артефакта эталонного мира после пересборки. - Размер артефакта эталонного мира после пересборки.
@@ -628,8 +640,16 @@ v2, этап 0).
`count()`, а в бою — нет; частый вопрос на собеседованиях; `count()`, а в бою — нет; частый вопрос на собеседованиях;
- два способа принять топик, рядом на одном стенде: сырьё байтами с разбором - два способа принять топик, рядом на одном стенде: сырьё байтами с разбором
функциями (`hits`, [ADR 0005](../adr/0005-event-ingestion.md)) против функциями (`hits`, [ADR 0005](../adr/0005-event-ingestion.md)) против
типизированного чтеца с `kafka_handle_error_mode` (заказы, этап 3) — типизированного чтеца с `kafka_handle_error_mode` — сравнение цены и
сравнение цены и наблюдаемости как задание; наблюдаемости как задание. **Развилка этапа 3, не решена**: типизированный
чтец идёт заказам только вместе с ответом на вопрос, нужен ли им слой сырья.
Нужен — и чтецов на один топик станет два, а эту схему ADR 0005 отверг; не
нужен — и слои перестают быть единообразными. Разбирать грилингом, когда
дойдём до заказов;
- матвью как рабочий механизм, а не диковина: их видно на приёме и на сборке
ODS, а пакетная работа начинается выше. Отдельным заданием — как читать из ODS
последние версии, через `FINAL` или оконной функцией: что нагляднее, решаем на
месте;
- лекция про идентичность «как в бою»: `setUserID` и first-party id, - лекция про идентичность «как в бою»: `setUserID` и first-party id,
детерминированная против вероятностной склейки, identity graph, детерминированная против вероятностной склейки, identity graph,
кросс-девайс, CDP — с рамкой «мы склеили через транзакции, потому что трекер кросс-девайс, CDP — с рамкой «мы склеили через транзакции, потому что трекер