docs(storage): конвенции хранилища и приём событий перед этапом 2 #52
@@ -125,3 +125,7 @@ Airflow) и названия из кода. Если для понятия ес
|
||||
- Имена файлов в `docs/adr/` — `NNNN-краткое-имя.md`: сквозной номер из четырёх
|
||||
цифр и слаг (`0001-stand-services.md`). Решения нумеруются подряд, дата
|
||||
в имени не нужна.
|
||||
- Имена файлов в `docs/architecture/` — слаг строчными латинскими буквами через
|
||||
дефис (`storage.md`). Здесь живут рабочие справочники по зонам
|
||||
ответственности: не событие истории и не решение, а текущее устройство —
|
||||
один файл на зону, правится по мере постройки.
|
||||
|
||||
@@ -204,6 +204,11 @@ Superset закреплён на 6.1.0; драйвер `clickhouse-connect`, ф
|
||||
## Документация
|
||||
|
||||
- [docs/specs/](docs/specs/) — спеки: источник истины о задуманном.
|
||||
- [docs/adr/](docs/adr/) — принятые решения с доводами и отвергнутыми
|
||||
вариантами: почему сделано так, а не иначе.
|
||||
- [docs/architecture/](docs/architecture/) — рабочие справочники по зонам
|
||||
ответственности; сейчас это [хранилище](docs/architecture/storage.md):
|
||||
конвенции имён, раскладка по шардам, приём событий, карта таблиц.
|
||||
- [docs/research/](docs/research/) — исследования; сейчас это формат
|
||||
кликстрима Яндекса, по которому строится модель события.
|
||||
- [docs/formats/](docs/formats/) — описания форматов источников: по ним
|
||||
|
||||
@@ -0,0 +1,184 @@
|
||||
# ADR 0005. Приём событий: сырьё в STG байтами, разбор функциями в ODS
|
||||
|
||||
Дата: 3 августа 2026 года. Статус: принято.
|
||||
|
||||
## Решение
|
||||
|
||||
Топик `hits` читает одна Kafka-таблица формата `RawBLOB`: сообщение ложится в
|
||||
`stg.hits_raw_dist` строкой, как пришло, рядом с метаданными доставки. Ни
|
||||
типизации, ни проверки на этом шаге нет — слой сырья ничего не интерпретирует.
|
||||
У формата есть следствие для DDL: он читает вход в одно значение и рассчитан на
|
||||
таблицу с единственной колонкой, поэтому у чтеца она ровно одна — `raw`, а
|
||||
метаданные доставки берутся только из виртуальных колонок и добавить к чтецу
|
||||
что-либо своё нельзя.
|
||||
|
||||
Типизированный слой наполняют две матвью, привязанные к `stg.hits_raw_dist`.
|
||||
Поля достаются `JSONExtract`. В таблицу ошибок уходят три класса брака:
|
||||
сообщение, не являющееся объектом JSON; объект, чей набор ключей разошёлся с
|
||||
контрактным; объект, у которого не разобрался ключевой идентификатор или метка
|
||||
времени. Остальное — в событие. Первый класс проверяется именно на объект, а не
|
||||
на валидность: `isValidJSON('123')` возвращает единицу, скаляр — тоже законный
|
||||
JSON.
|
||||
|
||||
Классы пересекаются: скаляр проваливает заодно и сверку ключей, потому что
|
||||
`JSONExtractKeys` от него даёт пустой массив. Поэтому они проверяются по порядку,
|
||||
а в колонку `error_class` пишется первый совпавший — `not_an_object`,
|
||||
`keyset_mismatch`, `key_field_unparsed`. Приём тот же, что у `mismatch_class` в
|
||||
витрине сверки: пересекающиеся классы плюс объявленный приоритет.
|
||||
|
||||
Присутствие полей целиком держит сверка набора ключей — одно сравнение
|
||||
`arraySort(JSONExtractKeys(raw))` с контрактным списком, завёрнутым в тот же
|
||||
`arraySort`. Обёртка с обеих сторон стоит ноль и снимает ошибку, которая иначе
|
||||
увела бы в брак вообще всё: сорок семь CamelCase-имён, выписанных руками ровно в
|
||||
байтовом порядке. Обязательны все сорок
|
||||
семь полей: генератор шлёт их все в каждом событии, а «пусто» по контракту —
|
||||
пустое значение, а не отсутствие ключа. Этим же закрыт критерий #43 про опечатку
|
||||
в имени.
|
||||
|
||||
Тип проверяется не у всех колонок, а у пяти: `WatchID`, `VisitID`, `ClientID`,
|
||||
`EventDate`, `UTCEventTime` разбираются в `Nullable` и дают NULL, если значение
|
||||
не той природы. Остальные сорок две достаются обычными типами. Соотношение
|
||||
цены и пользы: единственный производитель топика — собственный генератор,
|
||||
сериализующий из контракта по объявленным типам, поэтому неверный тип может
|
||||
прийти только из руки, а сорок семь проверок на NULL превратили бы матвью в
|
||||
простыню. Пять выбраны по последствию: это идентификаторы события, визита и
|
||||
посетителя, дата партиции и метка времени, по которой события упорядочиваются
|
||||
внутри сессии, — порча любой отравляет всё ниже по течению. `CounterID`
|
||||
формально тоже входит в ключ сортировки, но на стенде он константа, и NULL там
|
||||
взяться неоткуда.
|
||||
|
||||
Присутствие иначе и не проверить. `Nullable`
|
||||
в ClickHouse не оборачивает составные типы: `Nullable(Array)` запрещён, а
|
||||
`Array(Nullable(T))` при пропавшем ключе даёт пустой массив, неотличимый от
|
||||
пустого по смыслу. Таких колонок в контракте двенадцать из сорока семи.
|
||||
|
||||
Типизированная Kafka-таблица не используется. Режим `kafka_handle_error_mode =
|
||||
'stream'` не используется тоже: при чтении байтами на входе нечему ломаться, и
|
||||
ошибке разбора взяться неоткуда.
|
||||
|
||||
Отсюда ограничение на форму выражений разбора: они собираются только из функций,
|
||||
которые не бросают исключений. `JSONExtract` и родственные возвращают значение по
|
||||
умолчанию или NULL, но не падают, — и правило репозитория «грязные записи не
|
||||
валят пайплайн» держится теперь именно на этом. Исключение в матвью не ошибка
|
||||
формата, его не перехватит никакой режим Kafka-движка: вставка упадёт, офсеты не
|
||||
закоммитятся, блок пойдёт читаться снова. Ломается это громко и чинится без
|
||||
потерь — поправил матвью, потребление продолжилось с некоммиченного офсета, — но
|
||||
пока не починено, топик стоит. Поэтому приведение `Nullable` к необнуляемому типу
|
||||
и любая арифметика в этих выражениях живут за предикатом, который NULL уже отсёк.
|
||||
|
||||
Этим решение снимает ограничение, записанное в постановке #37: «ошибки разбора
|
||||
должны рождаться на шаге Kafka-движка». Оно ставилось как условие достижимости
|
||||
критерия #43 про громкую ошибку на опечатку в имени поля — критерий достижим и
|
||||
без него, средствами матвью.
|
||||
|
||||
Требование строгого приёма из раздела 6 спеки остаётся в силе, меняется его
|
||||
механизм: настройка формата `input_format_skip_unknown_fields` при чтении
|
||||
байтами беспредметна, раскладки по полям на этом шаге нет вовсе.
|
||||
|
||||
## Почему
|
||||
|
||||
Типизированный чтец не отдаёт сырьё. При `kafka_handle_error_mode = 'stream'`
|
||||
виртуальные колонки `_raw_message` и `_error` заполняются только при ошибке
|
||||
разбора и пусты у разобранных сообщений. Один типизированный чтец оставил бы в
|
||||
STG сырьё исключительно от брака: слой сырых строк хранил бы всё, кроме того,
|
||||
что доехало.
|
||||
|
||||
Слои идут цепочкой. Kafka → STG → ODS → DDS → DM, каждый читает предыдущий.
|
||||
Отвергнутая схема с двумя чтецами одного топика — сырым для STG и типизированным
|
||||
для ODS — в бою обычна и дала бы строгий приём даром, средствами самого движка.
|
||||
Но ODS перестал бы быть надстройкой над STG и стал бы вторым входом с шины, а на
|
||||
стенде, где слои и есть предмет изучения, это дороже сэкономленного. Побочная
|
||||
выгода скромнее, чем кажется на первый взгляд: слой сырья при двух чтецах
|
||||
существовал бы точно так же, но SQL разбора не существовал бы нигде, и для
|
||||
переобработки дня X его пришлось бы писать заново. В цепочке он уже есть в
|
||||
матвью, и ручная вставка получается его копией.
|
||||
|
||||
Слой сырья ничего не проверяет, и формат чтеца выбран под это. Соседний вариант,
|
||||
`JSONAsString`, требует, чтобы каждое сообщение было корректным объектом JSON, и
|
||||
на некорректном падает: потребление встаёт, хотя правило репозитория гласит, что
|
||||
грязные записи не должны валить пайплайн. Лечится это двумя способами — включить
|
||||
на чтеце режим обработки ошибок или не проверять на входе вовсе. Второе честнее:
|
||||
частичная проверка в слое, чья работа — не проверять, спорит сама с собой, а
|
||||
настоящий разбор всё равно идёт ниже. При `RawBLOB` ломаться нечему, доезжают
|
||||
любые байты, и оба класса брака разбираются в одном месте.
|
||||
|
||||
Архив нужен буквальный и читаемый глазами. `RawBLOB` — это про способ чтения, а
|
||||
не про тип колонки: на диске лежит обычный `String` с текстом события, и менти
|
||||
открывает его в обычном клиенте и читает. В этом и смысл слоя — увидеть, что
|
||||
реально пришло. Мусор при этом лежит в той же колонке, что и целые события, а не
|
||||
в отдельной. Отвергнутый боевой вариант — собрать слой сырья на движке `Null`,
|
||||
тогда данные не хранятся вовсе, а матвью работает чистым триггером. Отлаживать
|
||||
так неудобно, поэтому здесь окно в трое суток, а не ноль.
|
||||
|
||||
Что при этом потеряно и чем закрыто. Настройка `input_format_skip_unknown_fields
|
||||
= 0` ловила два случая: ожидаемое поле пропало и появилось лишнее, неизвестное.
|
||||
Оба берёт на себя сверка набора ключей: опечатка в имени — это одновременно
|
||||
пропавшее ожидаемое и появившееся лишнее, и сравнение множеств видит и то, и
|
||||
другое. Она же единственная защита у колонок-массивов и она же ловит молчаливое
|
||||
расширение контракта, когда в сообщении появляется поле, о котором хранилище не
|
||||
знает. `Nullable`-разбор к этой работе отношения не имеет — он про тип, а не про
|
||||
присутствие, и потому сужен до пяти колонок.
|
||||
|
||||
Цена решения. Контракт получает ещё два места: типы сорока семи колонок в
|
||||
выражениях матвью и список тех же имён для сверки ключей. Раздел 1.4 спеки
|
||||
предупреждает, что, повторяясь примерно в семи местах, они расходятся молча, а
|
||||
contract-тест из #43 сюда не дотягивается — он сравнивает `system.columns`
|
||||
целевой таблицы со схемой генератора и о выражениях матвью ничего не знает.
|
||||
Смягчение работает не везде: опечатка в имени скалярного поля уводит строки в
|
||||
таблицу ошибок пачкой и видна сразу, а опечатка в имени массива даёт пустой
|
||||
массив тихо — сверка ключей проверяет ключи сообщения, а не выражения матвью.
|
||||
Эти двенадцать колонок сторожит smoke: известное событие с товарами обязано
|
||||
доезжать с непустыми массивами.
|
||||
|
||||
Второе — разбор функциями дороже разбора форматом. На объёмах стенда это
|
||||
несущественно; если станет заметно, сорок семь вызовов сворачиваются в один
|
||||
`JSONExtract` в именованный кортеж, и строка разбирается однократно.
|
||||
|
||||
## Что проверено
|
||||
|
||||
По документации ClickHouse через MCP Context7: основная сверка — 3 августа
|
||||
2026 года, перепроверка после правок — 5 августа. Датировка важна: 5 августа
|
||||
утверждение про `Nullable(Tuple)` развернулось на противоположное.
|
||||
|
||||
При режиме `stream` движок отдаёт `_raw_message` и `_error` только для
|
||||
сообщений, которые не разобрались, и оставляет их пустыми для разобранных.
|
||||
Отсюда весь довод о том, что типизированный чтец не может наполнить слой сырья.
|
||||
|
||||
Составные типы `Nullable` не оборачивает: `Nullable(Array)` и `Nullable(Map)` не
|
||||
поддерживаются, а `Nullable` внутри них — да. Отсюда слепота `Nullable`-разбора
|
||||
на двенадцати колонках и необходимость сверки ключей. Оговорка про кортеж:
|
||||
`Nullable(Tuple)` в ClickHouse всё-таки есть, но за настройкой
|
||||
`enable_nullable_tuple_type` и в статусе беты; на довод это не влияет — кортежей
|
||||
в контракте нет. `Nullable` внутри кортежа разрешён без всяких флагов, поэтому
|
||||
запасной вариант со свёрткой в именованный кортеж жив.
|
||||
|
||||
Оттуда же: у Kafka-движка есть третий режим обработки ошибок —
|
||||
`dead_letter_queue` с записью в системную таблицу. Он отвергнут независимо от
|
||||
остального: спека требует свои `*_errors`, системная таблица их не заменяет.
|
||||
|
||||
Форматы `RawBLOB` и `LineAsString` в ClickHouse есть, и Kafka-движок
|
||||
поддерживает все форматы. `RawBLOB` читает вход в одно значение и рассчитан на
|
||||
таблицу с единственным полем `String` — отсюда ограничение на форму чтеца.
|
||||
|
||||
Матвью с источником-`Distributed` срабатывает на вставку именно в эту
|
||||
распределённую таблицу — блок она видит до разрезания по шардам. Проверено
|
||||
владельцем на рабочих проектах; на стенде подтверждается заодно с приёмкой #37.
|
||||
|
||||
На живом стенде проверяется при исполнении #37. Первые два внесены в раздел 11
|
||||
спеки; остальные — однострочные `SELECT`, их довольно прогнать заодно:
|
||||
|
||||
- `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение. Проверять это
|
||||
нужно первым и до написания DDL: формулировка «читает вход в одно значение»
|
||||
про файл понятна, а про пачку сообщений из топика — нет, и если сообщения
|
||||
склеятся, переделывать придётся решение целиком, а не DDL. Опыт стоит трёх
|
||||
сообщений и одного `count()`. Запасной вариант — `LineAsString`: он режет по
|
||||
переводу строки, а события у нас однострочные; цена запасного — сообщение с
|
||||
переводом строки внутри даст две строки вместо одной;
|
||||
- форма именованного кортежа в `JSONExtract` с `Nullable`-членами — нужна для
|
||||
свёртки сорока семи вызовов в один, если разбор окажется дорогим;
|
||||
- `isValidJSON('123')` возвращает единицу, а `JSONExtractKeys` от скаляра —
|
||||
пустой массив. На обоих стоят классы брака и их приоритет, а документация
|
||||
поведение на не-объекте не описывает: два `SELECT` закрывают вопрос;
|
||||
- `JSONAsString` действительно падает на некорректном JSON, а не пропускает
|
||||
строку. На этом стоит отказ от него в пользу `RawBLOB`; документация про
|
||||
ошибочный ввод молчит.
|
||||
@@ -0,0 +1,69 @@
|
||||
# ADR 0006. Имена объектов хранилища: суффикс вида
|
||||
|
||||
Дата: 4 августа 2026 года. Статус: принято.
|
||||
|
||||
## Решение
|
||||
|
||||
Имя объекта в ClickHouse заканчивается тем, что это за объект: `_rep` —
|
||||
локальная таблица шарда, `_dist` — `Distributed` поверх неё, `_kafka` — чтец
|
||||
топика, `_mv` — материализованное представление, `_v` — обычное представление.
|
||||
|
||||
Суффикс носит каждый физический объект, поэтому голого имени у таблицы не
|
||||
существует: запрос к `stg.hits_raw` даёт ошибку «нет такой таблицы».
|
||||
Единственное исключение — словари: воплощение у них одно, шардировать нечего, и
|
||||
суффикс ничего не различал бы. В документах и разговоре голое имя означает
|
||||
сущность, у которой этих объектов несколько.
|
||||
|
||||
Конвенция применена к мастер-спеке тем же решением: `stg.kafka_hits` стал
|
||||
`stg.hits_raw_kafka`, `dds.v_event` — `dds.event_v`, витрины `dm.v_*` —
|
||||
`dm.*_v`. Полная таблица суффиксов и следствия для DDL — в
|
||||
[доке хранилища](../architecture/storage.md).
|
||||
|
||||
## Почему
|
||||
|
||||
Главный довод — как ошибается забытый суффикс. Распространённая конвенция, где
|
||||
голое имя означает локальную таблицу, а распределённая получает `_all`,
|
||||
ошибается молча: запрос к голому имени вернёт данные одного шарда и никак об
|
||||
этом не скажет. Правило спеки «проверки и контрольные суммы — только по
|
||||
`Distributed`» такую тишину не переживает: контрольная сумма, посчитанная по
|
||||
половине кластера, выглядит как честное число. При суффиксе вида голого имени
|
||||
нет вовсе, и та же ошибка становится громкой.
|
||||
|
||||
Второй довод — что группируется в списке по алфавиту. Суффикс группирует объекты
|
||||
по сущности: все четыре объекта топика `hits` стоят рядом, потому что различаются
|
||||
хвостом. Префиксный стиль, которым пользовался стенд-предшественник (`kafka_*`,
|
||||
`mv_*`, `v_*`), группирует по технологии, и объекты одной сущности
|
||||
расползаются по алфавиту. В хранилище, где у одной сущности живёт по три-четыре
|
||||
воплощения, полезнее первое. Довод про дерево в клиенте и про `ORDER BY name`:
|
||||
порядок выдачи `SHOW TABLES` документация не оговаривает, так что на него здесь
|
||||
опираться нельзя.
|
||||
|
||||
Третий — преемственность: `_rep` и `_dist` уже используются владельцем в других
|
||||
хранилищах на ClickHouse, и общий словарь между стендами стоит больше, чем
|
||||
локальная стройность.
|
||||
|
||||
Отвергнуты, кроме `_all`: префиксный стиль предшественника — он не покрывает
|
||||
пару локальная/распределённая, для неё префикса просто нет; голое имя как
|
||||
`Distributed` с суффиксом `_local` у локальной — привычное имя ведёт в
|
||||
правильную таблицу, но конвенция расходится с другими стендами владельца;
|
||||
раскладка пары по разным базам (`stg` и `stg_dist`) — удваивает число баз в
|
||||
каждом слое и разъезжается с таблицей слоёв спеки.
|
||||
|
||||
Цена решения — правка принятой спеки задним числом. Имена представлений и
|
||||
витрин переехали с префикса на суффикс, хотя сами объекты спроектированы не
|
||||
полностью и появятся только на этапах 4 и дальше. Размен принят осознанно:
|
||||
конвенция, введённая после того, как по ней написан первый слой, обходится
|
||||
дороже.
|
||||
|
||||
## Что проверено
|
||||
|
||||
Проверять здесь нечем — это соглашение, а не поведение системы. Вместо проверки
|
||||
конвенция прогнана по карте таблиц спеки, раздел 7: суффикс выводится для всех
|
||||
объектов слоёв STG, ODS, DDS и DM, включая пары локальная/распределённая,
|
||||
представления и матвью. Единственным объектом без выводимого суффикса оказался
|
||||
словарь `products` — отсюда исключение в решении.
|
||||
|
||||
Первый прогон был неполным: четыре витрины из восьми остались с префиксом, и
|
||||
заметило это холодное ревью, а не автор. Имена приведены в порядок 5 августа
|
||||
2026 года. Урок не про имена: «прогнал по документу» — такое же утверждение,
|
||||
как утверждение о поведении системы, и проверять его надо так же.
|
||||
@@ -0,0 +1,387 @@
|
||||
# Хранилище: слои и конвенции
|
||||
|
||||
Документ описывает сторону ClickHouse: как называются объекты, какие служебные
|
||||
колонки у них общие, чем нарезаны и сколько живут данные, как устроен приём и из
|
||||
каких файлов собирается DDL. Здесь же карта таблиц, которая растёт с этапами, и
|
||||
раздел «Что проверено» — чему в этом тексте верить и на каком основании.
|
||||
|
||||
**Что здесь описано и чего ещё нет.** Собран этап 1: кластер из двух шардов,
|
||||
keeper, Kafka, каркас сервисов. Объекты хранилища и механизм применения DDL
|
||||
закладывает этап 2 — на момент написания их в репозитории нет. Дальше по тексту
|
||||
устройство описано так, как оно проектируется; построенное от заложенного
|
||||
отличает карта таблиц в конце.
|
||||
|
||||
Зона ответственности у документа одна — хранилище. Генератор описан отдельно:
|
||||
его замысел — в [спеке генератора](../specs/2026-08-01-generator.md), формат
|
||||
события — в [описании выгрузки](../formats/clickstream-event.md). Хранилище
|
||||
строится по описанию выгрузки, как в бою строится по документации источника.
|
||||
Целевая картина всего стенда — [спека «Боевой реализм стенда
|
||||
(v2)»](../specs/2026-07-30-stand-v2-realism.md).
|
||||
|
||||
## Имена объектов
|
||||
|
||||
Имя объекта заканчивается тем, что это за объект:
|
||||
|
||||
| Суффикс | Что это |
|
||||
|---|---|
|
||||
| `_rep` | локальная таблица шарда, движок семейства `Replicated*` |
|
||||
| `_dist` | `Distributed` поверх одноимённой локальной |
|
||||
| `_kafka` | таблица на движке `Kafka` |
|
||||
| `_mv` | материализованное представление |
|
||||
| `_v` | обычное представление |
|
||||
|
||||
Суффикс носит каждый физический объект. Голого имени у таблицы не существует:
|
||||
запрос к `stg.hits_raw` даёт громкую ошибку «нет такой таблицы» — а под голым
|
||||
именем в документах и разговоре понимается сущность, у которой этих объектов
|
||||
несколько. Единственное исключение — словари: у них воплощение одно, шардировать
|
||||
нечего, и суффикс ничего не различал бы.
|
||||
|
||||
Распространённая конвенция, где голое имя означает локальную таблицу, а
|
||||
распределённая получает суффикс `_all`, ошибается иначе: забытый суффикс тихо
|
||||
возвращает данные одного шарда. Правило спеки «проверки и контрольные суммы —
|
||||
только по `Distributed`» такую тишину переживает плохо.
|
||||
|
||||
Второе свойство — в списке по алфавиту объекты группируются по сущности, а не по
|
||||
технологии: все четыре объекта топика `hits` стоят рядом, потому что различаются
|
||||
хвостом, а не началом имени. Речь про дерево в клиенте и про `ORDER BY name`:
|
||||
порядок выдачи `SHOW TABLES` документация не оговаривает.
|
||||
|
||||
Начало имени — сущность, и берётся она в разных слоях из разных мест. В STG имя
|
||||
приходит от транспорта: слой хранит то, что доехало по топику, и зовётся именем
|
||||
топика — `hits_raw` от `hits`, `orders_raw` от `orders`. В типизированных слоях
|
||||
имя приходит от предметной области и стоит в единственном числе: `event`,
|
||||
`order_snapshot`, `session`. Граница между «как привезли» и «что это такое»
|
||||
проходит по STG, и имена её показывают.
|
||||
|
||||
## Раскладка по шардам
|
||||
|
||||
Пишем только в `_dist`. Локальные таблицы остаются для чтения и обслуживания —
|
||||
операций с партициями, ручной переобработки. Правило не про удобство: при записи
|
||||
через распределённую таблицу раскладку определяет ключ шардирования, то есть
|
||||
свойство данных, а при записи в локальную — то, какая нода случайно выполняла
|
||||
код. Отсюда урок стенда: какая нода читала топик, меняется между прогонами
|
||||
(видно в колонке `consumer_host`), а куда легли данные — нет.
|
||||
|
||||
Ключи шардирования: `cityHash64(ClientID)` у событий, `cityHash64(order_id)` у
|
||||
заказов, `cityHash64` сырой строки у STG и у таблицы ошибок. У первых двух хеш
|
||||
выбран против перекоса: структурированный числовой идентификатор распределяется
|
||||
по остатку от деления неравномерно. У сырья выбора нет — строку иначе не
|
||||
разложишь; там хеш даёт другое свойство, одинаковые сообщения ложатся на один
|
||||
шард.
|
||||
|
||||
Ключи у слоёв разные, и это имеет наблюдаемое следствие: сырая строка и
|
||||
разобранное из неё событие почти всегда оказываются на разных шардах.
|
||||
Пошардовые счётчики STG и ODS поэтому не сходятся и сходиться не должны —
|
||||
сверять слои можно только через `_dist`.
|
||||
|
||||
Ключи ко-локации названы заранее, потому что на них стоит политика соединений из
|
||||
раздела 6 спеки: обычное соединение разрешено только по ключу ко-локации, всё
|
||||
прочее — через `GLOBAL`. Значит `dds.session` и `dds.identity_map` шардируются по
|
||||
`cityHash64(ClientID)`, а `dds.order` и производные от заказа — по
|
||||
`cityHash64(order_id)`. Ключи витрин появятся вместе с самими витринами.
|
||||
|
||||
Открытый вопрос на будущее — не сама замена партиций: операции с ними по
|
||||
локальным таблицам правило разрешает прямо. Вопрос в шаге до неё. Партиция-донор
|
||||
должна быть уже разложена по шардам по тому же ключу, а разложить её можно
|
||||
только вставкой через распределённую таблицу — значит у каждой пакетной сущности
|
||||
появится вторая пара объектов, и имени для неё конвенция пока не даёт. Решать
|
||||
это вместе со сборкой DDS, а не задним числом.
|
||||
|
||||
## Служебные колонки
|
||||
|
||||
Собственные колонки не повторяют имён виртуальных. Виртуальные даёт движок:
|
||||
`_topic`, `_partition`, `_offset`, `_timestamp` у Kafka, `_shard_num` у
|
||||
`Distributed` и прочие. Если положить на диск колонку с таким же именем, в
|
||||
матвью перестанет читаться, что дано движком, а что положено нами, — а это
|
||||
ровно то различие, ради которого метаданные доставки и хранятся. Поэтому они
|
||||
ложатся под именами `kafka_topic`, `kafka_partition`, `kafka_offset`,
|
||||
`kafka_timestamp`, рядом — `consumer_host`, имя читавшей ноды: виртуальные
|
||||
колонки его не несут, а после записи в `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` тем, чем пришло: чтец читает байты и ничего
|
||||
не проверяет, поэтому там оказываются и целые события, и мусор. Разбирается всё
|
||||
это ниже, в матвью ODS — см. [ADR 0005](../adr/0005-event-ingestion.md).
|
||||
|
||||
Движок таблицы сырья — обычный `ReplicatedMergeTree`, `ORDER BY (kafka_partition,
|
||||
kafka_offset)`: разбор полётов идёт от «какое сообщение», другого ключа у сырья и
|
||||
нет. Замену версий сюда ставить нельзя — она отменила бы свойство слоя, ради
|
||||
которого он заведён: повтор доставки в сырье обязан быть виден.
|
||||
|
||||
Метка времени загрузки зовётся `_load_ts`, тип `DateTime64(3)`. Ставится она
|
||||
один раз, в матвью приёма, и дальше переносится из STG в ODS как есть: колонка
|
||||
отвечает на вопрос «когда строка приехала в хранилище», а не «когда её
|
||||
разобрали». В ODS она же служит колонкой версии `ReplacingMergeTree`, и работа у
|
||||
этой версии ровно одна — схлопнуть повтор доставки. Содержимое у повтора то же
|
||||
самое, отличается только метка, поэтому какая из двух строк переживёт мерж,
|
||||
безразлично. Пакетной переобработки у ODS нет: слой наполняет матвью, а не
|
||||
задание Airflow, и работа с партициями начинается выше. Переделать разобранное
|
||||
руками можно — вставкой из сырья с фильтром по `_load_ts`, в пределах
|
||||
трёхсуточного окна; ничья по версии разрешается в пользу вставленного позже.
|
||||
|
||||
Имя согласовано с каноном служебных полей соседнего учебного стенда на
|
||||
Greenplum, чтобы словарь был общим у двух хранилищ; ведущее подчёркивание у
|
||||
технических колонок — распространённая запись,
|
||||
её же используют Fivetran, Airbyte и Stitch. С правилом выше это не спорит:
|
||||
запрещено совпадать с именами виртуальных колонок, а не носить подчёркивание.
|
||||
|
||||
Идентификатора пачки загрузки (`_load_id`) пока нет. В STG и ODS данные приезжают
|
||||
потоком через матвью, у которого нет ни батча, ни `run_id`, и колонка была бы
|
||||
пустой формальностью. В слоях, которые наполняет Airflow, `run_id` появится
|
||||
по-настоящему — тогда и заведём, тем же стилем имени.
|
||||
|
||||
## Путь реплицированных таблиц в keeper
|
||||
|
||||
Шаблон — `/clickhouse/tables/{shard}/{database}/{table}`. База в пути
|
||||
обязательна: без неё одноимённые таблицы разных слоёв получат один и тот же узел
|
||||
в keeper и подерутся. Макрос `{uuid}` не используем, хотя он тоже развёл бы
|
||||
пути: он завязан на движок базы `Atomic` и делает путь нечитаемым, а на учебном
|
||||
стенде возможность открыть `system.zookeeper` и увидеть осмысленный путь — сама
|
||||
по себе половина урока про то, чем занят keeper.
|
||||
|
||||
## Приём: поток и его свойства
|
||||
|
||||
Цепочка одна: чтец топика → матвью → сырьё STG → матвью разбора → событие и
|
||||
таблица ошибок ODS.
|
||||
|
||||
Чтец стоит на обеих нодах и читает одной группой потребителей — имя группы
|
||||
`clickstream_hits`, и оно одинаково на обеих нодах по построению: DDL идёт
|
||||
`ON CLUSTER` и макросов в имени не содержит. Разные группы дали бы каждой ноде
|
||||
полную копию топика, и это отдельная сцена для лабы, а не рабочий режим. Имя
|
||||
кластера в `ON CLUSTER` и в движке `Distributed` — `clickstream_cluster`, оно
|
||||
задано в `infra/clickhouse/config.d/cluster.xml`.
|
||||
|
||||
**Источник матвью разбора — `stg.hits_raw_dist`, а не локальная таблица.**
|
||||
У матвью две привязки: источник, на вставку в который она срабатывает, и цель,
|
||||
куда пишет. Распределённая таблица — лицо слоя, локальная — его хранилище;
|
||||
потребитель слоя цепляется к лицу. Практически это значит, что разбор идёт на
|
||||
той же ноде, что читала Kafka, в момент вставки первой матвью — до раскладки по
|
||||
шардам.
|
||||
|
||||
**Вставка фоновая, и окно потери мы принимаем.** Вставка в распределённую
|
||||
таблицу кладёт блок в локальный спул и сразу возвращает управление, а Kafka
|
||||
коммитит офсеты по факту работы матвью — то есть по факту записи в спул. Топик
|
||||
уже считает сообщение прочитанным, хотя на шарде его ещё нет: умри нода в этом
|
||||
промежутке — сообщения не перечитаются.
|
||||
|
||||
Закрывает окно настройка `distributed_foreground_insert = 1`, и на ETL-вставках
|
||||
Airflow она стоит — там это обычный `SETTINGS` у запроса. На пути приёма её нет,
|
||||
и по трём причинам. Вставку выполняет фоновый поток Kafka-движка, своего запроса
|
||||
у него не бывает, так что настройка уровня запроса доехала бы только профилем
|
||||
пользователя в конфигурации ноды. Синхронный режим связывает шарды: пока второй
|
||||
недоступен, вставка падает, офсеты не коммитятся, и приём встаёт целиком — тогда
|
||||
как при фоновом первая нода продолжает принимать и копит спул для соседа.
|
||||
Платится при этом не одно ожидание на блок, а три распределённые вставки — сырьё,
|
||||
событие, ошибки, — и все внутри потока-потребителя, что само по себе повод для
|
||||
ребаланса по таймауту сессии. Против всего этого — окно в сотню миллисекунд на
|
||||
ноутбуке, где мир пересобирается одной командой. Размен не в пользу настройки, а
|
||||
компромисс полезнее показать, чем спрятать за галочкой.
|
||||
|
||||
**Гарантии нет ни в одну сторону — есть два узких окна.** Окно потери описано
|
||||
выше: нода умерла между коммитом офсетов и сбросом спула. Окно дубля
|
||||
противоположное: нода умерла после записи на шард, но до коммита офсетов, и при
|
||||
перечитывании сообщение приедет второй раз. Сказать про такой приём «хотя бы
|
||||
один раз» нельзя — это обещало бы, что потерь не бывает, а они возможны.
|
||||
|
||||
Дубль ниже по течению ведёт себя по-разному. В ODS его схлопнет
|
||||
`ReplacingMergeTree`, а сырьё дедупа не имеет вовсе: перезаливка модельного дня
|
||||
честно удваивает `count()` в STG, и живёт эта пара до истечения срока хранения.
|
||||
Это свойство слоя, а не поломка, — но обещание идемпотентности конвейера к
|
||||
сырому слою не относится.
|
||||
|
||||
Оговорка к последнему: у семейства `Replicated*` есть своя дедупликация — блок с
|
||||
тем же хешем, вставленный повторно, отбрасывается (`insert_deduplicate`).
|
||||
Удвоение сырья проходит мимо неё только потому, что при повторном чтении
|
||||
`_load_ts` новый и хеш блока другой. Свойство слоя держится на этом, а не на
|
||||
отсутствии механизма.
|
||||
|
||||
**Матвью разбора две, и их условия обязаны делить поток без зазора и без
|
||||
нахлёста.** Одна забирает годные строки в `ods.event_dist`, вторая — брак в
|
||||
`ods.event_errors_dist`. Строка, подошедшая обеим, задвоится; не подошедшая ни
|
||||
одной — исчезнет молча. Держится это формой: второе условие пишется буквальным
|
||||
отрицанием первого, а сам предикат собирается только из функций, не возвращающих
|
||||
NULL, — иначе трёхзначная логика даст строку, которую не возьмёт ни `условие`,
|
||||
ни `NOT условие`. Те же функции не должны и бросать исключений: упавшая матвью
|
||||
роняет вставку и останавливает потребление до починки
|
||||
([ADR 0005](../adr/0005-event-ingestion.md)).
|
||||
|
||||
Постоянной сверки счётчиков при этом нет и не должно быть. У сырья срок жизни
|
||||
трое суток, а ODS хранит всё, поэтому равенство «сырьё = события + ошибки»
|
||||
разъедется само: сырьё истечёт раньше, а перезаливка модельного дня задвоит его,
|
||||
тогда как в ODS тот же повтор схлопнется. Правило, красное в норме, учит не
|
||||
смотреть на оповещения
|
||||
([ADR 0002](../adr/0002-monitoring-scope.md)). Равенство проверяется разово в
|
||||
smoke на управляемой пачке: отправили N сообщений — получили N строк сырья и N в
|
||||
сумме событий и ошибок. Счёт по ODS идёт через `FINAL`: голый `count()` по
|
||||
`ReplacingMergeTree` зависит от того, сколько мержей успело пройти, и спека это
|
||||
прямо запрещает (раздел 6).
|
||||
|
||||
## Срок жизни сырья
|
||||
|
||||
Сырьё в STG живёт трое суток реального времени и уходит само. Трое — это окно
|
||||
отладки: столько сырьё лежит на ноутбуке, чтобы менти успел разобрать полёты,
|
||||
после чего перестаёт занимать место. Нарезка — по дню загрузки, срок — по той же
|
||||
колонке `_load_ts`, снятие — целыми кусками (`ttl_only_drop_parts`). В DDL
|
||||
значение проставлено явно, чтобы поведение не зависело от умолчания версии.
|
||||
|
||||
Нарезать сырьё по модельному дню события было бы соблазнительно — он единица
|
||||
переобработки и он же ключ партиции в ODS, — но чистку это ломает. Кусок
|
||||
снимается целиком, только когда в нём истекли все строки, а в партиции модельного
|
||||
дня лежит приехавшее в разное реальное время: мерж склеит куски разного возраста,
|
||||
самая свежая строка удержит весь кусок, и данные переживут срок неограниченно.
|
||||
|
||||
В партиции дня загрузки склейка идёт точно так же, и строки в ней тоже разного
|
||||
возраста — но не более чем на сутки, потому что партицию закрывает календарный
|
||||
день. Отсюда и оценка: сырьё живёт трое суток плюс хвост до суток, а не ровно
|
||||
трое. Оговорка про сроки: TTL исполняется на мержах, а не по будильнику, так что
|
||||
«уходит само» здесь обещано, а «уходит вовремя» — нет.
|
||||
|
||||
Модельного дня среди колонок сырья нет вовсе. Он свойство содержимого, а
|
||||
содержимое разбирает ODS — там `EventDate` и живёт, ключом партиции. Сырьё
|
||||
режется своими координатами: и разбор полётов, и переобработка фильтруют по
|
||||
`_load_ts`, попадая в ключ партиции, а не идя сплошным проходом. Фильтр по
|
||||
модельному дню вдобавок пропускал бы битые строки — у них дата не извлекается.
|
||||
|
||||
Две оси времени тут не совпадают намеренно. Ось модельного времени начинается в
|
||||
D0 и к реальному календарю не привязана; пакетный режим проигрывает две недели
|
||||
модельного мира за минуты реальных. Поэтому у переобработки и у гигиены диска
|
||||
разные часы, и обслуживают их разные средства. Декларативный TTL по модельной
|
||||
дате был бы просто сломан: он отсчитывает срок от реального «сейчас» и удалял бы
|
||||
эталонные дни прямо на входе.
|
||||
|
||||
В бою слой сырья иногда собирают на движке `Null` — тогда он не хранится вовсе.
|
||||
Такой вариант отвергнут: на стенде сырьё нужно для отладки, поэтому окно, а не
|
||||
ноль.
|
||||
|
||||
## Таблица ошибок
|
||||
|
||||
`ods.event_errors` держит строки, не прошедшие строгий приём, вместе с их сырым
|
||||
текстом, метаданными доставки и классом брака. Ключи её собственные, потому что у
|
||||
брака нет разобранных полей: шардируется `cityHash64` сырой строки — `ClientID` у
|
||||
строки, которая не разобралась, взять неоткуда; нарезается по дню загрузки, как и
|
||||
сырьё; живёт месяц. Дольше сырья — намеренно: если брак истекает вместе с ним,
|
||||
разбираться к моменту разбирательства будет уже нечем.
|
||||
|
||||
Класс брака лежит в колонке `error_class` типа `LowCardinality(String)`. Без неё
|
||||
в таблице копятся строки «что-то не так» без ответа на «что именно», а витрине
|
||||
качества не на что опереться. Сами классы, их порядок и довод, почему порядок
|
||||
обязателен, — в [ADR 0005](../adr/0005-event-ingestion.md).
|
||||
|
||||
Движок — обычный `ReplicatedMergeTree`, без замены версий: схлопывать брак не по
|
||||
чему, у него нет ключа сущности. `ORDER BY` — `(error_class, kafka_partition,
|
||||
kafka_offset)`: смотрят такую таблицу от класса, а внутри класса — по координатам
|
||||
доставки.
|
||||
|
||||
## Раскладка DDL
|
||||
|
||||
Файлы лежат в `sql/ddl/` и применяются по порядку имён. Сначала все статичные
|
||||
объекты, потом матвью — тогда к моменту создания матвью её цель уже существует.
|
||||
|
||||
| Файл | Что в нём |
|
||||
|---|---|
|
||||
| `00-databases.sql` | базы слоёв |
|
||||
| `10-stg-tables.sql` | Kafka-таблица, локальная и распределённая таблицы сырья |
|
||||
| `20-ods-tables.sql` | типизированное событие и таблица ошибок |
|
||||
| `30-ods-views.sql` | матвью разбора: сырьё в событие и в ошибки |
|
||||
| `40-stg-views.sql` | матвью приёма: чтец в сырьё |
|
||||
|
||||
Порядок задают два правила. Первое: матвью принадлежит слою своей цели, а не
|
||||
источника, — разбор из STG в ODS лежит среди файлов ODS, потому что наполняет
|
||||
ODS. Второе: матвью приёма создаётся последней из всех, и потому нарушает
|
||||
нумерацию слоёв. Kafka-движок начинает читать топик ровно тогда, когда к нему
|
||||
привязывают первую матвью; создай её раньше разбора — и всё, что доедет в
|
||||
зазоре, ляжет в сырьё и не попадёт в ODS никуда, ни в событие, ни в ошибки. На
|
||||
пустом топике зазор безвреден, поэтому первый прогон о нём не скажет. Проснётся
|
||||
он, когда тома ClickHouse снесены, а данные Kafka целы, — то есть на обычной
|
||||
отладке.
|
||||
|
||||
Применение — двумя одноразовыми сервисами при `make up`, по образцу уже
|
||||
работающих `airflow-init` и `superset-init`. Сначала `kafka-init` создаёт топик
|
||||
`hits` с двумя партициями, затем `clickhouse-init` дожидается его завершения и
|
||||
применяет файлы с ноды 1, `ON CLUSTER`. Этот порядок страхует от автосоздания
|
||||
топика брокером с одной партицией: у потребителя librdkafka разрешение на
|
||||
автосоздание по умолчанию выключено, так что случиться это не обязано, но урок
|
||||
«обе ноды читают топик» умирает тихо, и полагаться на умолчание клиента здесь
|
||||
не стоит.
|
||||
|
||||
Образцы копируются не целиком, и в двух местах. `clickhouse-init` обязан ждать
|
||||
готовности **обеих** нод: `ON CLUSTER` ждёт исполнения на всех хостах и по
|
||||
таймауту бросает, а `airflow-init` ждёт только первую ноду, `superset-init` —
|
||||
только вторую. И второе: оба образца переживают `make up --wait` лишь потому, что
|
||||
от них зависят долгоживущие сервисы; у пары `kafka-init` / `clickhouse-init`
|
||||
таких зависимых нет, и как поведёт себя `--wait` с одноразовым сервисом без них —
|
||||
проверяется при исполнении #37.
|
||||
|
||||
Повторный `make up` поверх живого тома проходит зелёным: весь DDL идёт через
|
||||
`CREATE ... IF NOT EXISTS`. Оборотная сторона — изменённый объект тем же
|
||||
запуском не применяется, причём молча. Отдельного механизма для этого нет и не
|
||||
нужно: правка существующего DDL случается, только пока стенд пишут, а лекарство
|
||||
уже есть — `make clean && make up`. Мир регенерируется, сырьё живёт трое суток,
|
||||
терять нечего.
|
||||
|
||||
## Карта таблиц
|
||||
|
||||
Ниже — то, что закладывает этап 2. В репозитории этих объектов пока нет.
|
||||
|
||||
| Слой | Объект | Что это |
|
||||
|---|---|---|
|
||||
| STG | `stg.hits_raw_kafka` | чтец топика `hits`, формат `RawBLOB` |
|
||||
| STG | `stg.hits_raw_rep` / `_dist` | сырая строка сообщения плюс метаданные доставки |
|
||||
| STG | `stg.hits_raw_mv` | наполняет сырьё из чтеца |
|
||||
| ODS | `ods.event_rep` / `_dist` | типизированное широкое событие |
|
||||
| ODS | `ods.event_errors_rep` / `_dist` | строки, не прошедшие строгий приём |
|
||||
| ODS | `ods.event_mv`, `ods.event_errors_mv` | разбор сырья в событие и в ошибки |
|
||||
|
||||
Слои 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; до тех пор это предположения, а не
|
||||
знание.
|
||||
@@ -309,11 +309,23 @@ CSV в репозитории (`data/catalog/products.csv`: `sku`, `name`, `cate
|
||||
- **Приём Kafka**: Kafka-таблицы и MV — на обеих нодах, одна consumer group,
|
||||
2 партиции на топик; MV пишут в Distributed-цели. Раскладку решает ключ:
|
||||
события — по `cityHash64(ClientID)` (см. 1.3), заказы —
|
||||
`cityHash64(order_id)`, сырьё STG —
|
||||
`cityHash64(сырой строки)`. Урок: «какая нода читала топик — меняется между
|
||||
прогонами, куда легли данные — нет».
|
||||
- **Приём строгий**: `input_format_skip_unknown_fields = 0`, обязательные
|
||||
поля — без значений по умолчанию. Контракт присутствия: генератор выдаёт
|
||||
`cityHash64(order_id)`, сырьё STG — `cityHash64(сырой строки)`; полный
|
||||
список и доводы — в [доке хранилища](../architecture/storage.md). Урок:
|
||||
«какая нода читала топик — меняется между прогонами, куда легли данные — нет».
|
||||
- **Приём строгий**: пять опорных колонок — `WatchID`, `VisitID`, `ClientID`,
|
||||
`EventDate`, `UTCEventTime` — разбираются как `Nullable`, а набор ключей
|
||||
сообщения сверяется с контрактным; строка с NULL среди опорных колонок
|
||||
или с разошедшимся набором ключей уходит в `*_errors`. Опорными выбраны те,
|
||||
чья порча отравляет всё ниже по течению: идентификаторы события, визита и
|
||||
посетителя, дата партиции и метка времени, по которой события упорядочиваются
|
||||
внутри сессии. `CounterID` формально тоже в ключе сортировки, но на стенде он
|
||||
константа, и NULL там взяться неоткуда. Остальные сорок две
|
||||
достаются обычными типами — сорок семь проверок на NULL превратили бы матвью в
|
||||
простыню, а присутствие и так целиком закрыто сверкой ключей. Сверка ключей —
|
||||
не добавка: у массивов NULL не бывает, и пропавшее поле-массив иначе
|
||||
неотличимо от пустого по смыслу. На входе разбора нет вовсе — Kafka-таблица
|
||||
читает сообщение байтами, строгость целиком в матвью ODS
|
||||
([ADR 0005](../adr/0005-event-ingestion.md)). Контракт присутствия: генератор выдаёт
|
||||
**все 47 полей в каждом событии**; «пусто» — пустой массив, пустая строка
|
||||
или 0, а не отсутствие ключа в JSON. Так строгий приём уживается с
|
||||
полями, пустыми по смыслу (ecommerce у `pageview`, UTM у прямого захода). Несовпадение имени поля — громкая
|
||||
@@ -326,8 +338,9 @@ CSV в репозитории (`data/catalog/products.csv`: `sku`, `name`, `cate
|
||||
легитимная GLOBAL-витрина (заказы малы).
|
||||
- **Конвейер без TRUNCATE**: поток — append-only в ReplacingMergeTree (дедуп
|
||||
через argMax); батчевая переобработка — по дневным партициям
|
||||
(`DROP/REPLACE PARTITION ON CLUSTER`); `TRUNCATE ... ON CLUSTER` остаётся
|
||||
только в `make reset`. `DROP/REPLACE PARTITION` работает только по
|
||||
(`DROP/REPLACE PARTITION ON CLUSTER`); `TRUNCATE ... ON CLUSTER` в конвейере не
|
||||
применяется вовсе — полный сброс стенда делается `make clean && make up`, то
|
||||
есть вместе с томами. `DROP/REPLACE PARTITION` работает только по
|
||||
**локальным** таблицам ON CLUSTER, не по Distributed; замена через
|
||||
DROP+INSERT неатомарна — дашборд в середине прогона честно моргает (это
|
||||
осознанная цена, не баг).
|
||||
@@ -369,26 +382,41 @@ README.
|
||||
| Слой | Объект | Что это |
|
||||
|---|---|---|
|
||||
| Kafka | `hits`, `orders` | два топика, по 2 партиции |
|
||||
| STG | `stg.kafka_hits`, `stg.hits_raw` + MV; то же для orders | сырые строки, Kafka Engine на обеих нодах |
|
||||
| STG | `stg.hits_raw_kafka`, `stg.hits_raw` + MV; для orders — развилка этапа 3, не решена (ниже) | сырые строки, Kafka Engine на обеих нодах |
|
||||
| ODS | `ods.event` (+`_errors`) | типизированное широкое событие, ReplacingMergeTree |
|
||||
| ODS | `ods.order_snapshot` (+`_errors`) | слепки заказов как приехали, партиция по `snapshot_date`, без дедупа |
|
||||
| DDS | `dds.session` | сборка сессий из событий (наследник `dds.click`) |
|
||||
| DDS | `dds.v_event` | представление над `ods.event`: snake_case-имена, расшифровка кодов `DeviceCategory`; витрины DM читают его, а не ODS напрямую |
|
||||
| DDS | `dds.event_v` | представление над `ods.event`: snake_case-имена, расшифровка кодов `DeviceCategory`; витрины DM читают его, а не ODS напрямую |
|
||||
| DDS | `dds.order` | единственная дедуплицированная таблица заказа: партиция по дню заказа (`toDate(created_at)`), ReplacingMergeTree(`updated_at`), `ORDER BY order_id`, дедуп до последней версии — argMax в трансформации при сборке |
|
||||
| DDS | `dds.identity_map` | карта кука↔пользователь |
|
||||
| DDS | словарь `products` | каталог из CSV |
|
||||
| DM | витрины `dm.v_*`, `dm.dq_summary` | см. ниже |
|
||||
| DM | витрины `dm.*_v`, `dm.dq_summary` | см. ниже |
|
||||
|
||||
Служебные колонки: `ods.event` и `ods.order_snapshot` получают метку приёма
|
||||
`_ingested_at`; у `ods.event` та же колонка — колонка версии
|
||||
ReplacingMergeTree. Таблицы `stg.*_raw` хранят виртуальные колонки Kafka
|
||||
(`_topic`, `_partition`, `_offset`, `_timestamp`) — без них урок «какая нода
|
||||
читала топик» ненаблюдаем. `stg.hits_raw` дополнительно хранит извлечённый
|
||||
`event_date` — им кормится переобработка дня X (при исчерпании retention
|
||||
У каждой таблицы слоя — пара из локальной и распределённой, имена по конвенции
|
||||
суффиксов; она же задаёт служебные колонки, нарезку и срок хранения сырья —
|
||||
см. [доку хранилища](../architecture/storage.md).
|
||||
|
||||
Как принимаются заказы — развилка этапа 3, и она не решена. Событиям выбран
|
||||
приём сырья байтами с разбором функциями ([ADR 0005](../adr/0005-event-ingestion.md));
|
||||
заказам этот же способ идёт только вместе с ответом на вопрос, нужен ли им слой
|
||||
сырья вообще — у них слепок, а не поток. Нужен — и типизированный чтец даст двух
|
||||
чтецов на один топик, а такую схему ADR 0005 отверг; не нужен — и слои
|
||||
перестают быть единообразными. Разбирать грилингом, когда дойдём до заказов;
|
||||
как учебное сравнение двух способов приёма это записано и в опорных точках
|
||||
раздела 12.
|
||||
|
||||
Состав служебных колонок задаёт дока хранилища. Спеке важны два следствия:
|
||||
`ods.event` и `ods.order_snapshot` получают метку загрузки `_load_ts`, и у
|
||||
`ods.event` она же служит колонкой версии ReplacingMergeTree; а таблицы
|
||||
`stg.*_raw` хранят метаданные доставки Kafka вместе с именем читавшей ноды —
|
||||
без них урок «какая нода читала топик» ненаблюдаем.
|
||||
Модельного дня в STG нет: `EventDate` — свойство содержимого, а содержимое
|
||||
разбирает ODS, где эта колонка и служит ключом партиции. Переобработка режется
|
||||
по времени загрузки, координате самого слоя доставки. При исчерпании retention
|
||||
Kafka день переигрывается генератором заново: снимок — кэш чистой функции,
|
||||
см. [спеку генератора](2026-08-01-generator.md)).
|
||||
см. [спеку генератора](2026-08-01-generator.md).
|
||||
|
||||
`dds.v_event` — первый на стенде пример правила «слой — это контракт, а не
|
||||
`dds.event_v` — первый на стенде пример правила «слой — это контракт, а не
|
||||
обязательно копия данных».
|
||||
|
||||
Событие в DDS не дублируется: склейки четырёх источников больше нет, ODS уже
|
||||
@@ -398,11 +426,11 @@ Kafka день переигрывается генератором заново:
|
||||
|
||||
### Витрины DM
|
||||
|
||||
- **`v_revenue_daily`** (выручка, только от заказов): `report_date`,
|
||||
- **`revenue_daily_v`** (выручка, только от заказов): `report_date`,
|
||||
`product_category` (через `dictGet` каталога + ARRAY JOIN позиций),
|
||||
`orders`, `units`, `revenue`, `aov`. Считается по заказам в статусе
|
||||
`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`,
|
||||
`declared_revenue` (клиент), `items_total` (бэкенд, сравнимая база — не
|
||||
`total`: промокод и доставка клиенту не видны), `status`, `mismatch_class`
|
||||
@@ -416,15 +444,15 @@ Kafka день переигрывается генератором заново:
|
||||
(шестое значение `mismatch_class`, вне приоритетов расхождений). После
|
||||
закрытия окна K таких строк не остаётся — сироты исключены построением
|
||||
(раздел 4).
|
||||
- **`v_utm_effectiveness`** — остаётся клиентской (атрибуция по трекеру);
|
||||
- **`utm_effectiveness_v`** — остаётся клиентской (атрибуция по трекеру);
|
||||
счётчики `purchases`/`add_to_carts` оживают из таксономии, добавляется
|
||||
`declared_revenue` по UTM.
|
||||
- **`v_daily_traffic`** — расширяется парой «посетители» / «известные
|
||||
- **`daily_traffic_v`** — расширяется парой «посетители» / «известные
|
||||
пользователи» (обогащение через `dds.identity_map`, локальное соединение
|
||||
по ключу ко-локации).
|
||||
- `v_events_enriched`, `v_top_pages_daily`, `v_session_overview`,
|
||||
`v_dq_errors_daily` — переезжают на новую модель без смены роли: источник —
|
||||
`dds.v_event`, не `ods.event`.
|
||||
- `events_enriched_v`, `top_pages_daily_v`, `session_overview_v`,
|
||||
`dq_errors_daily_v` — переезжают на новую модель без смены роли: источник —
|
||||
`dds.event_v`, не `ods.event`.
|
||||
- `dm.dq_summary` переводится с TRUNCATE+INSERT на партиционную замену
|
||||
(политика «без TRUNCATE»).
|
||||
|
||||
@@ -465,7 +493,7 @@ v2 стартует пустым, поэтому объём ниже — это
|
||||
|---|---|---|
|
||||
| Генератор | с нуля: модель v1 не переносится (другая модель данных, плюс известные проблемы производительности v1); широкое событие, таксономия, анонимность, N:1, заказы слепками, расхождения A–D, каталог; масштаб — ~4–5 тыс. строк с тестами | L |
|
||||
| Инфраструктура | compose: 2 ноды CH + keeper + остальной стенд; конфиги кластера, макросы; make/скрипты | M — ~10–12 файлов |
|
||||
| SQL | 5 DDL-файлов (ON CLUSTER, Replicated*, Distributed) + трансформации событий, заказов, identity, сверки + словарь | L — ~12–15 файлов, главная сложность |
|
||||
| SQL | DDL по слоям и ролям (ON CLUSTER, Replicated*, Distributed; раскладка файлов — в доке хранилища) + трансформации событий, заказов, identity, сверки + словарь | L — ~12–15 файлов, главная сложность |
|
||||
| Airflow | DAG'и по образцу v1: etl_pipeline (партиционная переобработка, ожидание дневного батча заказов — сенсор/Datasets), world_init/next_day, helpers | M — ~5–6 файлов |
|
||||
| Superset | датасеты + дашборд с тремя новыми сюжетами | M — 2 файла |
|
||||
| Эталонный мир | пересборка снимка на месте, манифест-счётчики, чек-скрипты | M–L |
|
||||
@@ -487,7 +515,7 @@ v2 стартует пустым, поэтому объём ниже — это
|
||||
`cancelled`) — на статусе `paid` стоит «дыхание» выручки; сужение до пары
|
||||
`created`/`cancelled` — запасной ход, если генератор заказов окажется дороже
|
||||
ожиданий. Связка: если срез 1 сработает, определение выручки в
|
||||
`v_revenue_daily` придётся сменить с «заказы в статусе `paid`» на «все
|
||||
`revenue_daily_v` придётся сменить с «заказы в статусе `paid`» на «все
|
||||
неотменённые заказы».
|
||||
- **Отступление от порядка #15**: резолюция предписывала резать в порядке
|
||||
1 → 2 → 3, спека применяет 3 и 2, а 1 держит в резерве. Довод: срезы 3 и 2
|
||||
@@ -557,6 +585,18 @@ smoke-проверки, а не «дашборд зелёный». Это мин
|
||||
прогонами, отсутствие дублей при штатной работе.
|
||||
- Точная форма `ORDER BY` ODS-таблиц (выражение `intHash32` в ключе
|
||||
ReplacingMergeTree).
|
||||
- `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение (ADR 0005).
|
||||
Проверять первым, до написания DDL: если сообщения склеятся, переделывать
|
||||
придётся решение, а не запрос. Запасной формат — `LineAsString`.
|
||||
- Тип виртуальной колонки `_timestamp` у Kafka-движка: обнуляемость и
|
||||
разрядность (секунды против миллисекунд) — от этого зависит объявление
|
||||
`kafka_timestamp` в таблице сырья.
|
||||
- Матвью с источником-`Distributed` срабатывает на вставку именно в эту
|
||||
распределённую таблицу, до раскладки по шардам: на этом стоит цепочка
|
||||
STG → ODS (ADR 0005). Проверено владельцем на рабочих проектах, в документации
|
||||
ClickHouse этот случай не описан.
|
||||
- Форма именованного кортежа в `JSONExtract` с `Nullable`-членами — ею
|
||||
сворачиваются 47 вызовов в один, если разбор окажется дорогим (ADR 0005).
|
||||
- Размер артефакта эталонного мира после пересборки.
|
||||
- Спорные API (Airflow Datasets/сенсоры, ClickHouse DDL) — перед кодом
|
||||
сверять через MCP Context7 (правило AGENTS.md).
|
||||
@@ -610,6 +650,18 @@ v2, этап 0).
|
||||
синтетическая постановка — осознанный приём;
|
||||
- лекция «`Sign` и CollapsingMergeTree»: почему на стенде `sum(Sign)` =
|
||||
`count()`, а в бою — нет; частый вопрос на собеседованиях;
|
||||
- два способа принять топик, рядом на одном стенде: сырьё байтами с разбором
|
||||
функциями (`hits`, [ADR 0005](../adr/0005-event-ingestion.md)) против
|
||||
типизированного чтеца с `kafka_handle_error_mode` — сравнение цены и
|
||||
наблюдаемости как задание. **Развилка этапа 3, не решена**: типизированный
|
||||
чтец идёт заказам только вместе с ответом на вопрос, нужен ли им слой сырья.
|
||||
Нужен — и чтецов на один топик станет два, а эту схему ADR 0005 отверг; не
|
||||
нужен — и слои перестают быть единообразными. Разбирать грилингом, когда
|
||||
дойдём до заказов;
|
||||
- матвью как рабочий механизм, а не диковина: их видно на приёме и на сборке
|
||||
ODS, а пакетная работа начинается выше. Отдельным заданием — как читать из ODS
|
||||
последние версии, через `FINAL` или оконной функцией: что нагляднее, решаем на
|
||||
месте;
|
||||
- лекция про идентичность «как в бою»: `setUserID` и first-party id,
|
||||
детерминированная против вероятностной склейки, identity graph,
|
||||
кросс-девайс, CDP — с рамкой «мы склеили через транзакции, потому что трекер
|
||||
|
||||
@@ -157,7 +157,7 @@
|
||||
строк; сверка — по хешам манифеста, а два локально пересобранных снимка
|
||||
сравнимы обычным diff — пустой означает «ничего не изменилось». Вне
|
||||
обещания — транспорт: офсеты и партиции Kafka, какая нода прочитала,
|
||||
`_ingested_at`, темп живого дня.
|
||||
`_load_ts`, темп живого дня.
|
||||
- **Условия обещания.** Детерминизм держится при зафиксированном `uv.lock`
|
||||
и внутри канонического контейнера — то есть везде Linux, на маке и в WSL
|
||||
тоже; единственная переменная — архитектура CPU. Истина — CI на Linux; сходимость любой машины проверяет
|
||||
@@ -224,13 +224,13 @@
|
||||
менти несёт рендеренная таблица в доках, не модуль. Уточнение при
|
||||
исполнении (#36): нормализованное имя — имя источника, приведённое к
|
||||
нашему стилю, а не имя атрибута в модели данных. Слой DDS складывает свою
|
||||
модель и называет атрибуты по ней; `dds.v_event` эти имена берёт (раздел 7
|
||||
модель и называет атрибуты по ней; `dds.event_v` эти имена берёт (раздел 7
|
||||
мастер-спеки), но контракт их не диктует и тестами не сторожит.
|
||||
- **Сторона хранилища пишется по документации, не генерируется.** DDL
|
||||
`ods.event`, SELECT матвью, `dds.v_event`, трансформации — работа
|
||||
`ods.event`, SELECT матвью, `dds.event_v`, трансформации — работа
|
||||
следующих этапов по «описанию выгрузки», как в бою хранилище адаптируется
|
||||
к источнику. Границу сторожат два боевых механизма: строгий приём
|
||||
(`input_format_skip_unknown_fields = 0`, таблицы `*_errors` — раздел 6
|
||||
(`Nullable`-разбор со сверкой набора ключей, таблицы `*_errors` — раздел 6
|
||||
мастер-спеки) и contract-тест в smoke — сравнение `system.columns`
|
||||
поднятого стенда со схемой генератора.
|
||||
|
||||
@@ -254,6 +254,24 @@
|
||||
превращается в JSON. Приёмники не знают о содержимом: файл (локальный кэш
|
||||
для пересборки и проверок манифеста), Kafka пачкой — пакетный режим,
|
||||
Kafka с темпом ×60 — живой день. Новых топиков нет.
|
||||
- **Одно событие — одно сообщение Kafka.** «Пачкой» относится к темпу
|
||||
отправки, а не к упаковке: приёмник шлёт события подряд без пауз, но каждое
|
||||
отдельным сообщением. Контракт транспорта, не деталь реализации — сторона
|
||||
хранилища читает топик байтами и кладёт сообщение строкой
|
||||
([ADR 0005](../adr/0005-event-ingestion.md)), поэтому склейка нескольких
|
||||
событий в одно сообщение сломала бы разбор целиком. Сторожится тестом
|
||||
приёмника.
|
||||
- **Даты и время на проводе — ISO-8601.** `EventDate` уезжает как `2026-06-01`,
|
||||
`UTCEventTime` — как `2026-06-01T12:34:56Z`. Довод — читаемость сырья: весь
|
||||
смысл слоя STG в том, что менти открывает колонку `raw` в обычном клиенте и
|
||||
разбирает событие глазами, а число эпохи этот урок убивает. Разбору это
|
||||
ничего не стоит: `JSONExtract(raw, 'UTCEventTime', 'Nullable(DateTime)')`
|
||||
принимает ISO без плясок. Колонка `ecommerce` — строка, внутри которой лежит
|
||||
экранированный JSON, как отдаёт Метрика.
|
||||
|
||||
Форму пинит хранилище (#43) как первый потребитель, сериализатор (#41)
|
||||
её соблюдает. Порядок тикетов обратный порядку зависимости, поэтому здесь она
|
||||
и записана — иначе каждый выберет своё, и разойдётся это уже после приёмки.
|
||||
- **Рабочий выбор сериализатора — orjson**: быстрее stdlib json в 5–14 раз,
|
||||
numpy-массивы и datetime сериализует нативно (заметка исследования #31).
|
||||
Смена библиотеки меняет канонические байты, поэтому проходит как
|
||||
|
||||
Reference in New Issue
Block a user