From daf13384a8c3ef4de24e0c02e8e9535983ed2913 Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Thu, 6 Aug 2026 08:04:25 +0300 Subject: [PATCH] =?UTF-8?q?feat(stg):=20DDL-=D0=B1=D1=83=D1=82=D1=81=D1=82?= =?UTF-8?q?=D1=80=D0=B0=D0=BF,=20=D1=82=D0=BE=D0=BF=D0=B8=D0=BA=20hits=20?= =?UTF-8?q?=D0=B8=20=D0=BF=D1=80=D0=B8=D1=91=D0=BC=20=D1=81=D1=8B=D1=80?= =?UTF-8?q?=D1=8C=D1=8F=20=D0=BE=D0=B1=D0=B5=D0=B8=D0=BC=D0=B8=20=D0=BD?= =?UTF-8?q?=D0=BE=D0=B4=D0=B0=D0=BC=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Зачем: стенду нужен воспроизводимый холодный старт, при котором схема хранилища и топик появляются сами, а сырьё из Kafka доезжает в STG обеими нодами кластера — без ручных шагов между `make clean` и рабочим приёмом. Что: - `sql/ddl/` — три файла, применяются по порядку имён: базы `stg` и `ods`, Kafka-чтец `hits_raw_kafka` формата RawBLOB, реплицируемая `hits_raw_rep` с окном TTL в трое суток, распределённая `hits_raw_dist` и матвью `hits_raw_mv`, переносящая сырьё вместе с метаданными доставки. - `compose.yaml` — службы `kafka-init` (топик `hits` на две партиции, с ремонтом уже созданного однопартиционного) и `clickhouse-init` (применяет `/ddl/*.sql`); `hostname:` у обеих нод, чтобы `hostName()` отдавал имя узла, а не идентификатор контейнера; `airflow-init` зависит от `clickhouse-init` — без зависимого успешный одноразовый сервис считается упавшим для `--wait`. - Доки: конвенции и раздел «Что проверено» в справочнике хранилища, указатели и границы обещаний в ADR 0005, снятые пункты в разделе 11 спеки. Проверка: `make lint`, `make typecheck`, `make config-test`, `make smoke` (25 проверок), `make smoke-guards` — зелёные. Приёмочный прогон с чистого тома подтвердил все пять критериев #37: холодный старт и идемпотентный повтор, две партиции у `hits`, метаданные доставки у доехавшего сообщения, обе партиции на обеих потребляющих нодах в одном прогоне, некорректный JSON лежит сырым и приём не встаёт. Известная граница: RawBLOB молча теряет запись с пустым значением и запись-надгробие; принято как свойство, замер и довод — в справочнике хранилища. Co-Authored-By: Claude Opus 5 --- compose.yaml | 79 ++++++++++++ docs/adr/0005-event-ingestion.md | 37 +++--- docs/architecture/storage.md | 139 +++++++++++++++------- docs/specs/2026-07-30-stand-v2-realism.md | 15 +-- sql/ddl/00-databases.sql | 15 +++ sql/ddl/10-stg-tables.sql | 87 ++++++++++++++ sql/ddl/40-stg-views.sql | 35 ++++++ 7 files changed, 344 insertions(+), 63 deletions(-) create mode 100644 sql/ddl/00-databases.sql create mode 100644 sql/ddl/10-stg-tables.sql create mode 100644 sql/ddl/40-stg-views.sql diff --git a/compose.yaml b/compose.yaml index 40f638d..74edbfc 100644 --- a/compose.yaml +++ b/compose.yaml @@ -94,6 +94,7 @@ services: # Инициатор DDL и будущее подключение Airflow. clickhouse-01: <<: *clickhouse-common + hostname: clickhouse-01 ports: - "127.0.0.1:${CLICKHOUSE_01_HTTP_PORT:-28123}:8123" - "127.0.0.1:${CLICKHOUSE_01_TCP_PORT:-29000}:9000" @@ -106,6 +107,7 @@ services: # Будущее подключение Superset. clickhouse-02: <<: *clickhouse-common + hostname: clickhouse-02 ports: - "127.0.0.1:${CLICKHOUSE_02_HTTP_PORT:-28124}:8123" - "127.0.0.1:${CLICKHOUSE_02_TCP_PORT:-29001}:9000" @@ -150,6 +152,79 @@ services: retries: 20 start_period: 30s + kafka-init: + image: ${KAFKA_IMAGE:-apache/kafka:4.3.1} + restart: "no" + depends_on: + kafka: + condition: service_healthy + entrypoint: ["/bin/bash", "-ec"] + command: + - | + topic_partition_count() { + /opt/kafka/bin/kafka-topics.sh \ + --bootstrap-server kafka:9092 \ + --describe \ + --topic hits | \ + sed -n '1s/.*PartitionCount: \([0-9][0-9]*\).*/\1/p' + } + + /opt/kafka/bin/kafka-topics.sh \ + --bootstrap-server kafka:9092 \ + --create \ + --if-not-exists \ + --topic hits \ + --partitions 2 \ + --replication-factor 1 + + partition_count="$$(topic_partition_count)" + # Брокер разрешает автосоздание, хотя потребитель ClickHouse его не запрашивает. + # На сохранённом томе другой клиент мог уже создать hits с одним разделом. + if [[ "$$partition_count" == "1" ]]; then + /opt/kafka/bin/kafka-topics.sh \ + --bootstrap-server kafka:9092 \ + --alter \ + --topic hits \ + --partitions 2 + partition_count="$$(topic_partition_count)" + fi + if [[ "$$partition_count" != "2" ]]; then + printf 'ОШИБКА: у топика hits разделов: %s, ожидалось: 2.\n' \ + "$${partition_count:-неизвестно}" >&2 + exit 1 + fi + + # DDL применяется с ноды 1 по порядку имён файлов после готовности всего кластера. + clickhouse-init: + image: ${CLICKHOUSE_IMAGE:-clickhouse/clickhouse-server:26.3.17.56} + restart: "no" + depends_on: + kafka-init: + condition: service_completed_successfully + clickhouse-01: + condition: service_healthy + clickhouse-02: + condition: service_healthy + environment: + LC_ALL: C + entrypoint: ["/bin/bash", "-ec"] + command: + - | + set -- /ddl/*.sql + if [[ ! -e "$$1" ]]; then + printf 'ОШИБКА: в /ddl нет файлов SQL.\n' >&2 + exit 1 + fi + for ddl_file do + printf 'Применяем %s.\n' "$${ddl_file##*/}" + clickhouse-client \ + --host clickhouse-01 \ + --multiquery \ + --queries-file "$$ddl_file" + done + volumes: + - ./sql/ddl:/ddl:ro + postgres-metadata: image: ${POSTGRES_IMAGE:-postgres:16-alpine} restart: unless-stopped @@ -187,6 +262,10 @@ services: condition: service_healthy clickhouse-01: condition: service_healthy + # Не убирать: без зависимого успешный clickhouse-init считается + # упавшим для --wait. + clickhouse-init: + condition: service_completed_successfully entrypoint: ["/bin/bash"] command: ["/opt/airflow/init.sh"] diff --git a/docs/adr/0005-event-ingestion.md b/docs/adr/0005-event-ingestion.md index a45ee71..00ada90 100644 --- a/docs/adr/0005-event-ingestion.md +++ b/docs/adr/0005-event-ingestion.md @@ -4,14 +4,18 @@ ## Решение -Топик `hits` читает одна Kafka-таблица формата `RawBLOB`: сообщение ложится в -`stg.hits_raw_dist` строкой, как пришло, рядом с метаданными доставки. Ни -типизации, ни проверки на этом шаге нет — слой сырья ничего не интерпретирует. +Топик `hits` читает одна Kafka-таблица формата `RawBLOB`: непустое сообщение +ложится в `stg.hits_raw_dist` строкой, как пришло, рядом с метаданными доставки. +Ни типизации, ни проверки на этом шаге нет — слой сырья ничего не интерпретирует. У формата есть следствие для DDL: он читает вход в одно значение и рассчитан на таблицу с единственной колонкой, поэтому у чтеца она ровно одна — `raw`, а метаданные доставки берутся только из виртуальных колонок и добавить к чтецу что-либо своё нельзя. +Слово «непустое» приписано позже самого решения: у обещания нашлась измеренная +граница, и она датирована ниже, в разделе «Что проверено». Решения она не +меняет — от того, рождает ли пустая запись строку, выбор формата не зависит. + Типизированный слой наполняют две матвью, привязанные к `stg.hits_raw_dist`. Поля достаются `JSONExtract`. В таблицу ошибок уходят три класса брака: сообщение, не являющееся объектом JSON; объект, чей набор ключей разошёлся с @@ -99,8 +103,9 @@ STG сырьё исключительно от брака: слой сырых грязные записи не должны валить пайплайн. Лечится это двумя способами — включить на чтеце режим обработки ошибок или не проверять на входе вовсе. Второе честнее: частичная проверка в слое, чья работа — не проверять, спорит сама с собой, а -настоящий разбор всё равно идёт ниже. При `RawBLOB` ломаться нечему, доезжают -любые байты, и оба класса брака разбираются в одном месте. +настоящий разбор всё равно идёт ниже. При `RawBLOB` ломаться нечему, доезжает +любое непустое сообщение, каким бы мусором оно ни было (про «непустое» — там же, +в «Что проверено»), и оба класса брака разбираются в одном месте. Архив нужен буквальный и читаемый глазами. `RawBLOB` — это про способ чтения, а не про тип колонки: на диске лежит обычный `String` с текстом события, и менти @@ -162,18 +167,20 @@ contract-тест из #43 сюда не дотягивается — он ср Матвью с источником-`Distributed` срабатывает на вставку именно в эту распределённую таблицу — блок она видит до разрезания по шардам. Проверено -владельцем на рабочих проектах; на стенде подтверждается заодно с приёмкой #37. +владельцем на рабочих проектах; на стенде проверяется вместе с матвью разбора, +то есть при исполнении #43. -На живом стенде проверяется при исполнении #37. Первые два внесены в раздел 11 -спеки; остальные — однострочные `SELECT`, их довольно прогнать заодно: +Пункт про `RawBLOB`, стоявший в списке ниже первым, закрыт при исполнении #37, +на стенде 5 августа 2026 года: формат даёт ровно одну строку на каждое непустое +сообщение, пачка продюсера границы сообщений не стирает, запасной `LineAsString` +не понадобился. +Тогда же нашлась и граница обещания — запись с пустым значением и +запись-надгробие не дают строки вовсе; из-за неё в «Решении» и в «Почему» +приписано слово «непустое». Замер целиком — в [доке +хранилища](../architecture/storage.md), раздел «Что проверено». Остальные три +по-прежнему ждут живого стенда; это однострочные `SELECT`, их довольно прогнать +заодно: -- `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение. Проверять это - нужно первым и до написания DDL: формулировка «читает вход в одно значение» - про файл понятна, а про пачку сообщений из топика — нет, и если сообщения - склеятся, переделывать придётся решение целиком, а не DDL. Опыт стоит трёх - сообщений и одного `count()`. Запасной вариант — `LineAsString`: он режет по - переводу строки, а события у нас однострочные; цена запасного — сообщение с - переводом строки внутри даст две строки вместо одной; - форма именованного кортежа в `JSONExtract` с `Nullable`-членами — нужна для свёртки сорока семи вызовов в один, если разбор окажется дорогим; - `isValidJSON('123')` возвращает единицу, а `JSONExtractKeys` от скаляра — diff --git a/docs/architecture/storage.md b/docs/architecture/storage.md index 436943d..929ebd7 100644 --- a/docs/architecture/storage.md +++ b/docs/architecture/storage.md @@ -6,10 +6,10 @@ раздел «Что проверено» — чему в этом тексте верить и на каком основании. **Что здесь описано и чего ещё нет.** Собран этап 1: кластер из двух шардов, -keeper, Kafka, каркас сервисов. Объекты хранилища и механизм применения DDL -закладывает этап 2 — на момент написания их в репозитории нет. Дальше по тексту -устройство описано так, как оно проектируется; построенное от заложенного -отличает карта таблиц в конце. +keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/` уже лежат базы слоёв +и объекты STG — чтец топика `hits`, таблицы сырья и матвью приёма. Объектов ODS +в репозитории пока нет. Дальше по тексту устройство описано так, как оно +проектируется; построенное от заложенного отличает карта таблиц в конце. Зона ответственности у документа одна — хранилище. Генератор описан отдельно: его замысел — в [спеке генератора](../specs/2026-08-01-generator.md), формат @@ -100,16 +100,19 @@ keeper, Kafka, каркас сервисов. Объекты хранилища Типы у них такие: `kafka_topic` и `consumer_host` — `LowCardinality(String)`, значений мало и они повторяются; `kafka_partition` и `kafka_offset` — `UInt64`. -С `kafka_timestamp` сложнее: виртуальная колонка `_timestamp` заполнена не -всегда, а разрядность у секундной и миллисекундной версий разная. Поэтому колонка -объявляется `Nullable(DateTime)`, а точная форма проверяется на стенде -(раздел 11 спеки) — записать её в необнуляемый тип значит либо уронить приём на -первом сообщении, либо получить тихие нули за 1970 год. +С `kafka_timestamp` сложнее, и форма его решена на стенде. Меток времени движок +даёт две: `_timestamp` — `Nullable(DateTime)`, то есть секунды, и +`_timestamp_ms` — `Nullable(DateTime64(3))`, миллисекунды. Колонка объявлена +`Nullable(DateTime64(3))` и заполняется из `_timestamp_ms`: у брокера метка +миллисекундная, соседняя `_load_ts` тоже `DateTime64(3)`, а слой сырья хранит +приехавшее, и округлять ему нечего. Обнуляемость нужна отдельно от разрядности: +брокер метку заполняет не всегда, а необнуляемый тип значил бы либо падение +приёма на первом сообщении, либо тихие нули за 1970 год. Заполняются все они выражением в `SELECT` матвью приёма, а не `DEFAULT` в -таблице. Для `consumer_host` это обязательно: `DEFAULT hostName()` -сработал бы на шарде-получателе и назвал бы не ту ноду, которая читала топик, — -то есть колонка молча отвечала бы на другой вопрос. +таблице. Для `consumer_host` это обязательно: `DEFAULT hostName()` вычисляется +на шарде-получателе, то есть назвал бы не ту ноду, которая читала топик, — +колонка молча отвечала бы на другой вопрос. Само сообщение лежит в колонке `raw` тем, чем пришло: чтец читает байты и ничего не проверяет, поэтому там оказываются и целые события, и мусор. Разбирается всё @@ -189,6 +192,10 @@ Airflow она стоит — там это обычный `SETTINGS` у зап ноутбуке, где мир пересобирается одной командой. Размен не в пользу настройки, а компромисс полезнее показать, чем спрятать за галочкой. +Ещё одно место, где сырьё хранит не всё приехавшее: запись с пустым значением и +запись-надгробие проходят молча, не оставляя строки. Свойство измерено и принято +осознанно — подробности в разделе «Что проверено». + **Гарантии нет ни в одну сторону — есть два узких окна.** Окно потери описано выше: нода умерла между коммитом офсетов и сбросом спула. Окно дубля противоположное: нода умерла после записи на шард, но до коммита офсетов, и при @@ -311,18 +318,23 @@ ODS. Второе: матвью приёма создаётся последне работающих `airflow-init` и `superset-init`. Сначала `kafka-init` создаёт топик `hits` с двумя партициями, затем `clickhouse-init` дожидается его завершения и применяет файлы с ноды 1, `ON CLUSTER`. Этот порядок страхует от автосоздания -топика брокером с одной партицией: у потребителя librdkafka разрешение на -автосоздание по умолчанию выключено, так что случиться это не обязано, но урок -«обе ноды читают топик» умирает тихо, и полагаться на умолчание клиента здесь -не стоит. +топика с одной партицией. Переключателей тут два, и путать их не надо: брокер +автосоздание разрешает, а потребитель librdkafka внутри ClickHouse его не +просит — оба конца измерены, см. «Что проверено». То есть стенд держится на +умолчании клиента, а урок «обе ноды читают топик» умирает тихо, поэтому топик и +создаётся явно, до применения DDL. Образцы копируются не целиком, и в двух местах. `clickhouse-init` обязан ждать готовности **обеих** нод: `ON CLUSTER` ждёт исполнения на всех хостах и по таймауту бросает, а `airflow-init` ждёт только первую ноду, `superset-init` — только вторую. И второе: оба образца переживают `make up --wait` лишь потому, что -от них зависят долгоживущие сервисы; у пары `kafka-init` / `clickhouse-init` -таких зависимых нет, и как поведёт себя `--wait` с одноразовым сервисом без них — -проверяется при исполнении #37. +от них зависят долгоживущие сервисы. Что делает `--wait` с одноразовым сервисом +без зависимых, проверено при исполнении #37 на Docker Compose 2.40.3: считает +его упавшим и возвращает единицу, хотя контейнер вышел с нулём. Поэтому +зависимые есть и у новой пары: `clickhouse-init` ждёт `kafka-init`, а +`airflow-init` — `clickhouse-init`, и цепочка упирается в долгоживущий Airflow. +Побочная выгода важнее обхода `--wait`: к моменту старта Airflow DDL заведомо +применён. Повторный `make up` поверх живого тома проходит зелёным: весь DDL идёт через `CREATE ... IF NOT EXISTS`. Оборотная сторона — изменённый объект тем же @@ -333,7 +345,8 @@ ODS. Второе: матвью приёма создаётся последне ## Карта таблиц -Ниже — то, что закладывает этап 2. В репозитории этих объектов пока нет. +Ниже — то, что закладывает этап 2. DDL слоя STG уже лежит в `sql/ddl/`; +объектов ODS в репозитории пока нет. | Слой | Объект | Что это | |---|---|---| @@ -349,12 +362,12 @@ ODS. Второе: матвью приёма создаётся последне ## Что проверено -Документ описывает устройство, которого в репозитории ещё нет, и на каждом шагу -опирается на поведение ClickHouse. Поэтому утверждения о движке разведены на три -группы: насколько фразе можно верить, должно быть видно из текста, а не зависеть -от того, хорошо ли автор помнит документацию. Сверка — через MCP Context7, -5 августа 2026 года; то же разведение для механики приёма — в -[ADR 0005](../adr/0005-event-ingestion.md). +Документ на каждом шагу опирается на поведение ClickHouse, а местами и Kafka. +Поэтому утверждения о них разведены на три группы: насколько фразе можно верить, +должно быть видно из текста, а не зависеть от того, хорошо ли автор помнит +документацию. Сверка с +документацией — через MCP Context7, 5 августа 2026 года; то же разведение для +механики приёма — в [ADR 0005](../adr/0005-event-ingestion.md). **Сверено с документацией.** Собственная колонка с именем виртуальной делает виртуальную недоступной. При вставке в `Distributed` шард выбирается по ключу @@ -369,19 +382,63 @@ ODS. Второе: матвью приёма создаётся последне завязан на движок базы `Atomic`. `ON CLUSTER` ждёт все хосты и бросает по таймауту; `CREATE ... IF NOT EXISTS` на существующем объекте не бросает. -**Записано как проверка на стенде** — раздел 11 спеки и ADR 0005. Срабатывание -матвью с источником-`Distributed` на вставку именно в неё: этого случая в -документации нет вовсе, утверждение держится на опыте владельца. Одна строка на -сообщение у `RawBLOB`. Обнуляемость и разрядность виртуальной колонки -`_timestamp`. +**Проверено на стенде.** Опыты прогнаны на живом кластере при исполнении #37: +четыре — 5 августа 2026 года, пятый — 6 августа. Все подтвердили то, что здесь +написано. -**Сказано по памяти, проверки пока нет.** Что `DROP/REPLACE PARTITION` не -работает по `Distributed` — прямого запрета в документации нет, все примеры даны -для семейства MergeTree. Что `DEFAULT hostName()` вычислился бы на -шарде-получателе, а не на вставляющей ноде, и что имя читавшей ноды после записи -в `Distributed` уже невосстановимо. Что у потребителя librdkafka автосоздание -топиков по умолчанию выключено. Что упавшая матвью роняет вставку и -останавливает потребление до починки — на этой фразе держится правило «грязные -записи не валят пайплайн», и стоит она пока на одном рассуждении. Проверяются -все пятеро дёшево и заодно с приёмкой #37; до тех пор это предположения, а не -знание. +- `RawBLOB` даёт ровно одну строку на каждое непустое сообщение. Три сообщения + с ключами, поставленные в очередь до одного сброса продюсера, стали тремя + строками с тремя разными офсетами: пачка продюсера границы сообщений не + стирает. Запасной формат `LineAsString` не понадобился. Про границу «непустое» + — сразу ниже. +- Меток времени у Kafka-движка две, и разрядность у них разная: `_timestamp` — + `Nullable(DateTime)`, `_timestamp_ms` — `Nullable(DateTime64(3))`. Отсюда + форма `kafka_timestamp` в разделе о служебных колонках. +- `DEFAULT hostName()`, объявленный только на локальной таблице, вычисляется на + шарде-получателе. Строки, вставленные с ноды 1 и уехавшие на ноду 2, несут в + этой колонке ноду 2, а в соседней, заполненной явным `hostName()` во + вставляющем `SELECT`, — ноду 1. Довод за то, чтобы служебные колонки + заполняла матвью выражением, держится. Отсюда же и невосстановимость: после + записи в `Distributed` имя читавшей ноды взять больше неоткуда — своей + колонкой оно не сохранено, а умолчание назовёт получателя. +- `DROP PARTITION` и `REPLACE PARTITION` по `Distributed` не работают: обе + операции отвечают кодом 48, «Table engine Distributed doesn't support + partitioning», и оба шарда остаются нетронутыми. +- Автосоздания топиков ClickHouse не просит, и переключателей здесь два. На + брокере автосоздание разрешено: `auto.create.topics.enable=true`, и это + умолчание образа, а не наша настройка. На клиенте — выключено: чтец с матвью, + наведённые на несуществующий топик, ждали его с ошибкой «Broker: Unknown topic + or partition», и топик не появился. Значит, от тихого топика с одной партицией + стенд бережёт клиентское умолчание, а не брокер. + +Оговорка к первому опыту, и она измеренная: граница проходит по пустоте. Запись +Kafka с пустым значением (ноль байт) и запись-надгробие (значение `null`) +читаются, двигают офсет потребителя и не дают строки вовсе — ни в сырьё, ни в +таблицу ошибок; приём при этом не останавливается. Проверено 5 августа +2026 года: в топик ушли четыре записи — пустая, обычная, надгробие, обычная, — +в сырьё приехали две, а офсет группы сдвинулся на все четыре. Формат тут ни при +чём: `LineAsString`, запасной по [ADR 0005](../adr/0005-event-ingestion.md), на +том же наборе даёт ровно те же две строки и тот же офсет. + +Свойство принято осознанно и переделкой не закрывается. Генератор стенда пустых +сообщений не шлёт, приём от них не встаёт, а менять решённый формат из-за +случая, которого стенд не производит, — размен не в ту сторону. Знать о нём +стоит ровно затем, чтобы не искать пропавшую строку глазами: пропуск виден в +самом сырье, `kafka_offset` лежит там колонкой, и дырка в офсетах — это он. + +Оговорка к третьему опыту, и она сама непроверенная: `CREATE ... Distributed AS +<локальная>` копирует умолчания колонок, а значит, умолчание, оказавшееся заодно +и на распределённой таблице, могло бы вычислиться до раскладки по шардам. То +есть вывод опыта надёжен именно для умолчания, живущего только на локальной +таблице. Это вычитано, а не измерено. На устройство приёма оговорка не влияет: +служебные колонки мы заполняем выражением при любом ответе. + +**Сказано по памяти, проверки пока нет.** Осталось два утверждения, и оба ждут +одного и того же — матвью разбора, а она приходит с #43. + +- Матвью с источником-`Distributed` срабатывает на вставку именно в эту + распределённую таблицу, до раскладки по шардам. В документации случая нет + вовсе, утверждение держится на опыте владельца. +- Упавшая матвью роняет вставку и останавливает потребление до починки. На этой + фразе стоит правило «грязные записи не валят пайплайн», а сама она стоит пока + на одном рассуждении. diff --git a/docs/specs/2026-07-30-stand-v2-realism.md b/docs/specs/2026-07-30-stand-v2-realism.md index e5c3e5c..a1d7964 100644 --- a/docs/specs/2026-07-30-stand-v2-realism.md +++ b/docs/specs/2026-07-30-stand-v2-realism.md @@ -579,18 +579,19 @@ smoke-проверки, а не «дашборд зелёный». Это мин ## 11. Проверить при исполнении +Список убывает по мере постройки: проверенное уходит отсюда, а ответ с датой +остаётся там, где на него опираются. Формат чтеца и форма виртуальной метки +времени закрыты при исполнении #37 — см. [доку +хранилища](../architecture/storage.md), раздел «Что проверено». + - Поведение соединения двух Distributed-таблиц и `distributed_product_mode` — эмпирически на стенде (хвост #14). - Kafka Engine на двух нодах в одной consumer group: ребаланс партиций между - прогонами, отсутствие дублей при штатной работе. + прогонами, отсутствие дублей при штатной работе. Закрыто пока наполовину: что + обе ноды читают топик и обе партиции доезжают, показал #37; что дублей нет и + как раскладка меняется между прогонами — нет. - Точная форма `ORDER BY` ODS-таблиц (выражение `intHash32` в ключе ReplacingMergeTree). -- `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение (ADR 0005). - Проверять первым, до написания DDL: если сообщения склеятся, переделывать - придётся решение, а не запрос. Запасной формат — `LineAsString`. -- Тип виртуальной колонки `_timestamp` у Kafka-движка: обнуляемость и - разрядность (секунды против миллисекунд) — от этого зависит объявление - `kafka_timestamp` в таблице сырья. - Матвью с источником-`Distributed` срабатывает на вставку именно в эту распределённую таблицу, до раскладки по шардам: на этом стоит цепочка STG → ODS (ADR 0005). Проверено владельцем на рабочих проектах, в документации diff --git a/sql/ddl/00-databases.sql b/sql/ddl/00-databases.sql new file mode 100644 index 0000000..8a15c8d --- /dev/null +++ b/sql/ddl/00-databases.sql @@ -0,0 +1,15 @@ +-- Базы слоёв хранилища. +-- +-- Файлы этой папки применяются по порядку имён с ноды 1 и всегда ON CLUSTER: +-- объекты обязаны появиться на обеих нодах, иначе распределённая таблица +-- окажется лицом половины кластера. Имя кластера — clickstream_cluster, +-- задано в infra/clickhouse/config.d/cluster.xml. +-- +-- Идемпотентность везде через IF NOT EXISTS: make up применяет эти файлы и +-- поверх живого тома. + +CREATE DATABASE IF NOT EXISTS stg ON CLUSTER clickstream_cluster; + +-- База ODS заводится здесь же, хотя её объекты приносит #43: базы дёшевы, а +-- порядок файлов от этого не зависит. +CREATE DATABASE IF NOT EXISTS ods ON CLUSTER clickstream_cluster; diff --git a/sql/ddl/10-stg-tables.sql b/sql/ddl/10-stg-tables.sql new file mode 100644 index 0000000..73a7c10 --- /dev/null +++ b/sql/ddl/10-stg-tables.sql @@ -0,0 +1,87 @@ +-- STG: чтец топика hits и таблицы сырья. +-- +-- Слой сырья ничего не интерпретирует: сообщение ложится строкой, как пришло, +-- рядом с метаданными доставки. Довод целиком — ADR 0005, конвенции колонок и +-- сроков — docs/architecture/storage.md. + +-- Чтец топика. Колонка ровно одна: формат RawBLOB читает вход в одно значение +-- и рассчитан на таблицу с единственным полем String, метаданные доставки +-- берутся только из виртуальных колонок, своего к чтецу добавить нельзя. +-- Проверено на стенде 5 августа 2026 года: одно непустое сообщение Kafka даёт +-- ровно одну строку, сообщения в одной пачке продюсера не склеиваются. Граница +-- у обещания есть: запись с пустым значением и запись-надгробие читаются, +-- двигают офсет и строки не дают вовсе. Свойство принято осознанно — генератор +-- таких сообщений не шлёт; замер и довод — в доке хранилища, «Что проверено». +-- +-- Чтец стоит на обеих нодах и читает одной группой потребителей. Имя группы +-- одинаково на обеих по построению: DDL идёт ON CLUSTER и макросов в имени +-- нет. Разные группы дали бы каждой ноде полную копию топика. +-- +-- Читать топик движок начинает не сейчас, а в момент создания матвью +-- (40-stg-views.sql) — см. комментарий там. +CREATE TABLE IF NOT EXISTS stg.hits_raw_kafka ON CLUSTER clickstream_cluster +( + raw String +) +ENGINE = Kafka +SETTINGS + kafka_broker_list = 'kafka:9092', + kafka_topic_list = 'hits', + kafka_group_name = 'clickstream_hits', + kafka_format = 'RawBLOB'; + +-- Локальная таблица сырья. +-- +-- Движок именно ReplicatedMergeTree, не Replacing: повтор доставки в сырье +-- обязан быть виден — ради этого слой и заведён. +-- +-- Служебные колонки не повторяют имён виртуальных (_topic, _partition, +-- _offset, _timestamp у Kafka), иначе в матвью перестанет читаться, что дано +-- движком, а что положено нами. consumer_host — имя читавшей топик ноды: +-- виртуальные колонки его не несут, а после записи в Distributed он уже +-- невосстановим. +-- +-- kafka_timestamp — Nullable(DateTime64(3)), и заполняется из виртуальной +-- колонки _timestamp_ms, а не из _timestamp. Измерено на стенде 5 августа +-- 2026 года: _timestamp — Nullable(DateTime), то есть секунды; _timestamp_ms — +-- Nullable(DateTime64(3)). Взяты миллисекунды: у брокера метка миллисекундная, +-- _load_ts рядом тоже DateTime64(3), а слой сырья хранит то, что приехало, и +-- округлять ему нечего. Обнуляемость обязательна: метку брокер заполняет не +-- всегда, а необнуляемый тип дал бы либо падение приёма, либо тихий 1970 год. +-- +-- Нарезка и срок жизни — по _load_ts, то есть по реальному времени загрузки: +-- модельный день события живёт в ODS, а по нему TTL был бы просто сломан. +-- Срок — трое суток плюс хвост до суток: куски снимаются целиком +-- (ttl_only_drop_parts), а партицию закрывает календарный день. Значение +-- настройки проставлено явно, чтобы поведение не зависело от умолчания версии. +-- +-- Путь в keeper — с базой и без {uuid}: одноимённые таблицы разных слоёв иначе +-- подерутся за один узел, а читаемый путь на учебном стенде сам по себе +-- половина урока про keeper. +CREATE TABLE IF NOT EXISTS stg.hits_raw_rep ON CLUSTER clickstream_cluster +( + raw String, + kafka_topic LowCardinality(String), + kafka_partition UInt64, + kafka_offset UInt64, + kafka_timestamp Nullable(DateTime64(3)), + consumer_host LowCardinality(String), + _load_ts DateTime64(3) +) +ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/{database}/{table}', '{replica}') +PARTITION BY toDate(_load_ts) +ORDER BY (kafka_partition, kafka_offset) +TTL toDateTime(_load_ts) + INTERVAL 3 DAY +SETTINGS ttl_only_drop_parts = 1; + +-- Лицо слоя: пишем и читаем через него, локальная таблица остаётся для +-- обслуживания. Ключ шардирования — хеш сырой строки: разложить строку иначе +-- нечем, зато одинаковые сообщения ложатся на один шард. +-- +-- Операции с партициями по этой таблице не работают: проверено на стенде +-- 5 августа 2026 года, и DROP PARTITION, и REPLACE PARTITION отвечают +-- кодом 48 «Table engine Distributed doesn't support partitioning». Партиции +-- снимаются по локальным таблицам, ON CLUSTER. +CREATE TABLE IF NOT EXISTS stg.hits_raw_dist ON CLUSTER clickstream_cluster +AS stg.hits_raw_rep +ENGINE = Distributed('clickstream_cluster', 'stg', 'hits_raw_rep', cityHash64(raw)); diff --git a/sql/ddl/40-stg-views.sql b/sql/ddl/40-stg-views.sql new file mode 100644 index 0000000..bb8db44 --- /dev/null +++ b/sql/ddl/40-stg-views.sql @@ -0,0 +1,35 @@ +-- STG: матвью приёма — из чтеца в сырьё. +-- +-- Номер 40, а не 11, и пропуск в нумерации намеренный. Файлы 20-ods-tables.sql +-- и 30-ods-views.sql приносит #43, и матвью приёма обязана создаваться после +-- матвью разбора: Kafka-движок начинает читать топик ровно тогда, когда к нему +-- привязывают первую матвью. Создай её раньше разбора — и всё, что доедет в +-- зазоре, ляжет в сырьё и не попадёт в ODS никуда, ни в событие, ни в ошибки. +-- На пустом топике зазор безвреден, поэтому первый прогон о нём не скажет: +-- проснётся он, когда тома ClickHouse снесены, а данные Kafka целы, то есть на +-- обычной отладке. Нумерацию здесь не «приводить в порядок». +-- +-- Пишем в stg.hits_raw_dist, а не в локальную таблицу: раскладку по шардам +-- обязан определять ключ шардирования, а не то, какая нода случайно читала +-- топик. Вставка при этом фоновая — окно потери принято осознанно, довод +-- целиком в docs/architecture/storage.md, раздел «Приём». +-- +-- Служебные колонки заполняются выражением здесь, а не DEFAULT в таблице. +-- Для consumer_host это обязательно: проверено на стенде 5 августа 2026 года — +-- DEFAULT hostName() вычисляется на шарде-получателе и назвал бы не ту ноду, +-- которая читала топик. hostName() в SELECT снимается на вставляющей ноде, +-- то есть отвечает ровно на нужный вопрос. +-- +-- Порядок колонок в SELECT совпадает с порядком в целевой таблице. +CREATE MATERIALIZED VIEW IF NOT EXISTS stg.hits_raw_mv ON CLUSTER clickstream_cluster +TO stg.hits_raw_dist +AS +SELECT + raw, + _topic AS kafka_topic, + _partition AS kafka_partition, + _offset AS kafka_offset, + _timestamp_ms AS kafka_timestamp, + hostName() AS consumer_host, + now64(3) AS _load_ts +FROM stg.hits_raw_kafka;