docs(storage): конвенции хранилища и приём событий перед этапом 2 #52

Merged
ddmitry merged 5 commits from docs/37-storage-conventions into main 2026-08-05 22:09:41 +03:00
5 changed files with 490 additions and 24 deletions
Showing only changes of commit 319db308bf - Show all commits
+141
View File
@@ -0,0 +1,141 @@
# ADR 0005. Приём событий: сырьё в STG байтами, разбор функциями в ODS
Дата: 3 августа 2026 года. Статус: принято.
## Решение
Топик `hits` читает одна Kafka-таблица формата `RawBLOB`: сообщение ложится в
`stg.hits_raw_rep` строкой, как пришло, рядом с метаданными доставки. Ни
типизации, ни проверки на этом шаге нет — слой сырья ничего не интерпретирует.
Типизированный слой наполняют две матвью, привязанные к `stg.hits_raw_dist`.
Поля достаются `JSONExtract`. В таблицу ошибок уходят три класса брака:
сообщение, не являющееся объектом JSON; объект, чей набор ключей разошёлся с
контрактным; объект, у которого не разобрался ключевой идентификатор или метка
времени. Остальное — в событие. Первый класс проверяется именно на объект, а не
на валидность: `isValidJSON('123')` возвращает единицу, скаляр — тоже законный
JSON.
Присутствие полей целиком держит сверка набора ключей — одно сравнение
отсортированного `JSONExtractKeys` с контрактным списком. Обязательны все сорок
семь полей: генератор шлёт их все в каждом событии, а «пусто» по контракту —
пустое значение, а не отсутствие ключа. Этим же закрыт критерий #43 про опечатку
в имени.
Тип проверяется не у всех колонок, а у пяти: `WatchID`, `VisitID`, `ClientID`,
`EventDate`, `UTCEventTime` разбираются в `Nullable` и дают NULL, если значение
не той природы. Остальные сорок две достаются обычными типами. Соотношение
цены и пользы: единственный производитель топика — собственный генератор,
сериализующий из контракта по объявленным типам, поэтому неверный тип может
прийти только из руки, а сорок семь проверок на NULL превратили бы матвью в
простыню. Пять выбраны по последствию: на них стоят ключ сортировки, партиция и
дедупликация, и их порча отравляет всё ниже по течению.
Присутствие иначе и не проверить. `Nullable`
в ClickHouse не оборачивает составные типы: `Nullable(Array)` запрещён, а
`Array(Nullable(T))` при пропавшем ключе даёт пустой массив, неотличимый от
пустого по смыслу. Таких колонок в контракте двенадцать из сорока семи.
Типизированная Kafka-таблица не используется. Режим `kafka_handle_error_mode =
'stream'` не используется тоже: при чтении байтами на входе нечему ломаться, и
ошибке разбора взяться неоткуда.
Этим решение снимает ограничение, записанное в постановке #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 года.
При режиме `stream` движок отдаёт `_raw_message` и `_error` только для
сообщений, которые не разобрались, и оставляет их пустыми для разобранных.
Отсюда весь довод о том, что типизированный чтец не может наполнить слой сырья.
Составные типы `Nullable` не оборачивает: `Nullable(Array)`, `Nullable(Map)` и
`Nullable(Tuple)` не поддерживаются, а `Nullable` внутри них — да. Отсюда
слепота `Nullable`-разбора на двенадцати колонках и необходимость сверки ключей.
`Nullable` внутри кортежа разрешён, поэтому запасной вариант со свёрткой в
именованный кортеж жив.
Оттуда же: у Kafka-движка есть третий режим обработки ошибок —
`dead_letter_queue` с записью в системную таблицу. Он отвергнут независимо от
остального: спека требует свои `*_errors`, системная таблица их не заменяет.
Форматы `RawBLOB` и `LineAsString` в ClickHouse есть, и Kafka-движок
поддерживает все форматы.
Матвью с источником-`Distributed` срабатывает на вставку именно в эту
распределённую таблицу — блок она видит до разрезания по шардам. Проверено
владельцем на рабочих проектах; на стенде подтверждается заодно с приёмкой #37.
На живом стенде проверяется при исполнении #37, и то же внесено в раздел 11
спеки:
- `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение;
- форма именованного кортежа в `JSONExtract` с `Nullable`-членами — нужна для
свёртки сорока семи вызовов в один, если разбор окажется дорогим.
+62
View File
@@ -0,0 +1,62 @@
# 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`» такую тишину не переживает: контрольная сумма, посчитанная по
половине кластера, выглядит как честное число. При суффиксе вида голого имени
нет вовсе, и та же ошибка становится громкой.
Второй довод — что группируется в `SHOW TABLES`. Суффикс группирует объекты по
сущности: все четыре объекта топика `hits` стоят рядом, потому что различаются
хвостом. Префиксный стиль, которым пользовался стенд-предшественник (`kafka_*`,
`mv_*`, `v_*`), группирует по технологии, и объекты одной сущности
расползаются по алфавиту. В хранилище, где у одной сущности живёт по три-четыре
воплощения, полезнее первое.
Третий — преемственность: `_rep` и `_dist` уже используются владельцем в других
хранилищах на ClickHouse, и общий словарь между стендами стоит больше, чем
локальная стройность.
Отвергнуты, кроме `_all`: префиксный стиль предшественника — он не покрывает
пару локальная/распределённая, для неё префикса просто нет; голое имя как
`Distributed` с суффиксом `_local` у локальной — привычное имя ведёт в
правильную таблицу, но конвенция расходится с другими стендами владельца;
раскладка пары по разным базам (`stg` и `stg_dist`) — удваивает число баз в
каждом слое и разъезжается с таблицей слоёв спеки.
Цена решения — правка принятой спеки задним числом. Имена представлений и
витрин переехали с префикса на суффикс, хотя сами объекты спроектированы не
полностью и появятся только на этапах 4 и дальше. Размен принят осознанно:
конвенция, введённая после того, как по ней написан первый слой, обходится
дороже.
## Что проверено
Проверять здесь нечем — это соглашение, а не поведение системы. Вместо проверки
конвенция прогнана по карте таблиц спеки, раздел 7: суффикс выводится для всех
объектов слоёв STG, ODS, DDS и DM, включая пары локальная/распределённая,
представления и матвью. Единственным объектом без выводимого суффикса оказался
словарь `products` — отсюда исключение в решении.
+236
View File
@@ -0,0 +1,236 @@
# Хранилище: слои и конвенции
Документ описывает сторону 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`» такую тишину переживает плохо.
Второе свойство — `SHOW TABLES` группирует объекты по сущности, а не по
технологии: все четыре объекта топика `hits` стоят рядом, потому что различаются
хвостом, а не началом имени.
## Раскладка по шардам
Пишем только в `_dist`. Локальные таблицы остаются для чтения и обслуживания —
операций с партициями, ручной переобработки. Правило не про удобство: при записи
через распределённую таблицу раскладку определяет ключ шардирования, то есть
свойство данных, а при записи в локальную — то, какая нода случайно выполняла
код. Отсюда урок стенда: какая нода читала топик, меняется между прогонами
(видно в колонке `consumer_host`), а куда легли данные — нет.
Ключи шардирования: `cityHash64(ClientID)` у событий, `cityHash64(order_id)` у
заказов, `cityHash64` сырой строки у STG и у таблицы ошибок. У первых двух хеш
выбран против перекоса: структурированный числовой идентификатор распределяется
по остатку от деления неравномерно. У сырья выбора нет — строку иначе не
разложишь; там хеш даёт другое свойство, одинаковые сообщения ложатся на один
шард.
Ключи у слоёв разные, и это имеет наблюдаемое следствие: сырая строка и
разобранное из неё событие почти всегда оказываются на разных шардах.
Пошардовые счётчики STG и ODS поэтому не сходятся и сходиться не должны —
сверять слои можно только через `_dist`.
## Служебные колонки
Собственные колонки не повторяют имён виртуальных. Виртуальные даёт движок:
`_topic`, `_partition`, `_offset`, `_timestamp` у Kafka, `_shard_num` у
`Distributed` и прочие. Если положить на диск колонку с таким же именем, в
матвью перестанет читаться, что дано движком, а что положено нами, — а это
ровно то различие, ради которого метаданные доставки и хранятся. Поэтому они
ложатся под именами `kafka_topic`, `kafka_partition`, `kafka_offset`,
`kafka_timestamp`, рядом — `consumer_host`, имя читавшей ноды: виртуальные
колонки его не несут, а после записи в `Distributed` он уже невосстановим.
Само сообщение лежит в колонке `raw` тем, чем пришло: чтец читает байты и ничего
не проверяет, поэтому там оказываются и целые события, и мусор. Разбирается всё
это ниже, в матвью ODS — см. [ADR 0005](../adr/0005-event-ingestion.md).
Метка времени загрузки зовётся `_load_ts`, тип `DateTime64(3)`. В ODS она же
служит колонкой версии `ReplacingMergeTree`. Имя согласовано с каноном служебных
полей соседнего учебного стенда на 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.
**Источник матвью разбора — `stg.hits_raw_dist`, а не локальная таблица.**
У матвью две привязки: источник, на вставку в который она срабатывает, и цель,
куда пишет. Распределённая таблица — лицо слоя, локальная — его хранилище;
потребитель слоя цепляется к лицу. Практически это значит, что разбор идёт на
той же ноде, что читала Kafka, в момент вставки первой матвью — до раскладки по
шардам.
**Вставка синхронная.** На пути приёма стоит
`distributed_foreground_insert = 1`. По умолчанию вставка в распределённую
таблицу кладёт блок в локальный спул и сразу возвращает управление, а Kafka
коммитит офсеты по факту работы матвью — то есть по факту записи в спул. Топик
уже считает сообщение прочитанным, хотя на шарде его нет. Синхронный режим не
добавляет работы, он её не прячет: кусок на шарде будет записан всё равно,
вопрос лишь в том, ждём ли мы этого внутри вставки. Платится ожидание один раз
на блок Kafka в десятки тысяч строк, а не на событие.
**Гарантия — «хотя бы один раз», не транзакция.** Падение после записи на шард,
но до коммита офсетов даёт повтор при перечитывании. В ODS повтор схлопнет
`ReplacingMergeTree`, а сырьё дедупа не имеет вовсе: перезаливка модельного дня
честно удваивает `count()` в STG, и живёт эта пара до истечения срока хранения.
Это свойство слоя, а не поломка, — но обещание идемпотентности конвейера к
сырому слою не относится.
**Матвью разбора две, и их условия обязаны делить поток без зазора и без
нахлёста.** Одна забирает годные строки в `ods.event_dist`, вторая — брак в
`ods.event_errors_dist`. Строка, подошедшая обеим, задвоится; не подошедшая ни
одной — исчезнет молча. Держится это формой: второе условие пишется буквальным
отрицанием первого, а сам предикат собирается только из функций, не возвращающих
NULL, — иначе трёхзначная логика даст строку, которую не возьмёт ни `условие`,
ни `NOT условие`.
Постоянной сверки счётчиков при этом нет и не должно быть. У сырья срок жизни
трое суток, а ODS хранит всё, поэтому равенство «сырьё = события + ошибки»
разъедется само; переобработка добавит строк в ODS, перезаливка задвоит сырьё, а
ODS её схлопнет. Правило, красное в норме, учит не смотреть на оповещения
([ADR 0002](../adr/0002-monitoring-scope.md)). Равенство проверяется разово в
smoke на управляемой пачке: отправили N сообщений — получили N строк сырья и N в
сумме событий и ошибок.
## Срок жизни сырья
Сырьё в STG живёт трое суток реального времени и уходит само. Трое — это окно
отладки: столько сырьё лежит на ноутбуке, чтобы менти успел разобрать полёты,
после чего перестаёт занимать место. Нарезка — по дню загрузки, срок — по той же
колонке `_load_ts`, снятие — целыми кусками (`ttl_only_drop_parts`). Значение
этой настройки между версиями ClickHouse менялось, поэтому в DDL оно проставлено
явно.
Нарезать сырьё по модельному дню события было бы соблазнительно — он единица
переобработки и он же ключ партиции в ODS, — но чистку это ломает. Кусок
снимается целиком, только когда в нём истекли все строки, а обычный мерж внутри
партиции модельного дня склеит куски разного возраста, и данные переживут срок.
В партиции дня загрузки склеивать нечего: все строки в ней одного возраста, и
куски уходят по мере того, как истекает самый свежий из них — то есть сырьё
живёт трое суток с небольшим хвостом, а не ровно трое.
Модельного дня среди колонок сырья нет вовсе. Он свойство содержимого, а
содержимое разбирает ODS — там `EventDate` и живёт, ключом партиции. Сырьё
режется своими координатами: и разбор полётов, и переобработка фильтруют по
`_load_ts`, попадая в ключ партиции, а не идя сплошным проходом. Фильтр по
модельному дню вдобавок пропускал бы битые строки — у них дата не извлекается.
Две оси времени тут не совпадают намеренно. Ось модельного времени начинается в
D0 и к реальному календарю не привязана; пакетный режим проигрывает две недели
модельного мира за минуты реальных. Поэтому у переобработки и у гигиены диска
разные часы, и обслуживают их разные средства. Декларативный TTL по модельной
дате был бы просто сломан: он отсчитывает срок от реального «сейчас» и удалял бы
эталонные дни прямо на входе.
В бою слой сырья иногда собирают на движке `Null` — тогда он не хранится вовсе.
Такой вариант отвергнут: на стенде сырьё нужно для отладки, поэтому окно, а не
ноль.
## Таблица ошибок
`ods.event_errors` держит строки, не прошедшие строгий приём, вместе с их сырым
текстом и метаданными доставки. Ключи её собственные, потому что у брака нет
разобранных полей: шардируется `cityHash64` сырой строки — `ClientID` у строки,
которая не разобралась, взять неоткуда; нарезается по дню загрузки, как и сырьё;
живёт месяц. Дольше сырья — намеренно: если брак истекает вместе с ним,
разбираться к моменту разбирательства будет уже нечем.
## Раскладка DDL
Файлы лежат в `sql/ddl/` и применяются по порядку имён. Один файл — это слой и
роль: статичные объекты отдельно от матвью.
| Файл | Что в нём |
|---|---|
| `00-databases.sql` | базы слоёв |
| `10-stg-tables.sql` | Kafka-таблица, локальная и распределённая таблицы сырья |
| `11-stg-views.sql` | матвью, наполняющая сырьё из Kafka |
| `20-ods-tables.sql` | типизированное событие и таблица ошибок |
| `21-ods-views.sql` | матвью разбора: сырьё в событие и в ошибки |
Матвью принадлежит слою своей цели, а не источника: разбор из STG в ODS лежит
среди файлов ODS, потому что наполняет ODS.
Применение — двумя одноразовыми сервисами при `make up`, по образцу уже
работающих `airflow-init` и `superset-init`. Сначала `kafka-init` создаёт топик
`hits` с двумя партициями, затем `clickhouse-init` дожидается его завершения и
применяет файлы с ноды 1, `ON CLUSTER`. Порядок страхует от автосоздания топика
брокером с одной партицией: у потребителя librdkafka разрешение на автосоздание
по умолчанию выключено, так что случиться это не обязано, но урок «обе ноды
читают топик» умирает тихо, и полагаться на умолчание клиента здесь не стоит.
Повторный `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
мастер-спеки и переносится сюда по мере постройки.
+40 -20
View File
@@ -309,11 +309,16 @@ 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). Урок:
«какая нода читала топик — меняется между прогонами, куда легли данные — нет».
- **Приём строгий**: обязательные поля разбираются как `Nullable`, а набор
ключей сообщения сверяется с контрактным; строка с NULL среди обязательных
полей или с разошедшимся набором ключей уходит в `*_errors`. Сверка ключей —
не добавка: у массивов NULL не бывает, и пропавшее поле-массив иначе
неотличимо от пустого по смыслу. На входе разбора нет вовсе — Kafka-таблица
читает сообщение байтами, строгость целиком в матвью ODS
([ADR 0005](../adr/0005-event-ingestion.md)). Контракт присутствия: генератор выдаёт
**все 47 полей в каждом событии**; «пусто» — пустой массив, пустая строка
или 0, а не отсутствие ключа в JSON. Так строгий приём уживается с
полями, пустыми по смыслу (ecommerce у `pageview`, UTM у прямого захода). Несовпадение имени поля — громкая
@@ -369,26 +374,32 @@ 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 | сырые строки, 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).
Состав служебных колонок задаёт дока хранилища. Спеке важны два следствия:
`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 уже
@@ -422,9 +433,9 @@ Kafka день переигрывается генератором заново:
- **`v_daily_traffic`** — расширяется парой «посетители» / «известные
пользователи» (обогащение через `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 +476,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 — ~56 файлов |
| Superset | датасеты + дашборд с тремя новыми сюжетами | M — 2 файла |
| Эталонный мир | пересборка снимка на месте, манифест-счётчики, чек-скрипты | M–L |
@@ -557,6 +568,11 @@ smoke-проверки, а не «дашборд зелёный». Это мин
прогонами, отсутствие дублей при штатной работе.
- Точная форма `ORDER BY` ODS-таблиц (выражение `intHash32` в ключе
ReplacingMergeTree).
- `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение (ADR 0005).
- Матвью, привязанная к локальной таблице, срабатывает, когда строки приходят
вставкой через `Distributed`: на этом стоит цепочка STG → ODS (ADR 0005).
- Форма именованного кортежа в `JSONExtract` с `Nullable`-членами — ею
сворачиваются 47 вызовов в один, если разбор окажется дорогим (ADR 0005).
- Размер артефакта эталонного мира после пересборки.
- Спорные API (Airflow Datasets/сенсоры, ClickHouse DDL) — перед кодом
сверять через MCP Context7 (правило AGENTS.md).
@@ -610,6 +626,10 @@ v2, этап 0).
синтетическая постановка — осознанный приём;
- лекция «`Sign` и CollapsingMergeTree»: почему на стенде `sum(Sign)` =
`count()`, а в бою — нет; частый вопрос на собеседованиях;
- два способа принять топик, рядом на одном стенде: сырьё байтами с разбором
функциями (`hits`, [ADR 0005](../adr/0005-event-ingestion.md)) против
типизированного чтеца с `kafka_handle_error_mode` (заказы, этап 3) —
сравнение цены и наблюдаемости как задание;
- лекция про идентичность «как в бою»: `setUserID` и first-party id,
детерминированная против вероятностной склейки, identity graph,
кросс-девайс, CDP — с рамкой «мы склеили через транзакции, потому что трекер
+11 -4
View File
@@ -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,13 @@
превращается в JSON. Приёмники не знают о содержимом: файл (локальный кэш
для пересборки и проверок манифеста), Kafka пачкой — пакетный режим,
Kafka с темпом ×60 — живой день. Новых топиков нет.
- **Одно событие — одно сообщение Kafka.** «Пачкой» относится к темпу
отправки, а не к упаковке: приёмник шлёт события подряд без пауз, но каждое
отдельным сообщением. Контракт транспорта, не деталь реализации — сторона
хранилища читает топик байтами и кладёт сообщение строкой
([ADR 0005](../adr/0005-event-ingestion.md)), поэтому склейка нескольких
событий в одно сообщение сломала бы разбор целиком. Сторожится тестом
приёмника.
- **Рабочий выбор сериализатора — orjson**: быстрее stdlib json в 514 раз,
numpy-массивы и datetime сериализует нативно (заметка исследования #31).
Смена библиотеки меняет канонические байты, поэтому проходит как