Зачем: стенду нужен воспроизводимый холодный старт, при котором схема хранилища и топик появляются сами, а сырьё из 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 <noreply@anthropic.com>
88 lines
6.2 KiB
SQL
88 lines
6.2 KiB
SQL
-- 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));
|