From 319db308bfc78c615d7c1c34fc13abfd018b0084 Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Tue, 4 Aug 2026 00:14:13 +0300 Subject: [PATCH 1/5] =?UTF-8?q?docs(storage):=20=D0=BF=D1=80=D0=B8=D0=BD?= =?UTF-8?q?=D1=8F=D1=82=D1=8B=20=D1=80=D0=B5=D1=88=D0=B5=D0=BD=D0=B8=D1=8F?= =?UTF-8?q?=20=D0=BF=D0=BE=20=D0=BF=D1=80=D0=B8=D1=91=D0=BC=D1=83=20=D1=81?= =?UTF-8?q?=D0=BE=D0=B1=D1=8B=D1=82=D0=B8=D0=B9=20=D0=B8=20=D0=B8=D0=BC?= =?UTF-8?q?=D0=B5=D0=BD=D0=B0=D0=BC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - тикет #37 молча опирался на конвенции хранилища, которых в проекте не было; без них #43 и следующие этапы разъехались бы в именах, служебных колонках и механике приёма. - Что: - ADR 0005: топик читается байтами в STG, разбор идёт функциями в матвью ODS; строгий приём — сверка набора ключей плюс Nullable на пяти опорных колонках. - ADR 0006: суффикс вида в именах объектов (_rep, _dist, _kafka, _mv, _v). - docs/architecture/storage.md: конвенции имён и служебных колонок, путь в keeper, раскладка по шардам, срок жизни сырья, свойства приёма, раскладка файлов DDL и карта таблиц. - спеки приведены в соответствие: механизм строгого приёма, имена объектов, контракт транспорта «одно событие — одно сообщение Kafka», три проверки при исполнении. - Проверка: - make config-test --- docs/adr/0005-event-ingestion.md | 141 +++++++++++++ docs/adr/0006-object-naming.md | 62 ++++++ docs/architecture/storage.md | 236 ++++++++++++++++++++++ docs/specs/2026-07-30-stand-v2-realism.md | 60 ++++-- docs/specs/2026-08-01-generator.md | 15 +- 5 files changed, 490 insertions(+), 24 deletions(-) create mode 100644 docs/adr/0005-event-ingestion.md create mode 100644 docs/adr/0006-object-naming.md create mode 100644 docs/architecture/storage.md diff --git a/docs/adr/0005-event-ingestion.md b/docs/adr/0005-event-ingestion.md new file mode 100644 index 0000000..16ff1dc --- /dev/null +++ b/docs/adr/0005-event-ingestion.md @@ -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`-членами — нужна для + свёртки сорока семи вызовов в один, если разбор окажется дорогим. diff --git a/docs/adr/0006-object-naming.md b/docs/adr/0006-object-naming.md new file mode 100644 index 0000000..e741226 --- /dev/null +++ b/docs/adr/0006-object-naming.md @@ -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` — отсюда исключение в решении. diff --git a/docs/architecture/storage.md b/docs/architecture/storage.md new file mode 100644 index 0000000..1ba9ef4 --- /dev/null +++ b/docs/architecture/storage.md @@ -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 +мастер-спеки и переносится сюда по мере постройки. diff --git a/docs/specs/2026-07-30-stand-v2-realism.md b/docs/specs/2026-07-30-stand-v2-realism.md index cc0299b..378c43b 100644 --- a/docs/specs/2026-07-30-stand-v2-realism.md +++ b/docs/specs/2026-07-30-stand-v2-realism.md @@ -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 — ~5–6 файлов | | 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 — с рамкой «мы склеили через транзакции, потому что трекер diff --git a/docs/specs/2026-08-01-generator.md b/docs/specs/2026-08-01-generator.md index c4a70f0..0528318 100644 --- a/docs/specs/2026-08-01-generator.md +++ b/docs/specs/2026-08-01-generator.md @@ -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 в 5–14 раз, numpy-массивы и datetime сериализует нативно (заметка исследования #31). Смена библиотеки меняет канонические байты, поэтому проходит как -- 2.54.0 From dedb6c6739ccd21f437e7894d16e895d70e9ff68 Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Tue, 4 Aug 2026 00:21:43 +0300 Subject: [PATCH 2/5] =?UTF-8?q?docs(handoff):=20=D0=BF=D0=B5=D1=80=D0=B5?= =?UTF-8?q?=D0=B4=D0=B0=D1=87=D0=B0=20=D0=BA=D0=BE=D0=BD=D1=82=D0=B5=D0=BA?= =?UTF-8?q?=D1=81=D1=82=D0=B0=20=D0=BF=D0=BE=20=D0=BA=D0=BE=D0=BD=D0=B2?= =?UTF-8?q?=D0=B5=D0=BD=D1=86=D0=B8=D1=8F=D0=BC=20=D1=85=D1=80=D0=B0=D0=BD?= =?UTF-8?q?=D0=B8=D0=BB=D0=B8=D1=89=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - работа над #37 продолжится в новой сессии, а решения сессии разбросаны по двум ADR, доке хранилища, двум спекам и двум тикетам; свежему агенту нужен один вход с указателями и списком того, что переоткрывать не надо. - Что: - добавлен .scratch/handoffs/2026-08-04-storage-conventions.md: указатели на артефакты, перечень принятых развилок, открытые хвосты и особенности работы с трекером. - Проверка: - make config-test --- .../2026-08-04-storage-conventions.md | 91 +++++++++++++++++++ 1 file changed, 91 insertions(+) create mode 100644 .scratch/handoffs/2026-08-04-storage-conventions.md diff --git a/.scratch/handoffs/2026-08-04-storage-conventions.md b/.scratch/handoffs/2026-08-04-storage-conventions.md new file mode 100644 index 0000000..27a591e --- /dev/null +++ b/.scratch/handoffs/2026-08-04-storage-conventions.md @@ -0,0 +1,91 @@ +# Handoff: конвенции хранилища и приём событий (#37) + +Дата: 4 августа 2026 года. Ветка: `docs/37-storage-conventions`. + +## Что это было + +Сессия началась с вопроса «хватает ли информации, чтобы взять #37 в работу». +Оказалось, что нет: тикет молча опирался на конвенции хранилища, которых в +проекте не существовало. Дальше шёл грилинг пяти развилок, затем два холодных +ревью — по линии дефектов и по линии уместности, — затем правки по их находкам. + +Кода не написано ни строки. Итог сессии — принятые решения и приведённые в +соответствие постановки. + +## Где лежат решения + +Ничего из перечисленного здесь не пересказывается — читать по адресам: + +- `docs/adr/0005-event-ingestion.md` — как принимаем события: чтец читает топик + байтами, разбор идёт функциями в матвью ODS. Там же отвергнутые варианты и + цена решения. +- `docs/adr/0006-object-naming.md` — суффикс вида в именах объектов. +- `docs/architecture/storage.md` — рабочий справочник: конвенции имён и + служебных колонок, раскладка по шардам, путь в keeper, срок жизни сырья, + свойства приёма, раскладка файлов DDL, карта таблиц. +- `docs/specs/2026-07-30-stand-v2-realism.md` — правлены разделы 6, 7, 9, 11, + «Витрины DM» и опорные точки для лекций. +- `docs/specs/2026-08-01-generator.md` — раздел 4 получил контракт транспорта + «одно событие — одно сообщение Kafka». +- Тикеты #37 и #43 переписаны целиком; у каждого сверху комментарий с разбором + того, что изменилось против исходной постановки, — читать `tea issues 37 + --comments`. + +Коммит с документами — `319db30`, он же единственный на ветке помимо этого +файла. + +## Чего не переоткрывать + +Всё ниже прошло грилинг с веером вариантов и записано с доводами. Если решение +кажется странным — сначала прочитать довод, а не начинать заново: + +- имена объектов: суффикс вида, а не префикс и не `_all`; +- служебная колонка загрузки `_load_ts`, метаданные доставки без ведущего + подчёркивания; +- нарезка сырья по дню загрузки, а не по модельному дню; TTL трое суток; +- модельного дня в STG нет вовсе; +- формат чтеца `RawBLOB`, режим `kafka_handle_error_mode` не используется; +- строгий приём: сверка набора ключей плюс `Nullable` на пяти опорных колонках, + не на сорока семи; +- источник матвью разбора — `stg.hits_raw_dist`, пишем только в `_dist`, + `distributed_foreground_insert = 1`; +- путь в keeper `/clickhouse/tables/{shard}/{database}/{table}`, без `{uuid}`. + +## Что осталось + +По убыванию веса: + +1. **Проверки на живом стенде**, две из них внесены в раздел 11 мастер-спеки: + `RawBLOB` даёт ровно одну строку на сообщение; форма именованного кортежа в + `JSONExtract` с `Nullable`-членами. Не внесены, но всплыли в ревью: как + `make up --wait` поведёт себя с одноразовым сервисом, у которого нет + зависимых долгоживущих (существующие `airflow-init` и `superset-init` + переживают `--wait` только за счёт зависимостей), и как аккуратнее навесить + `distributed_foreground_insert` на путь приёма — это настройка уровня + запроса, вероятно через профиль пользователя в конфиге. +2. **PR не открыт.** Ветка запушена, тело PR писать с английским `Closes #NN`, + иначе Gitea задачу не закроет (см. `docs/agents/issue-tracker.md`). +3. **Третий холодный проход** — по желанию владельца. Оба документа переписаны + после ревью существенно, а разделы про keeper, таблицу ошибок и свойства + приёма ревьюеры не видели вовсе. +4. **Реализация #37** — собственно этап, ради которого всё затевалось. + +## Что стоит знать про ход работы + +- Владелец правит рекомендации по существу и часто оказывается прав: из пяти + развилок три пересматривались по его возражениям. Предлагать вариант с + доводом, а не спрашивать «как сделать», — и быть готовым, что довод разберут. +- Спека не выбита в камне: менять её аргументированно можно и нужно, но + обсуждая с владельцем, а не молча. +- Трекер — Gitea, только через `tea`, при проблемах с прокси префикс + `NO_PROXY='*'`. Тикеты читать с `--comments`: хвосты живут там. +- Сверка спорных API ClickHouse — через MCP Context7 до кода, не после. + +## Suggested skills + +- `brainstorm-with-docs` — если всплывёт новая развилка. Именно им шла эта + сессия; формат «веер вариантов, потом конвергенция» владельцу привычен. +- `claude-subagent-playbook` — когда #37 пойдёт в реализацию: конвейер с + делегированием механической части и слепым ревью. +- `conventional-commits` — обязателен при любом коммите в этом репозитории. +- `code-review` — для ревью изменений перед приёмкой. -- 2.54.0 From 31b274175a45ce573031e41a4328834ac6772e7a Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Wed, 5 Aug 2026 21:25:35 +0300 Subject: [PATCH 3/5] =?UTF-8?q?docs(storage):=20=D0=BA=D0=BE=D0=BD=D0=B2?= =?UTF-8?q?=D0=B5=D0=BD=D1=86=D0=B8=D0=B8=20=D0=B8=20=D0=BF=D1=80=D0=B8?= =?UTF-8?q?=D1=91=D0=BC=20=D1=81=D0=BE=D0=B1=D1=8B=D1=82=D0=B8=D0=B9=20?= =?UTF-8?q?=D0=B2=D1=8B=D0=BF=D1=80=D0=B0=D0=B2=D0=BB=D0=B5=D0=BD=D1=8B=20?= =?UTF-8?q?=D0=BF=D0=BE=D1=81=D0=BB=D0=B5=20=D1=80=D0=B5=D0=B2=D1=8C=D1=8E?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - три холодных ревью и сверка с документацией ClickHouse нашли противоречия между докой, ADR и спекой: исполнитель #37 получал два разных ответа на один вопрос, а два утверждения о движке оказались неверными. - Что: - раскладка файлов DDL перестроена — сначала таблицы, матвью приёма последней: иначе часть событий тихо минует ODS. - синхронная вставка снята с пути приёма: настройка недостижима для потока Kafka-движка и связывает шарды; на ETL-вставках осталась. - у таблицы ошибок появился класс брака с порядком проверки, у сырья и ошибок названы движки и ключи сортировки. - в доку добавлен раздел «Что проверено»: сверенное с документацией, проверяемое на стенде и сказанное по памяти разведены. - в спеке выправлены источник матвью разбора, пять опорных колонок, имена четырёх витрин и ссылка на несуществующую цель make. - Проверка: - make config-test - grep по устаревшим именам файлов DDL и витрин — пусто Co-Authored-By: Claude Opus 5 --- AGENTS.md | 4 + README.md | 5 + docs/adr/0005-event-ingestion.md | 63 +++++-- docs/adr/0006-object-naming.md | 13 +- docs/architecture/storage.md | 214 ++++++++++++++++++---- docs/specs/2026-07-30-stand-v2-realism.md | 50 +++-- 6 files changed, 284 insertions(+), 65 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index f39838a..99ac974 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -125,3 +125,7 @@ Airflow) и названия из кода. Если для понятия ес - Имена файлов в `docs/adr/` — `NNNN-краткое-имя.md`: сквозной номер из четырёх цифр и слаг (`0001-stand-services.md`). Решения нумеруются подряд, дата в имени не нужна. +- Имена файлов в `docs/architecture/` — слаг строчными латинскими буквами через + дефис (`storage.md`). Здесь живут рабочие справочники по зонам + ответственности: не событие истории и не решение, а текущее устройство — + один файл на зону, правится по мере постройки. diff --git a/README.md b/README.md index eb7475b..64a3285 100644 --- a/README.md +++ b/README.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/) — описания форматов источников: по ним diff --git a/docs/adr/0005-event-ingestion.md b/docs/adr/0005-event-ingestion.md index 16ff1dc..2959f83 100644 --- a/docs/adr/0005-event-ingestion.md +++ b/docs/adr/0005-event-ingestion.md @@ -5,8 +5,12 @@ ## Решение Топик `hits` читает одна Kafka-таблица формата `RawBLOB`: сообщение ложится в -`stg.hits_raw_rep` строкой, как пришло, рядом с метаданными доставки. Ни +`stg.hits_raw_dist` строкой, как пришло, рядом с метаданными доставки. Ни типизации, ни проверки на этом шаге нет — слой сырья ничего не интерпретирует. +У формата есть следствие для DDL: он читает вход в одно значение и рассчитан на +таблицу с единственной колонкой, поэтому у чтеца она ровно одна — `raw`, а +метаданные доставки берутся только из виртуальных колонок и добавить к чтецу +что-либо своё нельзя. Типизированный слой наполняют две матвью, привязанные к `stg.hits_raw_dist`. Поля достаются `JSONExtract`. В таблицу ошибок уходят три класса брака: @@ -16,8 +20,17 @@ на валидность: `isValidJSON('123')` возвращает единицу, скаляр — тоже законный JSON. +Классы пересекаются: скаляр проваливает заодно и сверку ключей, потому что +`JSONExtractKeys` от него даёт пустой массив. Поэтому они проверяются по порядку, +а в колонку `error_class` пишется первый совпавший — `not_an_object`, +`keyset_mismatch`, `key_field_unparsed`. Приём тот же, что у `mismatch_class` в +витрине сверки: пересекающиеся классы плюс объявленный приоритет. + Присутствие полей целиком держит сверка набора ключей — одно сравнение -отсортированного `JSONExtractKeys` с контрактным списком. Обязательны все сорок +`arraySort(JSONExtractKeys(raw))` с контрактным списком, завёрнутым в тот же +`arraySort`. Обёртка с обеих сторон стоит ноль и снимает ошибку, которая иначе +увела бы в брак вообще всё: сорок семь CamelCase-имён, выписанных руками ровно в +байтовом порядке. Обязательны все сорок семь полей: генератор шлёт их все в каждом событии, а «пусто» по контракту — пустое значение, а не отсутствие ключа. Этим же закрыт критерий #43 про опечатку в имени. @@ -40,6 +53,16 @@ JSON. 'stream'` не используется тоже: при чтении байтами на входе нечему ломаться, и ошибке разбора взяться неоткуда. +Отсюда ограничение на форму выражений разбора: они собираются только из функций, +которые не бросают исключений. `JSONExtract` и родственные возвращают значение по +умолчанию или NULL, но не падают, — и правило репозитория «грязные записи не +валят пайплайн» держится теперь именно на этом. Исключение в матвью не ошибка +формата, его не перехватит никакой режим Kafka-движка: вставка упадёт, офсеты не +закоммитятся, блок пойдёт читаться снова. Ломается это громко и чинится без +потерь — поправил матвью, потребление продолжилось с некоммиченного офсета, — но +пока не починено, топик стоит. Поэтому приведение `Nullable` к необнуляемому типу +и любая арифметика в этих выражениях живут за предикатом, который NULL уже отсёк. + Этим решение снимает ограничение, записанное в постановке #37: «ошибки разбора должны рождаться на шаге Kafka-движка». Оно ставилось как условие достижимости критерия #43 про громкую ошибку на опечатку в имени поля — критерий достижим и @@ -116,26 +139,42 @@ contract-тест из #43 сюда не дотягивается — он ср сообщений, которые не разобрались, и оставляет их пустыми для разобранных. Отсюда весь довод о том, что типизированный чтец не может наполнить слой сырья. -Составные типы `Nullable` не оборачивает: `Nullable(Array)`, `Nullable(Map)` и -`Nullable(Tuple)` не поддерживаются, а `Nullable` внутри них — да. Отсюда -слепота `Nullable`-разбора на двенадцати колонках и необходимость сверки ключей. -`Nullable` внутри кортежа разрешён, поэтому запасной вариант со свёрткой в -именованный кортеж жив. +Составные типы `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 -спеки: +На живом стенде проверяется при исполнении #37. Первые два пункта внесены в +раздел 11 спеки как несущие; остальные — однострочные `SELECT`, их довольно +прогнать заодно: -- `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение; +- `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение. Проверять это + нужно первым и до написания DDL: формулировка «читает вход в одно значение» + про файл понятна, а про пачку сообщений из топика — нет, и если сообщения + склеятся, переделывать придётся решение целиком, а не DDL. Опыт стоит трёх + сообщений и одного `count()`. Запасной вариант — `LineAsString`: он режет по + переводу строки, а события у нас однострочные; цена запасного — сообщение с + переводом строки внутри даст две строки вместо одной; - форма именованного кортежа в `JSONExtract` с `Nullable`-членами — нужна для - свёртки сорока семи вызовов в один, если разбор окажется дорогим. + свёртки сорока семи вызовов в один, если разбор окажется дорогим; +- `isValidJSON('123')` возвращает единицу, а `JSONExtractKeys` от скаляра — + пустой массив. На обоих стоят классы брака и их приоритет, а документация + поведение на не-объекте не описывает: два `SELECT` закрывают вопрос; +- `JSONAsString` действительно падает на некорректном JSON, а не пропускает + строку. На этом стоит отказ от него в пользу `RawBLOB`; документация про + ошибочный ввод молчит. diff --git a/docs/adr/0006-object-naming.md b/docs/adr/0006-object-naming.md index e741226..707c99b 100644 --- a/docs/adr/0006-object-naming.md +++ b/docs/adr/0006-object-naming.md @@ -29,12 +29,14 @@ половине кластера, выглядит как честное число. При суффиксе вида голого имени нет вовсе, и та же ошибка становится громкой. -Второй довод — что группируется в `SHOW TABLES`. Суффикс группирует объекты по -сущности: все четыре объекта топика `hits` стоят рядом, потому что различаются +Второй довод — что группируется в списке по алфавиту. Суффикс группирует объекты +по сущности: все четыре объекта топика `hits` стоят рядом, потому что различаются хвостом. Префиксный стиль, которым пользовался стенд-предшественник (`kafka_*`, `mv_*`, `v_*`), группирует по технологии, и объекты одной сущности расползаются по алфавиту. В хранилище, где у одной сущности живёт по три-четыре -воплощения, полезнее первое. +воплощения, полезнее первое. Довод про дерево в клиенте и про `ORDER BY name`: +порядок выдачи `SHOW TABLES` документация не оговаривает, так что на него здесь +опираться нельзя. Третий — преемственность: `_rep` и `_dist` уже используются владельцем в других хранилищах на ClickHouse, и общий словарь между стендами стоит больше, чем @@ -60,3 +62,8 @@ объектов слоёв STG, ODS, DDS и DM, включая пары локальная/распределённая, представления и матвью. Единственным объектом без выводимого суффикса оказался словарь `products` — отсюда исключение в решении. + +Первый прогон был неполным: четыре витрины из восьми остались с префиксом, и +заметило это холодное ревью, а не автор. Имена приведены в порядок 5 августа +2026 года. Урок не про имена: «прогнал по документу» — такое же утверждение, +как утверждение о поведении системы, и проверять его надо так же. diff --git a/docs/architecture/storage.md b/docs/architecture/storage.md index 1ba9ef4..47847eb 100644 --- a/docs/architecture/storage.md +++ b/docs/architecture/storage.md @@ -40,9 +40,17 @@ keeper, Kafka, каркас сервисов. Объекты хранилища возвращает данные одного шарда. Правило спеки «проверки и контрольные суммы — только по `Distributed`» такую тишину переживает плохо. -Второе свойство — `SHOW TABLES` группирует объекты по сущности, а не по +Второе свойство — в списке по алфавиту объекты группируются по сущности, а не по технологии: все четыре объекта топика `hits` стоят рядом, потому что различаются -хвостом, а не началом имени. +хвостом, а не началом имени. Речь про дерево в клиенте и про `ORDER BY name`: +порядок выдачи `SHOW TABLES` документация не оговаривает. + +Начало имени — сущность, и берётся она в разных слоях из разных мест. В STG имя +приходит от транспорта: слой хранит то, что доехало по топику, и зовётся именем +топика — `hits_raw` от `hits`, `orders_raw` от `orders`. В типизированных слоях +имя приходит от предметной области и стоит в единственном числе: `event`, +`order_snapshot`, `session`. Граница между «как привезли» и «что это такое» +проходит по STG, и имена её показывают. ## Раскладка по шардам @@ -65,6 +73,18 @@ keeper, Kafka, каркас сервисов. Объекты хранилища Пошардовые счётчики STG и ODS поэтому не сходятся и сходиться не должны — сверять слои можно только через `_dist`. +Ключи ко-локации названы заранее, потому что на них стоит политика соединений из +раздела 6 спеки: обычное соединение разрешено только по ключу ко-локации, всё +прочее — через `GLOBAL`. Значит `dds.session` и `dds.identity_map` шардируются по +`cityHash64(ClientID)`, а `dds.order` и производные от заказа — по +`cityHash64(order_id)`. Ключи витрин появятся вместе с самими витринами. + +Известное ограничение правила «пишем только в `_dist`»: пакетные слои собираются +заменой дневных партиций, а `DROP/REPLACE PARTITION` работает только по локальным +таблицам. Чем и как раскладывать партицию-донор по шардам до замены, здесь не +решено — вопрос встаёт вместе со сборкой DDS, и решать его нужно тогда, а не +задним числом. + ## Служебные колонки Собственные колонки не повторяют имён виртуальных. Виртуальные даёт движок: @@ -76,14 +96,40 @@ keeper, Kafka, каркас сервисов. Объекты хранилища `kafka_timestamp`, рядом — `consumer_host`, имя читавшей ноды: виртуальные колонки его не несут, а после записи в `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). -Метка времени загрузки зовётся `_load_ts`, тип `DateTime64(3)`. В ODS она же -служит колонкой версии `ReplacingMergeTree`. Имя согласовано с каноном служебных -полей соседнего учебного стенда на Greenplum, чтобы словарь был общим у двух -хранилищ; ведущее подчёркивание у технических колонок — распространённая запись, +Движок таблицы сырья — обычный `ReplicatedMergeTree`, `ORDER BY (kafka_partition, +kafka_offset)`: разбор полётов идёт от «какое сообщение», другого ключа у сырья и +нет. Замену версий сюда ставить нельзя — она отменила бы свойство слоя, ради +которого он заведён: повтор доставки в сырье обязан быть виден. + +Метка времени загрузки зовётся `_load_ts`, тип `DateTime64(3)`. Ставится она +один раз, в матвью приёма, и дальше переносится из STG в ODS как есть: колонка +отвечает на вопрос «когда строка приехала в хранилище», а не «когда её +разобрали». В ODS она же служит колонкой версии `ReplacingMergeTree`, и работа у +этой версии ровно одна — схлопнуть повтор доставки. Содержимое у повтора то же +самое, отличается только метка, поэтому какая из двух строк переживёт мерж, +безразлично. Переобработки как стадии у ODS нет вовсе: слой наполняет матвью, а +не пакетное задание, и пакетная работа с партициями начинается выше. + +Имя согласовано с каноном служебных полей соседнего учебного стенда на +Greenplum, чтобы словарь был общим у двух хранилищ; ведущее подчёркивание у +технических колонок — распространённая запись, её же используют Fivetran, Airbyte и Stitch. С правилом выше это не спорит: запрещено совпадать с именами виртуальных колонок, а не носить подчёркивание. @@ -106,6 +152,13 @@ keeper, Kafka, каркас сервисов. Объекты хранилища Цепочка одна: чтец топика → матвью → сырьё STG → матвью разбора → событие и таблица ошибок ODS. +Чтец стоит на обеих нодах и читает одной группой потребителей — имя группы +`clickstream_hits`, и оно одинаково на обеих нодах по построению: DDL идёт +`ON CLUSTER` и макросов в имени не содержит. Разные группы дали бы каждой ноде +полную копию топика, и это отдельная сцена для лабы, а не рабочий режим. Имя +кластера в `ON CLUSTER` и в движке `Distributed` — `clickstream_cluster`, оно +задано в `infra/clickhouse/config.d/cluster.xml`. + **Источник матвью разбора — `stg.hits_raw_dist`, а не локальная таблица.** У матвью две привязки: источник, на вставку в который она срабатывает, и цель, куда пишет. Распределённая таблица — лицо слоя, локальная — его хранилище; @@ -113,14 +166,24 @@ keeper, Kafka, каркас сервисов. Объекты хранилища той же ноде, что читала Kafka, в момент вставки первой матвью — до раскладки по шардам. -**Вставка синхронная.** На пути приёма стоит -`distributed_foreground_insert = 1`. По умолчанию вставка в распределённую +**Вставка фоновая, и окно потери мы принимаем.** Вставка в распределённую таблицу кладёт блок в локальный спул и сразу возвращает управление, а Kafka коммитит офсеты по факту работы матвью — то есть по факту записи в спул. Топик -уже считает сообщение прочитанным, хотя на шарде его нет. Синхронный режим не -добавляет работы, он её не прячет: кусок на шарде будет записан всё равно, -вопрос лишь в том, ждём ли мы этого внутри вставки. Платится ожидание один раз -на блок Kafka в десятки тысяч строк, а не на событие. +уже считает сообщение прочитанным, хотя на шарде его ещё нет: умри нода в этом +промежутке — сообщения не перечитаются. + +Закрывает окно настройка `distributed_foreground_insert = 1`, и на ETL-вставках +Airflow она стоит — там это обычный `SETTINGS` у запроса. На пути приёма её нет, +и по трём причинам. Вставку выполняет фоновый поток Kafka-движка, своего запроса +у него не бывает, так что настройка уровня запроса доехала бы только профилем +пользователя в конфигурации ноды. Синхронный режим связывает шарды: пока второй +недоступен, вставка падает, офсеты не коммитятся, и приём встаёт целиком — тогда +как при фоновом первая нода продолжает принимать и копит спул для соседа. +Платится при этом не одно ожидание на блок, а три распределённые вставки — сырьё, +событие, ошибки, — и все внутри потока-потребителя, что само по себе повод для +ребаланса по таймауту сессии. Против всего этого — окно в сотню миллисекунд на +ноутбуке, где мир пересобирается одной командой. Размен не в пользу настройки, а +компромисс полезнее показать, чем спрятать за галочкой. **Гарантия — «хотя бы один раз», не транзакция.** Падение после записи на шард, но до коммита офсетов даёт повтор при перечитывании. В ODS повтор схлопнет @@ -129,38 +192,52 @@ keeper, Kafka, каркас сервисов. Объекты хранилища Это свойство слоя, а не поломка, — но обещание идемпотентности конвейера к сырому слою не относится. +Оговорка к последнему: у семейства `Replicated*` есть своя дедупликация — блок с +тем же хешем, вставленный повторно, отбрасывается (`insert_deduplicate`). +Удвоение сырья проходит мимо неё только потому, что при повторном чтении +`_load_ts` новый и хеш блока другой. Свойство слоя держится на этом, а не на +отсутствии механизма. + **Матвью разбора две, и их условия обязаны делить поток без зазора и без нахлёста.** Одна забирает годные строки в `ods.event_dist`, вторая — брак в `ods.event_errors_dist`. Строка, подошедшая обеим, задвоится; не подошедшая ни одной — исчезнет молча. Держится это формой: второе условие пишется буквальным отрицанием первого, а сам предикат собирается только из функций, не возвращающих NULL, — иначе трёхзначная логика даст строку, которую не возьмёт ни `условие`, -ни `NOT условие`. +ни `NOT условие`. Те же функции не должны и бросать исключений: упавшая матвью +роняет вставку и останавливает потребление до починки +([ADR 0005](../adr/0005-event-ingestion.md)). Постоянной сверки счётчиков при этом нет и не должно быть. У сырья срок жизни трое суток, а ODS хранит всё, поэтому равенство «сырьё = события + ошибки» -разъедется само; переобработка добавит строк в ODS, перезаливка задвоит сырьё, а -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`). Значение -этой настройки между версиями ClickHouse менялось, поэтому в DDL оно проставлено -явно. +колонке `_load_ts`, снятие — целыми кусками (`ttl_only_drop_parts`). В DDL +значение проставлено явно, чтобы поведение не зависело от умолчания версии. Нарезать сырьё по модельному дню события было бы соблазнительно — он единица переобработки и он же ключ партиции в ODS, — но чистку это ломает. Кусок -снимается целиком, только когда в нём истекли все строки, а обычный мерж внутри -партиции модельного дня склеит куски разного возраста, и данные переживут срок. -В партиции дня загрузки склеивать нечего: все строки в ней одного возраста, и -куски уходят по мере того, как истекает самый свежий из них — то есть сырьё -живёт трое суток с небольшим хвостом, а не ровно трое. +снимается целиком, только когда в нём истекли все строки, а в партиции модельного +дня лежит приехавшее в разное реальное время: мерж склеит куски разного возраста, +самая свежая строка удержит весь кусок, и данные переживут срок неограниченно. + +В партиции дня загрузки склейка идёт точно так же, и строки в ней тоже разного +возраста — но не более чем на сутки, потому что партицию закрывает календарный +день. Отсюда и оценка: сырьё живёт трое суток плюс хвост до суток, а не ровно +трое. Оговорка про сроки: TTL исполняется на мержах, а не по будильнику, так что +«уходит само» здесь обещано, а «уходит вовремя» — нет. Модельного дня среди колонок сырья нет вовсе. Он свойство содержимого, а содержимое разбирает ODS — там `EventDate` и живёт, ключом партиции. Сырьё @@ -182,32 +259,60 @@ D0 и к реальному календарю не привязана; паке ## Таблица ошибок `ods.event_errors` держит строки, не прошедшие строгий приём, вместе с их сырым -текстом и метаданными доставки. Ключи её собственные, потому что у брака нет -разобранных полей: шардируется `cityHash64` сырой строки — `ClientID` у строки, -которая не разобралась, взять неоткуда; нарезается по дню загрузки, как и сырьё; -живёт месяц. Дольше сырья — намеренно: если брак истекает вместе с ним, +текстом, метаданными доставки и классом брака. Ключи её собственные, потому что у +брака нет разобранных полей: шардируется `cityHash64` сырой строки — `ClientID` у +строки, которая не разобралась, взять неоткуда; нарезается по дню загрузки, как и +сырьё; живёт месяц. Дольше сырья — намеренно: если брак истекает вместе с ним, разбираться к моменту разбирательства будет уже нечем. +Класс брака лежит в колонке `error_class` типа `LowCardinality(String)`. Без неё +в таблице копятся строки «что-то не так» без ответа на «что именно», а витрине +качества не на что опереться. Сами классы перечислены в +[ADR 0005](../adr/0005-event-ingestion.md) и проверяются по порядку, потому что +пересекаются: сообщение, не являющееся объектом JSON, проваливает заодно и сверку +набора ключей — `JSONExtractKeys` от скаляра даёт пустой массив. Побеждает первый +совпавший класс, тем же приёмом, что `mismatch_class` в витрине сверки. + +Движок — обычный `ReplicatedMergeTree`, без замены версий: схлопывать брак не по +чему, у него нет ключа сущности. `ORDER BY` — `(error_class, kafka_partition, +kafka_offset)`: смотрят такую таблицу от класса, а внутри класса — по координатам +доставки. + ## Раскладка DDL -Файлы лежат в `sql/ddl/` и применяются по порядку имён. Один файл — это слой и -роль: статичные объекты отдельно от матвью. +Файлы лежат в `sql/ddl/` и применяются по порядку имён. Сначала все статичные +объекты, потом матвью — тогда к моменту создания матвью её цель уже существует. | Файл | Что в нём | |---|---| | `00-databases.sql` | базы слоёв | | `10-stg-tables.sql` | Kafka-таблица, локальная и распределённая таблицы сырья | -| `11-stg-views.sql` | матвью, наполняющая сырьё из Kafka | | `20-ods-tables.sql` | типизированное событие и таблица ошибок | -| `21-ods-views.sql` | матвью разбора: сырьё в событие и в ошибки | +| `30-ods-views.sql` | матвью разбора: сырьё в событие и в ошибки | +| `40-stg-views.sql` | матвью приёма: чтец в сырьё | -Матвью принадлежит слою своей цели, а не источника: разбор из STG в ODS лежит -среди файлов ODS, потому что наполняет ODS. +Порядок задают два правила. Первое: матвью принадлежит слою своей цели, а не +источника, — разбор из STG в ODS лежит среди файлов ODS, потому что наполняет +ODS. Второе: матвью приёма создаётся последней из всех, и потому нарушает +нумерацию слоёв. Kafka-движок начинает читать топик ровно тогда, когда к нему +привязывают первую матвью; создай её раньше разбора — и всё, что доедет в +зазоре, ляжет в сырьё и не попадёт в ODS никуда, ни в событие, ни в ошибки. На +пустом топике зазор безвреден, поэтому первый прогон о нём не скажет. Проснётся +он, когда тома ClickHouse снесены, а данные Kafka целы, — то есть на обычной +отладке. Применение — двумя одноразовыми сервисами при `make up`, по образцу уже работающих `airflow-init` и `superset-init`. Сначала `kafka-init` создаёт топик `hits` с двумя партициями, затем `clickhouse-init` дожидается его завершения и -применяет файлы с ноды 1, `ON CLUSTER`. Порядок страхует от автосоздания топика +применяет файлы с ноды 1, `ON CLUSTER`. + +Образцы копируются не целиком, и в двух местах. `clickhouse-init` обязан ждать +готовности **обеих** нод: `ON CLUSTER` ждёт исполнения на всех хостах и по +таймауту бросает, а `airflow-init` ждёт только первую ноду, `superset-init` — +только вторую. И второе: оба образца переживают `make up --wait` лишь потому, что +от них зависят долгоживущие сервисы; у пары `kafka-init` / `clickhouse-init` +таких зависимых нет, и как поведёт себя `--wait` с одноразовым сервисом без них — +проверяется при исполнении #37. Порядок страхует от автосоздания топика брокером с одной партицией: у потребителя librdkafka разрешение на автосоздание по умолчанию выключено, так что случиться это не обязано, но урок «обе ноды читают топик» умирает тихо, и полагаться на умолчание клиента здесь не стоит. @@ -234,3 +339,42 @@ D0 и к реальному календарю не привязана; паке Слои 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; до тех пор это предположения, а не +знание. diff --git a/docs/specs/2026-07-30-stand-v2-realism.md b/docs/specs/2026-07-30-stand-v2-realism.md index 378c43b..d725223 100644 --- a/docs/specs/2026-07-30-stand-v2-realism.md +++ b/docs/specs/2026-07-30-stand-v2-realism.md @@ -312,9 +312,13 @@ CSV в репозитории (`data/catalog/products.csv`: `sku`, `name`, `cate `cityHash64(order_id)`, сырьё STG — `cityHash64(сырой строки)`; полный список и доводы — в [доке хранилища](../architecture/storage.md). Урок: «какая нода читала топик — меняется между прогонами, куда легли данные — нет». -- **Приём строгий**: обязательные поля разбираются как `Nullable`, а набор - ключей сообщения сверяется с контрактным; строка с NULL среди обязательных - полей или с разошедшимся набором ключей уходит в `*_errors`. Сверка ключей — +- **Приём строгий**: пять опорных колонок — `WatchID`, `VisitID`, `ClientID`, + `EventDate`, `UTCEventTime` — разбираются как `Nullable`, а набор ключей + сообщения сверяется с контрактным; строка с NULL среди опорных колонок + или с разошедшимся набором ключей уходит в `*_errors`. Опорными выбраны те, на + которых стоят ключ сортировки, партиция и дедупликация; остальные сорок две + достаются обычными типами — сорок семь проверок на NULL превратили бы матвью в + простыню, а присутствие и так целиком закрыто сверкой ключей. Сверка ключей — не добавка: у массивов NULL не бывает, и пропавшее поле-массив иначе неотличимо от пустого по смыслу. На входе разбора нет вовсе — Kafka-таблица читает сообщение байтами, строгость целиком в матвью ODS @@ -331,8 +335,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 неатомарна — дашборд в середине прогона честно моргает (это осознанная цена, не баг). @@ -374,7 +379,7 @@ README. | Слой | Объект | Что это | |---|---|---| | Kafka | `hits`, `orders` | два топика, по 2 партиции | -| STG | `stg.hits_raw_kafka`, `stg.hits_raw` + MV; то же для orders | сырые строки, Kafka Engine на обеих нодах | +| STG | `stg.hits_raw_kafka`, `stg.hits_raw` + MV; для orders — развилка этапа 3, см. раздел 12 | сырые строки, Kafka Engine на обеих нодах | | ODS | `ods.event` (+`_errors`) | типизированное широкое событие, ReplacingMergeTree | | ODS | `ods.order_snapshot` (+`_errors`) | слепки заказов как приехали, партиция по `snapshot_date`, без дедупа | | DDS | `dds.session` | сборка сессий из событий (наследник `dds.click`) | @@ -409,11 +414,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` @@ -427,10 +432,10 @@ 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`, локальное соединение по ключу ко-локации). - `events_enriched_v`, `top_pages_daily_v`, `session_overview_v`, @@ -498,7 +503,7 @@ v2 стартует пустым, поэтому объём ниже — это `cancelled`) — на статусе `paid` стоит «дыхание» выручки; сужение до пары `created`/`cancelled` — запасной ход, если генератор заказов окажется дороже ожиданий. Связка: если срез 1 сработает, определение выручки в - `v_revenue_daily` придётся сменить с «заказы в статусе `paid`» на «все + `revenue_daily_v` придётся сменить с «заказы в статусе `paid`» на «все неотменённые заказы». - **Отступление от порядка #15**: резолюция предписывала резать в порядке 1 → 2 → 3, спека применяет 3 и 2, а 1 держит в резерве. Довод: срезы 3 и 2 @@ -569,8 +574,15 @@ smoke-проверки, а не «дашборд зелёный». Это мин - Точная форма `ORDER BY` ODS-таблиц (выражение `intHash32` в ключе ReplacingMergeTree). - `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение (ADR 0005). -- Матвью, привязанная к локальной таблице, срабатывает, когда строки приходят - вставкой через `Distributed`: на этом стоит цепочка STG → ODS (ADR 0005). + Проверять первым, до написания DDL: если сообщения склеятся, переделывать + придётся решение, а не запрос. Запасной формат — `LineAsString`. +- Тип виртуальной колонки `_timestamp` у Kafka-движка: обнуляемость и + разрядность (секунды против миллисекунд) — от этого зависит объявление + `kafka_timestamp` в таблице сырья. +- Матвью с источником-`Distributed` срабатывает на вставку именно в эту + распределённую таблицу, до раскладки по шардам: на этом стоит цепочка + STG → ODS (ADR 0005). Проверено владельцем на рабочих проектах, в документации + ClickHouse этот случай не описан. - Форма именованного кортежа в `JSONExtract` с `Nullable`-членами — ею сворачиваются 47 вызовов в один, если разбор окажется дорогим (ADR 0005). - Размер артефакта эталонного мира после пересборки. @@ -628,8 +640,16 @@ v2, этап 0). `count()`, а в бою — нет; частый вопрос на собеседованиях; - два способа принять топик, рядом на одном стенде: сырьё байтами с разбором функциями (`hits`, [ADR 0005](../adr/0005-event-ingestion.md)) против - типизированного чтеца с `kafka_handle_error_mode` (заказы, этап 3) — - сравнение цены и наблюдаемости как задание; + типизированного чтеца с `kafka_handle_error_mode` — сравнение цены и + наблюдаемости как задание. **Развилка этапа 3, не решена**: типизированный + чтец идёт заказам только вместе с ответом на вопрос, нужен ли им слой сырья. + Нужен — и чтецов на один топик станет два, а эту схему ADR 0005 отверг; не + нужен — и слои перестают быть единообразными. Разбирать грилингом, когда + дойдём до заказов; +- матвью как рабочий механизм, а не диковина: их видно на приёме и на сборке + ODS, а пакетная работа начинается выше. Отдельным заданием — как читать из ODS + последние версии, через `FINAL` или оконной функцией: что нагляднее, решаем на + месте; - лекция про идентичность «как в бою»: `setUserID` и first-party id, детерминированная против вероятностной склейки, identity graph, кросс-девайс, CDP — с рамкой «мы склеили через транзакции, потому что трекер -- 2.54.0 From f1254ce79d91af60159d6f508bc3c0bb72edcbaf Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Wed, 5 Aug 2026 21:31:32 +0300 Subject: [PATCH 4/5] =?UTF-8?q?docs(handoff):=20=D0=BF=D0=B5=D1=80=D0=B5?= =?UTF-8?q?=D0=B4=D0=B0=D1=87=D0=B0=20=D0=BA=D0=BE=D0=BD=D1=82=D0=B5=D0=BA?= =?UTF-8?q?=D1=81=D1=82=D0=B0=20=D1=83=D0=B4=D0=B0=D0=BB=D0=B5=D0=BD=D0=B0?= =?UTF-8?q?=20=D0=BA=D0=B0=D0=BA=20=D0=BE=D1=82=D1=80=D0=B0=D0=B1=D0=BE?= =?UTF-8?q?=D1=82=D0=B0=D0=B2=D1=88=D0=B0=D1=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - список «чего не переоткрывать» устарел за день: снятая настройка синхронной вставки и прежняя раскладка DDL остались в нём как принятые решения, а документ, утверждающий решённым переехавшее, — ловушка. - Что: - удалён .scratch/handoffs/2026-08-04-storage-conventions.md; содержание живёт в ADR 0005 и 0006, доке хранилища и тикетах #37 и #43. - Проверка: - make config-test Co-Authored-By: Claude Opus 5 --- .../2026-08-04-storage-conventions.md | 91 ------------------- 1 file changed, 91 deletions(-) delete mode 100644 .scratch/handoffs/2026-08-04-storage-conventions.md diff --git a/.scratch/handoffs/2026-08-04-storage-conventions.md b/.scratch/handoffs/2026-08-04-storage-conventions.md deleted file mode 100644 index 27a591e..0000000 --- a/.scratch/handoffs/2026-08-04-storage-conventions.md +++ /dev/null @@ -1,91 +0,0 @@ -# Handoff: конвенции хранилища и приём событий (#37) - -Дата: 4 августа 2026 года. Ветка: `docs/37-storage-conventions`. - -## Что это было - -Сессия началась с вопроса «хватает ли информации, чтобы взять #37 в работу». -Оказалось, что нет: тикет молча опирался на конвенции хранилища, которых в -проекте не существовало. Дальше шёл грилинг пяти развилок, затем два холодных -ревью — по линии дефектов и по линии уместности, — затем правки по их находкам. - -Кода не написано ни строки. Итог сессии — принятые решения и приведённые в -соответствие постановки. - -## Где лежат решения - -Ничего из перечисленного здесь не пересказывается — читать по адресам: - -- `docs/adr/0005-event-ingestion.md` — как принимаем события: чтец читает топик - байтами, разбор идёт функциями в матвью ODS. Там же отвергнутые варианты и - цена решения. -- `docs/adr/0006-object-naming.md` — суффикс вида в именах объектов. -- `docs/architecture/storage.md` — рабочий справочник: конвенции имён и - служебных колонок, раскладка по шардам, путь в keeper, срок жизни сырья, - свойства приёма, раскладка файлов DDL, карта таблиц. -- `docs/specs/2026-07-30-stand-v2-realism.md` — правлены разделы 6, 7, 9, 11, - «Витрины DM» и опорные точки для лекций. -- `docs/specs/2026-08-01-generator.md` — раздел 4 получил контракт транспорта - «одно событие — одно сообщение Kafka». -- Тикеты #37 и #43 переписаны целиком; у каждого сверху комментарий с разбором - того, что изменилось против исходной постановки, — читать `tea issues 37 - --comments`. - -Коммит с документами — `319db30`, он же единственный на ветке помимо этого -файла. - -## Чего не переоткрывать - -Всё ниже прошло грилинг с веером вариантов и записано с доводами. Если решение -кажется странным — сначала прочитать довод, а не начинать заново: - -- имена объектов: суффикс вида, а не префикс и не `_all`; -- служебная колонка загрузки `_load_ts`, метаданные доставки без ведущего - подчёркивания; -- нарезка сырья по дню загрузки, а не по модельному дню; TTL трое суток; -- модельного дня в STG нет вовсе; -- формат чтеца `RawBLOB`, режим `kafka_handle_error_mode` не используется; -- строгий приём: сверка набора ключей плюс `Nullable` на пяти опорных колонках, - не на сорока семи; -- источник матвью разбора — `stg.hits_raw_dist`, пишем только в `_dist`, - `distributed_foreground_insert = 1`; -- путь в keeper `/clickhouse/tables/{shard}/{database}/{table}`, без `{uuid}`. - -## Что осталось - -По убыванию веса: - -1. **Проверки на живом стенде**, две из них внесены в раздел 11 мастер-спеки: - `RawBLOB` даёт ровно одну строку на сообщение; форма именованного кортежа в - `JSONExtract` с `Nullable`-членами. Не внесены, но всплыли в ревью: как - `make up --wait` поведёт себя с одноразовым сервисом, у которого нет - зависимых долгоживущих (существующие `airflow-init` и `superset-init` - переживают `--wait` только за счёт зависимостей), и как аккуратнее навесить - `distributed_foreground_insert` на путь приёма — это настройка уровня - запроса, вероятно через профиль пользователя в конфиге. -2. **PR не открыт.** Ветка запушена, тело PR писать с английским `Closes #NN`, - иначе Gitea задачу не закроет (см. `docs/agents/issue-tracker.md`). -3. **Третий холодный проход** — по желанию владельца. Оба документа переписаны - после ревью существенно, а разделы про keeper, таблицу ошибок и свойства - приёма ревьюеры не видели вовсе. -4. **Реализация #37** — собственно этап, ради которого всё затевалось. - -## Что стоит знать про ход работы - -- Владелец правит рекомендации по существу и часто оказывается прав: из пяти - развилок три пересматривались по его возражениям. Предлагать вариант с - доводом, а не спрашивать «как сделать», — и быть готовым, что довод разберут. -- Спека не выбита в камне: менять её аргументированно можно и нужно, но - обсуждая с владельцем, а не молча. -- Трекер — Gitea, только через `tea`, при проблемах с прокси префикс - `NO_PROXY='*'`. Тикеты читать с `--comments`: хвосты живут там. -- Сверка спорных API ClickHouse — через MCP Context7 до кода, не после. - -## Suggested skills - -- `brainstorm-with-docs` — если всплывёт новая развилка. Именно им шла эта - сессия; формат «веер вариантов, потом конвергенция» владельцу привычен. -- `claude-subagent-playbook` — когда #37 пойдёт в реализацию: конвейер с - делегированием механической части и слепым ревью. -- `conventional-commits` — обязателен при любом коммите в этом репозитории. -- `code-review` — для ревью изменений перед приёмкой. -- 2.54.0 From 923ebad80e21552e4afb5c96d5fd3be048041d14 Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Wed, 5 Aug 2026 21:53:03 +0300 Subject: [PATCH 5/5] =?UTF-8?q?docs(storage):=20=D1=81=D0=B2=D1=8F=D0=B7?= =?UTF-8?q?=D0=BD=D0=BE=D1=81=D1=82=D1=8C=20=D0=B2=D0=BE=D1=81=D1=81=D1=82?= =?UTF-8?q?=D0=B0=D0=BD=D0=BE=D0=B2=D0=BB=D0=B5=D0=BD=D0=B0,=20=D1=84?= =?UTF-8?q?=D0=BE=D1=80=D0=BC=D0=B0=20=D0=B4=D0=B0=D1=82=20=D0=BD=D0=B0=20?= =?UTF-8?q?=D0=BF=D1=80=D0=BE=D0=B2=D0=BE=D0=B4=D0=B5=20=D0=B7=D0=B0=D0=B4?= =?UTF-8?q?=D0=B0=D0=BD=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - холодное ревью связности нашло девять мест, где вставленный текст спорит с соседним; отдельно вскрылось, что представление дат в JSON не зафиксировано нигде, а #43 обязан его знать раньше, чем #41 напишет сериализатор. - Что: - гарантия приёма переписана: после снятия синхронной вставки «хотя бы один раз» стало неправдой — есть и окно потери, и окно дубля. - критерий выбора пяти опорных колонок приведён к списку, который он порождает; `CounterID` оговорён отдельно. - «переобработки у ODS нет вовсе» смягчено до пакетной: ручная вставка из сырья в пределах окна возможна. - в спеку генератора добавлена форма дат на проводе — ISO-8601, с доводом от читаемости слоя сырья. - убраны осиротевшая фраза про порядок сервисов, дубль порядка классов брака, устаревшая датировка сверки и ещё три следа вставок. - Проверка: - make config-test Co-Authored-By: Claude Opus 5 --- docs/adr/0005-event-ingestion.md | 16 ++++--- docs/architecture/storage.md | 51 +++++++++++++---------- docs/specs/2026-07-30-stand-v2-realism.md | 18 ++++++-- docs/specs/2026-08-01-generator.md | 11 +++++ 4 files changed, 65 insertions(+), 31 deletions(-) diff --git a/docs/adr/0005-event-ingestion.md b/docs/adr/0005-event-ingestion.md index 2959f83..a45ee71 100644 --- a/docs/adr/0005-event-ingestion.md +++ b/docs/adr/0005-event-ingestion.md @@ -41,8 +41,11 @@ JSON. цены и пользы: единственный производитель топика — собственный генератор, сериализующий из контракта по объявленным типам, поэтому неверный тип может прийти только из руки, а сорок семь проверок на NULL превратили бы матвью в -простыню. Пять выбраны по последствию: на них стоят ключ сортировки, партиция и -дедупликация, и их порча отравляет всё ниже по течению. +простыню. Пять выбраны по последствию: это идентификаторы события, визита и +посетителя, дата партиции и метка времени, по которой события упорядочиваются +внутри сессии, — порча любой отравляет всё ниже по течению. `CounterID` +формально тоже входит в ключ сортировки, но на стенде он константа, и NULL там +взяться неоткуда. Присутствие иначе и не проверить. `Nullable` в ClickHouse не оборачивает составные типы: `Nullable(Array)` запрещён, а @@ -133,7 +136,9 @@ contract-тест из #43 сюда не дотягивается — он ср ## Что проверено -По документации ClickHouse через MCP Context7, 3 августа 2026 года. +По документации ClickHouse через MCP Context7: основная сверка — 3 августа +2026 года, перепроверка после правок — 5 августа. Датировка важна: 5 августа +утверждение про `Nullable(Tuple)` развернулось на противоположное. При режиме `stream` движок отдаёт `_raw_message` и `_error` только для сообщений, которые не разобрались, и оставляет их пустыми для разобранных. @@ -159,9 +164,8 @@ contract-тест из #43 сюда не дотягивается — он ср распределённую таблицу — блок она видит до разрезания по шардам. Проверено владельцем на рабочих проектах; на стенде подтверждается заодно с приёмкой #37. -На живом стенде проверяется при исполнении #37. Первые два пункта внесены в -раздел 11 спеки как несущие; остальные — однострочные `SELECT`, их довольно -прогнать заодно: +На живом стенде проверяется при исполнении #37. Первые два внесены в раздел 11 +спеки; остальные — однострочные `SELECT`, их довольно прогнать заодно: - `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение. Проверять это нужно первым и до написания DDL: формулировка «читает вход в одно значение» diff --git a/docs/architecture/storage.md b/docs/architecture/storage.md index 47847eb..436943d 100644 --- a/docs/architecture/storage.md +++ b/docs/architecture/storage.md @@ -2,7 +2,8 @@ Документ описывает сторону ClickHouse: как называются объекты, какие служебные колонки у них общие, чем нарезаны и сколько живут данные, как устроен приём и из -каких файлов собирается DDL. Здесь же карта таблиц, которая растёт с этапами. +каких файлов собирается DDL. Здесь же карта таблиц, которая растёт с этапами, и +раздел «Что проверено» — чему в этом тексте верить и на каком основании. **Что здесь описано и чего ещё нет.** Собран этап 1: кластер из двух шардов, keeper, Kafka, каркас сервисов. Объекты хранилища и механизм применения DDL @@ -79,11 +80,12 @@ keeper, Kafka, каркас сервисов. Объекты хранилища `cityHash64(ClientID)`, а `dds.order` и производные от заказа — по `cityHash64(order_id)`. Ключи витрин появятся вместе с самими витринами. -Известное ограничение правила «пишем только в `_dist`»: пакетные слои собираются -заменой дневных партиций, а `DROP/REPLACE PARTITION` работает только по локальным -таблицам. Чем и как раскладывать партицию-донор по шардам до замены, здесь не -решено — вопрос встаёт вместе со сборкой DDS, и решать его нужно тогда, а не -задним числом. +Открытый вопрос на будущее — не сама замена партиций: операции с ними по +локальным таблицам правило разрешает прямо. Вопрос в шаге до неё. Партиция-донор +должна быть уже разложена по шардам по тому же ключу, а разложить её можно +только вставкой через распределённую таблицу — значит у каждой пакетной сущности +появится вторая пара объектов, и имени для неё конвенция пока не даёт. Решать +это вместе со сборкой DDS, а не задним числом. ## Служебные колонки @@ -104,8 +106,8 @@ keeper, Kafka, каркас сервисов. Объекты хранилища (раздел 11 спеки) — записать её в необнуляемый тип значит либо уронить приём на первом сообщении, либо получить тихие нули за 1970 год. -Заполняются обе группы колонок выражением в `SELECT` матвью приёма, а не -`DEFAULT` в таблице. Для `consumer_host` это обязательно: `DEFAULT hostName()` +Заполняются все они выражением в `SELECT` матвью приёма, а не `DEFAULT` в +таблице. Для `consumer_host` это обязательно: `DEFAULT hostName()` сработал бы на шарде-получателе и назвал бы не ту ноду, которая читала топик, — то есть колонка молча отвечала бы на другой вопрос. @@ -124,8 +126,10 @@ kafka_offset)`: разбор полётов идёт от «какое сооб разобрали». В ODS она же служит колонкой версии `ReplacingMergeTree`, и работа у этой версии ровно одна — схлопнуть повтор доставки. Содержимое у повтора то же самое, отличается только метка, поэтому какая из двух строк переживёт мерж, -безразлично. Переобработки как стадии у ODS нет вовсе: слой наполняет матвью, а -не пакетное задание, и пакетная работа с партициями начинается выше. +безразлично. Пакетной переобработки у ODS нет: слой наполняет матвью, а не +задание Airflow, и работа с партициями начинается выше. Переделать разобранное +руками можно — вставкой из сырья с фильтром по `_load_ts`, в пределах +трёхсуточного окна; ничья по версии разрешается в пользу вставленного позже. Имя согласовано с каноном служебных полей соседнего учебного стенда на Greenplum, чтобы словарь был общим у двух хранилищ; ведущее подчёркивание у @@ -185,8 +189,13 @@ Airflow она стоит — там это обычный `SETTINGS` у зап ноутбуке, где мир пересобирается одной командой. Размен не в пользу настройки, а компромисс полезнее показать, чем спрятать за галочкой. -**Гарантия — «хотя бы один раз», не транзакция.** Падение после записи на шард, -но до коммита офсетов даёт повтор при перечитывании. В ODS повтор схлопнет +**Гарантии нет ни в одну сторону — есть два узких окна.** Окно потери описано +выше: нода умерла между коммитом офсетов и сбросом спула. Окно дубля +противоположное: нода умерла после записи на шард, но до коммита офсетов, и при +перечитывании сообщение приедет второй раз. Сказать про такой приём «хотя бы +один раз» нельзя — это обещало бы, что потерь не бывает, а они возможны. + +Дубль ниже по течению ведёт себя по-разному. В ODS его схлопнет `ReplacingMergeTree`, а сырьё дедупа не имеет вовсе: перезаливка модельного дня честно удваивает `count()` в STG, и живёт эта пара до истечения срока хранения. Это свойство слоя, а не поломка, — но обещание идемпотентности конвейера к @@ -267,11 +276,8 @@ D0 и к реальному календарю не привязана; паке Класс брака лежит в колонке `error_class` типа `LowCardinality(String)`. Без неё в таблице копятся строки «что-то не так» без ответа на «что именно», а витрине -качества не на что опереться. Сами классы перечислены в -[ADR 0005](../adr/0005-event-ingestion.md) и проверяются по порядку, потому что -пересекаются: сообщение, не являющееся объектом JSON, проваливает заодно и сверку -набора ключей — `JSONExtractKeys` от скаляра даёт пустой массив. Побеждает первый -совпавший класс, тем же приёмом, что `mismatch_class` в витрине сверки. +качества не на что опереться. Сами классы, их порядок и довод, почему порядок +обязателен, — в [ADR 0005](../adr/0005-event-ingestion.md). Движок — обычный `ReplicatedMergeTree`, без замены версий: схлопывать брак не по чему, у него нет ключа сущности. `ORDER BY` — `(error_class, kafka_partition, @@ -304,7 +310,11 @@ ODS. Второе: матвью приёма создаётся последне Применение — двумя одноразовыми сервисами при `make up`, по образцу уже работающих `airflow-init` и `superset-init`. Сначала `kafka-init` создаёт топик `hits` с двумя партициями, затем `clickhouse-init` дожидается его завершения и -применяет файлы с ноды 1, `ON CLUSTER`. +применяет файлы с ноды 1, `ON CLUSTER`. Этот порядок страхует от автосоздания +топика брокером с одной партицией: у потребителя librdkafka разрешение на +автосоздание по умолчанию выключено, так что случиться это не обязано, но урок +«обе ноды читают топик» умирает тихо, и полагаться на умолчание клиента здесь +не стоит. Образцы копируются не целиком, и в двух местах. `clickhouse-init` обязан ждать готовности **обеих** нод: `ON CLUSTER` ждёт исполнения на всех хостах и по @@ -312,10 +322,7 @@ ODS. Второе: матвью приёма создаётся последне только вторую. И второе: оба образца переживают `make up --wait` лишь потому, что от них зависят долгоживущие сервисы; у пары `kafka-init` / `clickhouse-init` таких зависимых нет, и как поведёт себя `--wait` с одноразовым сервисом без них — -проверяется при исполнении #37. Порядок страхует от автосоздания топика -брокером с одной партицией: у потребителя librdkafka разрешение на автосоздание -по умолчанию выключено, так что случиться это не обязано, но урок «обе ноды -читают топик» умирает тихо, и полагаться на умолчание клиента здесь не стоит. +проверяется при исполнении #37. Повторный `make up` поверх живого тома проходит зелёным: весь DDL идёт через `CREATE ... IF NOT EXISTS`. Оборотная сторона — изменённый объект тем же diff --git a/docs/specs/2026-07-30-stand-v2-realism.md b/docs/specs/2026-07-30-stand-v2-realism.md index d725223..e5c3e5c 100644 --- a/docs/specs/2026-07-30-stand-v2-realism.md +++ b/docs/specs/2026-07-30-stand-v2-realism.md @@ -315,8 +315,11 @@ CSV в репозитории (`data/catalog/products.csv`: `sku`, `name`, `cate - **Приём строгий**: пять опорных колонок — `WatchID`, `VisitID`, `ClientID`, `EventDate`, `UTCEventTime` — разбираются как `Nullable`, а набор ключей сообщения сверяется с контрактным; строка с NULL среди опорных колонок - или с разошедшимся набором ключей уходит в `*_errors`. Опорными выбраны те, на - которых стоят ключ сортировки, партиция и дедупликация; остальные сорок две + или с разошедшимся набором ключей уходит в `*_errors`. Опорными выбраны те, + чья порча отравляет всё ниже по течению: идентификаторы события, визита и + посетителя, дата партиции и метка времени, по которой события упорядочиваются + внутри сессии. `CounterID` формально тоже в ключе сортировки, но на стенде он + константа, и NULL там взяться неоткуда. Остальные сорок две достаются обычными типами — сорок семь проверок на NULL превратили бы матвью в простыню, а присутствие и так целиком закрыто сверкой ключей. Сверка ключей — не добавка: у массивов NULL не бывает, и пропавшее поле-массив иначе @@ -379,7 +382,7 @@ README. | Слой | Объект | Что это | |---|---|---| | Kafka | `hits`, `orders` | два топика, по 2 партиции | -| STG | `stg.hits_raw_kafka`, `stg.hits_raw` + MV; для orders — развилка этапа 3, см. раздел 12 | сырые строки, 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`) | @@ -393,6 +396,15 @@ README. суффиксов; она же задаёт служебные колонки, нарезку и срок хранения сырья — см. [доку хранилища](../architecture/storage.md). +Как принимаются заказы — развилка этапа 3, и она не решена. Событиям выбран +приём сырья байтами с разбором функциями ([ADR 0005](../adr/0005-event-ingestion.md)); +заказам этот же способ идёт только вместе с ответом на вопрос, нужен ли им слой +сырья вообще — у них слепок, а не поток. Нужен — и типизированный чтец даст двух +чтецов на один топик, а такую схему ADR 0005 отверг; не нужен — и слои +перестают быть единообразными. Разбирать грилингом, когда дойдём до заказов; +как учебное сравнение двух способов приёма это записано и в опорных точках +раздела 12. + Состав служебных колонок задаёт дока хранилища. Спеке важны два следствия: `ods.event` и `ods.order_snapshot` получают метку загрузки `_load_ts`, и у `ods.event` она же служит колонкой версии ReplacingMergeTree; а таблицы diff --git a/docs/specs/2026-08-01-generator.md b/docs/specs/2026-08-01-generator.md index 0528318..1268582 100644 --- a/docs/specs/2026-08-01-generator.md +++ b/docs/specs/2026-08-01-generator.md @@ -261,6 +261,17 @@ ([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). Смена библиотеки меняет канонические байты, поэтому проходит как -- 2.54.0