Merge pull request 'feat(stg): DDL-бутстрап, топик hits и приём сырья обеими нодами' (#53) from feat/37-ddl-bootstrap-stg-ingest into main
Reviewed-on: #53
This commit was merged in pull request #53.
This commit is contained in:
@@ -94,6 +94,7 @@ services:
|
|||||||
# Инициатор DDL и будущее подключение Airflow.
|
# Инициатор DDL и будущее подключение Airflow.
|
||||||
clickhouse-01:
|
clickhouse-01:
|
||||||
<<: *clickhouse-common
|
<<: *clickhouse-common
|
||||||
|
hostname: clickhouse-01
|
||||||
ports:
|
ports:
|
||||||
- "127.0.0.1:${CLICKHOUSE_01_HTTP_PORT:-28123}:8123"
|
- "127.0.0.1:${CLICKHOUSE_01_HTTP_PORT:-28123}:8123"
|
||||||
- "127.0.0.1:${CLICKHOUSE_01_TCP_PORT:-29000}:9000"
|
- "127.0.0.1:${CLICKHOUSE_01_TCP_PORT:-29000}:9000"
|
||||||
@@ -106,6 +107,7 @@ services:
|
|||||||
# Будущее подключение Superset.
|
# Будущее подключение Superset.
|
||||||
clickhouse-02:
|
clickhouse-02:
|
||||||
<<: *clickhouse-common
|
<<: *clickhouse-common
|
||||||
|
hostname: clickhouse-02
|
||||||
ports:
|
ports:
|
||||||
- "127.0.0.1:${CLICKHOUSE_02_HTTP_PORT:-28124}:8123"
|
- "127.0.0.1:${CLICKHOUSE_02_HTTP_PORT:-28124}:8123"
|
||||||
- "127.0.0.1:${CLICKHOUSE_02_TCP_PORT:-29001}:9000"
|
- "127.0.0.1:${CLICKHOUSE_02_TCP_PORT:-29001}:9000"
|
||||||
@@ -150,6 +152,79 @@ services:
|
|||||||
retries: 20
|
retries: 20
|
||||||
start_period: 30s
|
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:
|
postgres-metadata:
|
||||||
image: ${POSTGRES_IMAGE:-postgres:16-alpine}
|
image: ${POSTGRES_IMAGE:-postgres:16-alpine}
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
@@ -187,6 +262,10 @@ services:
|
|||||||
condition: service_healthy
|
condition: service_healthy
|
||||||
clickhouse-01:
|
clickhouse-01:
|
||||||
condition: service_healthy
|
condition: service_healthy
|
||||||
|
# Не убирать: без зависимого успешный clickhouse-init считается
|
||||||
|
# упавшим для --wait.
|
||||||
|
clickhouse-init:
|
||||||
|
condition: service_completed_successfully
|
||||||
entrypoint: ["/bin/bash"]
|
entrypoint: ["/bin/bash"]
|
||||||
command: ["/opt/airflow/init.sh"]
|
command: ["/opt/airflow/init.sh"]
|
||||||
|
|
||||||
|
|||||||
@@ -4,14 +4,18 @@
|
|||||||
|
|
||||||
## Решение
|
## Решение
|
||||||
|
|
||||||
Топик `hits` читает одна Kafka-таблица формата `RawBLOB`: сообщение ложится в
|
Топик `hits` читает одна Kafka-таблица формата `RawBLOB`: непустое сообщение
|
||||||
`stg.hits_raw_dist` строкой, как пришло, рядом с метаданными доставки. Ни
|
ложится в `stg.hits_raw_dist` строкой, как пришло, рядом с метаданными доставки.
|
||||||
типизации, ни проверки на этом шаге нет — слой сырья ничего не интерпретирует.
|
Ни типизации, ни проверки на этом шаге нет — слой сырья ничего не интерпретирует.
|
||||||
У формата есть следствие для DDL: он читает вход в одно значение и рассчитан на
|
У формата есть следствие для DDL: он читает вход в одно значение и рассчитан на
|
||||||
таблицу с единственной колонкой, поэтому у чтеца она ровно одна — `raw`, а
|
таблицу с единственной колонкой, поэтому у чтеца она ровно одна — `raw`, а
|
||||||
метаданные доставки берутся только из виртуальных колонок и добавить к чтецу
|
метаданные доставки берутся только из виртуальных колонок и добавить к чтецу
|
||||||
что-либо своё нельзя.
|
что-либо своё нельзя.
|
||||||
|
|
||||||
|
Слово «непустое» приписано позже самого решения: у обещания нашлась измеренная
|
||||||
|
граница, и она датирована ниже, в разделе «Что проверено». Решения она не
|
||||||
|
меняет — от того, рождает ли пустая запись строку, выбор формата не зависит.
|
||||||
|
|
||||||
Типизированный слой наполняют две матвью, привязанные к `stg.hits_raw_dist`.
|
Типизированный слой наполняют две матвью, привязанные к `stg.hits_raw_dist`.
|
||||||
Поля достаются `JSONExtract`. В таблицу ошибок уходят три класса брака:
|
Поля достаются `JSONExtract`. В таблицу ошибок уходят три класса брака:
|
||||||
сообщение, не являющееся объектом JSON; объект, чей набор ключей разошёлся с
|
сообщение, не являющееся объектом JSON; объект, чей набор ключей разошёлся с
|
||||||
@@ -99,8 +103,9 @@ STG сырьё исключительно от брака: слой сырых
|
|||||||
грязные записи не должны валить пайплайн. Лечится это двумя способами — включить
|
грязные записи не должны валить пайплайн. Лечится это двумя способами — включить
|
||||||
на чтеце режим обработки ошибок или не проверять на входе вовсе. Второе честнее:
|
на чтеце режим обработки ошибок или не проверять на входе вовсе. Второе честнее:
|
||||||
частичная проверка в слое, чья работа — не проверять, спорит сама с собой, а
|
частичная проверка в слое, чья работа — не проверять, спорит сама с собой, а
|
||||||
настоящий разбор всё равно идёт ниже. При `RawBLOB` ломаться нечему, доезжают
|
настоящий разбор всё равно идёт ниже. При `RawBLOB` ломаться нечему, доезжает
|
||||||
любые байты, и оба класса брака разбираются в одном месте.
|
любое непустое сообщение, каким бы мусором оно ни было (про «непустое» — там же,
|
||||||
|
в «Что проверено»), и оба класса брака разбираются в одном месте.
|
||||||
|
|
||||||
Архив нужен буквальный и читаемый глазами. `RawBLOB` — это про способ чтения, а
|
Архив нужен буквальный и читаемый глазами. `RawBLOB` — это про способ чтения, а
|
||||||
не про тип колонки: на диске лежит обычный `String` с текстом события, и менти
|
не про тип колонки: на диске лежит обычный `String` с текстом события, и менти
|
||||||
@@ -162,18 +167,20 @@ contract-тест из #43 сюда не дотягивается — он ср
|
|||||||
|
|
||||||
Матвью с источником-`Distributed` срабатывает на вставку именно в эту
|
Матвью с источником-`Distributed` срабатывает на вставку именно в эту
|
||||||
распределённую таблицу — блок она видит до разрезания по шардам. Проверено
|
распределённую таблицу — блок она видит до разрезания по шардам. Проверено
|
||||||
владельцем на рабочих проектах; на стенде подтверждается заодно с приёмкой #37.
|
владельцем на рабочих проектах; на стенде проверяется вместе с матвью разбора,
|
||||||
|
то есть при исполнении #43.
|
||||||
|
|
||||||
На живом стенде проверяется при исполнении #37. Первые два внесены в раздел 11
|
Пункт про `RawBLOB`, стоявший в списке ниже первым, закрыт при исполнении #37,
|
||||||
спеки; остальные — однострочные `SELECT`, их довольно прогнать заодно:
|
на стенде 5 августа 2026 года: формат даёт ровно одну строку на каждое непустое
|
||||||
|
сообщение, пачка продюсера границы сообщений не стирает, запасной `LineAsString`
|
||||||
|
не понадобился.
|
||||||
|
Тогда же нашлась и граница обещания — запись с пустым значением и
|
||||||
|
запись-надгробие не дают строки вовсе; из-за неё в «Решении» и в «Почему»
|
||||||
|
приписано слово «непустое». Замер целиком — в [доке
|
||||||
|
хранилища](../architecture/storage.md), раздел «Что проверено». Остальные три
|
||||||
|
по-прежнему ждут живого стенда; это однострочные `SELECT`, их довольно прогнать
|
||||||
|
заодно:
|
||||||
|
|
||||||
- `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение. Проверять это
|
|
||||||
нужно первым и до написания DDL: формулировка «читает вход в одно значение»
|
|
||||||
про файл понятна, а про пачку сообщений из топика — нет, и если сообщения
|
|
||||||
склеятся, переделывать придётся решение целиком, а не DDL. Опыт стоит трёх
|
|
||||||
сообщений и одного `count()`. Запасной вариант — `LineAsString`: он режет по
|
|
||||||
переводу строки, а события у нас однострочные; цена запасного — сообщение с
|
|
||||||
переводом строки внутри даст две строки вместо одной;
|
|
||||||
- форма именованного кортежа в `JSONExtract` с `Nullable`-членами — нужна для
|
- форма именованного кортежа в `JSONExtract` с `Nullable`-членами — нужна для
|
||||||
свёртки сорока семи вызовов в один, если разбор окажется дорогим;
|
свёртки сорока семи вызовов в один, если разбор окажется дорогим;
|
||||||
- `isValidJSON('123')` возвращает единицу, а `JSONExtractKeys` от скаляра —
|
- `isValidJSON('123')` возвращает единицу, а `JSONExtractKeys` от скаляра —
|
||||||
|
|||||||
@@ -6,10 +6,10 @@
|
|||||||
раздел «Что проверено» — чему в этом тексте верить и на каком основании.
|
раздел «Что проверено» — чему в этом тексте верить и на каком основании.
|
||||||
|
|
||||||
**Что здесь описано и чего ещё нет.** Собран этап 1: кластер из двух шардов,
|
**Что здесь описано и чего ещё нет.** Собран этап 1: кластер из двух шардов,
|
||||||
keeper, Kafka, каркас сервисов. Объекты хранилища и механизм применения DDL
|
keeper, Kafka, каркас сервисов. Этап 2 идёт: в `sql/ddl/` уже лежат базы слоёв
|
||||||
закладывает этап 2 — на момент написания их в репозитории нет. Дальше по тексту
|
и объекты STG — чтец топика `hits`, таблицы сырья и матвью приёма. Объектов ODS
|
||||||
устройство описано так, как оно проектируется; построенное от заложенного
|
в репозитории пока нет. Дальше по тексту устройство описано так, как оно
|
||||||
отличает карта таблиц в конце.
|
проектируется; построенное от заложенного отличает карта таблиц в конце.
|
||||||
|
|
||||||
Зона ответственности у документа одна — хранилище. Генератор описан отдельно:
|
Зона ответственности у документа одна — хранилище. Генератор описан отдельно:
|
||||||
его замысел — в [спеке генератора](../specs/2026-08-01-generator.md), формат
|
его замысел — в [спеке генератора](../specs/2026-08-01-generator.md), формат
|
||||||
@@ -100,16 +100,19 @@ keeper, Kafka, каркас сервисов. Объекты хранилища
|
|||||||
|
|
||||||
Типы у них такие: `kafka_topic` и `consumer_host` — `LowCardinality(String)`,
|
Типы у них такие: `kafka_topic` и `consumer_host` — `LowCardinality(String)`,
|
||||||
значений мало и они повторяются; `kafka_partition` и `kafka_offset` — `UInt64`.
|
значений мало и они повторяются; `kafka_partition` и `kafka_offset` — `UInt64`.
|
||||||
С `kafka_timestamp` сложнее: виртуальная колонка `_timestamp` заполнена не
|
С `kafka_timestamp` сложнее, и форма его решена на стенде. Меток времени движок
|
||||||
всегда, а разрядность у секундной и миллисекундной версий разная. Поэтому колонка
|
даёт две: `_timestamp` — `Nullable(DateTime)`, то есть секунды, и
|
||||||
объявляется `Nullable(DateTime)`, а точная форма проверяется на стенде
|
`_timestamp_ms` — `Nullable(DateTime64(3))`, миллисекунды. Колонка объявлена
|
||||||
(раздел 11 спеки) — записать её в необнуляемый тип значит либо уронить приём на
|
`Nullable(DateTime64(3))` и заполняется из `_timestamp_ms`: у брокера метка
|
||||||
первом сообщении, либо получить тихие нули за 1970 год.
|
миллисекундная, соседняя `_load_ts` тоже `DateTime64(3)`, а слой сырья хранит
|
||||||
|
приехавшее, и округлять ему нечего. Обнуляемость нужна отдельно от разрядности:
|
||||||
|
брокер метку заполняет не всегда, а необнуляемый тип значил бы либо падение
|
||||||
|
приёма на первом сообщении, либо тихие нули за 1970 год.
|
||||||
|
|
||||||
Заполняются все они выражением в `SELECT` матвью приёма, а не `DEFAULT` в
|
Заполняются все они выражением в `SELECT` матвью приёма, а не `DEFAULT` в
|
||||||
таблице. Для `consumer_host` это обязательно: `DEFAULT hostName()`
|
таблице. Для `consumer_host` это обязательно: `DEFAULT hostName()` вычисляется
|
||||||
сработал бы на шарде-получателе и назвал бы не ту ноду, которая читала топик, —
|
на шарде-получателе, то есть назвал бы не ту ноду, которая читала топик, —
|
||||||
то есть колонка молча отвечала бы на другой вопрос.
|
колонка молча отвечала бы на другой вопрос.
|
||||||
|
|
||||||
Само сообщение лежит в колонке `raw` тем, чем пришло: чтец читает байты и ничего
|
Само сообщение лежит в колонке `raw` тем, чем пришло: чтец читает байты и ничего
|
||||||
не проверяет, поэтому там оказываются и целые события, и мусор. Разбирается всё
|
не проверяет, поэтому там оказываются и целые события, и мусор. Разбирается всё
|
||||||
@@ -189,6 +192,10 @@ Airflow она стоит — там это обычный `SETTINGS` у зап
|
|||||||
ноутбуке, где мир пересобирается одной командой. Размен не в пользу настройки, а
|
ноутбуке, где мир пересобирается одной командой. Размен не в пользу настройки, а
|
||||||
компромисс полезнее показать, чем спрятать за галочкой.
|
компромисс полезнее показать, чем спрятать за галочкой.
|
||||||
|
|
||||||
|
Ещё одно место, где сырьё хранит не всё приехавшее: запись с пустым значением и
|
||||||
|
запись-надгробие проходят молча, не оставляя строки. Свойство измерено и принято
|
||||||
|
осознанно — подробности в разделе «Что проверено».
|
||||||
|
|
||||||
**Гарантии нет ни в одну сторону — есть два узких окна.** Окно потери описано
|
**Гарантии нет ни в одну сторону — есть два узких окна.** Окно потери описано
|
||||||
выше: нода умерла между коммитом офсетов и сбросом спула. Окно дубля
|
выше: нода умерла между коммитом офсетов и сбросом спула. Окно дубля
|
||||||
противоположное: нода умерла после записи на шард, но до коммита офсетов, и при
|
противоположное: нода умерла после записи на шард, но до коммита офсетов, и при
|
||||||
@@ -311,18 +318,23 @@ ODS. Второе: матвью приёма создаётся последне
|
|||||||
работающих `airflow-init` и `superset-init`. Сначала `kafka-init` создаёт топик
|
работающих `airflow-init` и `superset-init`. Сначала `kafka-init` создаёт топик
|
||||||
`hits` с двумя партициями, затем `clickhouse-init` дожидается его завершения и
|
`hits` с двумя партициями, затем `clickhouse-init` дожидается его завершения и
|
||||||
применяет файлы с ноды 1, `ON CLUSTER`. Этот порядок страхует от автосоздания
|
применяет файлы с ноды 1, `ON CLUSTER`. Этот порядок страхует от автосоздания
|
||||||
топика брокером с одной партицией: у потребителя librdkafka разрешение на
|
топика с одной партицией. Переключателей тут два, и путать их не надо: брокер
|
||||||
автосоздание по умолчанию выключено, так что случиться это не обязано, но урок
|
автосоздание разрешает, а потребитель librdkafka внутри ClickHouse его не
|
||||||
«обе ноды читают топик» умирает тихо, и полагаться на умолчание клиента здесь
|
просит — оба конца измерены, см. «Что проверено». То есть стенд держится на
|
||||||
не стоит.
|
умолчании клиента, а урок «обе ноды читают топик» умирает тихо, поэтому топик и
|
||||||
|
создаётся явно, до применения DDL.
|
||||||
|
|
||||||
Образцы копируются не целиком, и в двух местах. `clickhouse-init` обязан ждать
|
Образцы копируются не целиком, и в двух местах. `clickhouse-init` обязан ждать
|
||||||
готовности **обеих** нод: `ON CLUSTER` ждёт исполнения на всех хостах и по
|
готовности **обеих** нод: `ON CLUSTER` ждёт исполнения на всех хостах и по
|
||||||
таймауту бросает, а `airflow-init` ждёт только первую ноду, `superset-init` —
|
таймауту бросает, а `airflow-init` ждёт только первую ноду, `superset-init` —
|
||||||
только вторую. И второе: оба образца переживают `make up --wait` лишь потому, что
|
только вторую. И второе: оба образца переживают `make up --wait` лишь потому, что
|
||||||
от них зависят долгоживущие сервисы; у пары `kafka-init` / `clickhouse-init`
|
от них зависят долгоживущие сервисы. Что делает `--wait` с одноразовым сервисом
|
||||||
таких зависимых нет, и как поведёт себя `--wait` с одноразовым сервисом без них —
|
без зависимых, проверено при исполнении #37 на Docker Compose 2.40.3: считает
|
||||||
проверяется при исполнении #37.
|
его упавшим и возвращает единицу, хотя контейнер вышел с нулём. Поэтому
|
||||||
|
зависимые есть и у новой пары: `clickhouse-init` ждёт `kafka-init`, а
|
||||||
|
`airflow-init` — `clickhouse-init`, и цепочка упирается в долгоживущий Airflow.
|
||||||
|
Побочная выгода важнее обхода `--wait`: к моменту старта Airflow DDL заведомо
|
||||||
|
применён.
|
||||||
|
|
||||||
Повторный `make up` поверх живого тома проходит зелёным: весь DDL идёт через
|
Повторный `make up` поверх живого тома проходит зелёным: весь DDL идёт через
|
||||||
`CREATE ... IF NOT EXISTS`. Оборотная сторона — изменённый объект тем же
|
`CREATE ... IF NOT EXISTS`. Оборотная сторона — изменённый объект тем же
|
||||||
@@ -333,7 +345,8 @@ ODS. Второе: матвью приёма создаётся последне
|
|||||||
|
|
||||||
## Карта таблиц
|
## Карта таблиц
|
||||||
|
|
||||||
Ниже — то, что закладывает этап 2. В репозитории этих объектов пока нет.
|
Ниже — то, что закладывает этап 2. DDL слоя STG уже лежит в `sql/ddl/`;
|
||||||
|
объектов ODS в репозитории пока нет.
|
||||||
|
|
||||||
| Слой | Объект | Что это |
|
| Слой | Объект | Что это |
|
||||||
|---|---|---|
|
|---|---|---|
|
||||||
@@ -349,12 +362,12 @@ ODS. Второе: матвью приёма создаётся последне
|
|||||||
|
|
||||||
## Что проверено
|
## Что проверено
|
||||||
|
|
||||||
Документ описывает устройство, которого в репозитории ещё нет, и на каждом шагу
|
Документ на каждом шагу опирается на поведение ClickHouse, а местами и Kafka.
|
||||||
опирается на поведение ClickHouse. Поэтому утверждения о движке разведены на три
|
Поэтому утверждения о них разведены на три группы: насколько фразе можно верить,
|
||||||
группы: насколько фразе можно верить, должно быть видно из текста, а не зависеть
|
должно быть видно из текста, а не зависеть от того, хорошо ли автор помнит
|
||||||
от того, хорошо ли автор помнит документацию. Сверка — через MCP Context7,
|
документацию. Сверка с
|
||||||
5 августа 2026 года; то же разведение для механики приёма — в
|
документацией — через MCP Context7, 5 августа 2026 года; то же разведение для
|
||||||
[ADR 0005](../adr/0005-event-ingestion.md).
|
механики приёма — в [ADR 0005](../adr/0005-event-ingestion.md).
|
||||||
|
|
||||||
**Сверено с документацией.** Собственная колонка с именем виртуальной делает
|
**Сверено с документацией.** Собственная колонка с именем виртуальной делает
|
||||||
виртуальную недоступной. При вставке в `Distributed` шард выбирается по ключу
|
виртуальную недоступной. При вставке в `Distributed` шард выбирается по ключу
|
||||||
@@ -369,19 +382,63 @@ ODS. Второе: матвью приёма создаётся последне
|
|||||||
завязан на движок базы `Atomic`. `ON CLUSTER` ждёт все хосты и бросает по
|
завязан на движок базы `Atomic`. `ON CLUSTER` ждёт все хосты и бросает по
|
||||||
таймауту; `CREATE ... IF NOT EXISTS` на существующем объекте не бросает.
|
таймауту; `CREATE ... IF NOT EXISTS` на существующем объекте не бросает.
|
||||||
|
|
||||||
**Записано как проверка на стенде** — раздел 11 спеки и ADR 0005. Срабатывание
|
**Проверено на стенде.** Опыты прогнаны на живом кластере при исполнении #37:
|
||||||
матвью с источником-`Distributed` на вставку именно в неё: этого случая в
|
четыре — 5 августа 2026 года, пятый — 6 августа. Все подтвердили то, что здесь
|
||||||
документации нет вовсе, утверждение держится на опыте владельца. Одна строка на
|
написано.
|
||||||
сообщение у `RawBLOB`. Обнуляемость и разрядность виртуальной колонки
|
|
||||||
`_timestamp`.
|
|
||||||
|
|
||||||
**Сказано по памяти, проверки пока нет.** Что `DROP/REPLACE PARTITION` не
|
- `RawBLOB` даёт ровно одну строку на каждое непустое сообщение. Три сообщения
|
||||||
работает по `Distributed` — прямого запрета в документации нет, все примеры даны
|
с ключами, поставленные в очередь до одного сброса продюсера, стали тремя
|
||||||
для семейства MergeTree. Что `DEFAULT hostName()` вычислился бы на
|
строками с тремя разными офсетами: пачка продюсера границы сообщений не
|
||||||
шарде-получателе, а не на вставляющей ноде, и что имя читавшей ноды после записи
|
стирает. Запасной формат `LineAsString` не понадобился. Про границу «непустое»
|
||||||
в `Distributed` уже невосстановимо. Что у потребителя librdkafka автосоздание
|
— сразу ниже.
|
||||||
топиков по умолчанию выключено. Что упавшая матвью роняет вставку и
|
- Меток времени у Kafka-движка две, и разрядность у них разная: `_timestamp` —
|
||||||
останавливает потребление до починки — на этой фразе держится правило «грязные
|
`Nullable(DateTime)`, `_timestamp_ms` — `Nullable(DateTime64(3))`. Отсюда
|
||||||
записи не валят пайплайн», и стоит она пока на одном рассуждении. Проверяются
|
форма `kafka_timestamp` в разделе о служебных колонках.
|
||||||
все пятеро дёшево и заодно с приёмкой #37; до тех пор это предположения, а не
|
- `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` срабатывает на вставку именно в эту
|
||||||
|
распределённую таблицу, до раскладки по шардам. В документации случая нет
|
||||||
|
вовсе, утверждение держится на опыте владельца.
|
||||||
|
- Упавшая матвью роняет вставку и останавливает потребление до починки. На этой
|
||||||
|
фразе стоит правило «грязные записи не валят пайплайн», а сама она стоит пока
|
||||||
|
на одном рассуждении.
|
||||||
|
|||||||
@@ -579,18 +579,19 @@ smoke-проверки, а не «дашборд зелёный». Это мин
|
|||||||
|
|
||||||
## 11. Проверить при исполнении
|
## 11. Проверить при исполнении
|
||||||
|
|
||||||
|
Список убывает по мере постройки: проверенное уходит отсюда, а ответ с датой
|
||||||
|
остаётся там, где на него опираются. Формат чтеца и форма виртуальной метки
|
||||||
|
времени закрыты при исполнении #37 — см. [доку
|
||||||
|
хранилища](../architecture/storage.md), раздел «Что проверено».
|
||||||
|
|
||||||
- Поведение соединения двух Distributed-таблиц и `distributed_product_mode` —
|
- Поведение соединения двух Distributed-таблиц и `distributed_product_mode` —
|
||||||
эмпирически на стенде (хвост #14).
|
эмпирически на стенде (хвост #14).
|
||||||
- Kafka Engine на двух нодах в одной consumer group: ребаланс партиций между
|
- Kafka Engine на двух нодах в одной consumer group: ребаланс партиций между
|
||||||
прогонами, отсутствие дублей при штатной работе.
|
прогонами, отсутствие дублей при штатной работе. Закрыто пока наполовину: что
|
||||||
|
обе ноды читают топик и обе партиции доезжают, показал #37; что дублей нет и
|
||||||
|
как раскладка меняется между прогонами — нет.
|
||||||
- Точная форма `ORDER BY` ODS-таблиц (выражение `intHash32` в ключе
|
- Точная форма `ORDER BY` ODS-таблиц (выражение `intHash32` в ключе
|
||||||
ReplacingMergeTree).
|
ReplacingMergeTree).
|
||||||
- `RawBLOB` в Kafka-движке даёт ровно одну строку на сообщение (ADR 0005).
|
|
||||||
Проверять первым, до написания DDL: если сообщения склеятся, переделывать
|
|
||||||
придётся решение, а не запрос. Запасной формат — `LineAsString`.
|
|
||||||
- Тип виртуальной колонки `_timestamp` у Kafka-движка: обнуляемость и
|
|
||||||
разрядность (секунды против миллисекунд) — от этого зависит объявление
|
|
||||||
`kafka_timestamp` в таблице сырья.
|
|
||||||
- Матвью с источником-`Distributed` срабатывает на вставку именно в эту
|
- Матвью с источником-`Distributed` срабатывает на вставку именно в эту
|
||||||
распределённую таблицу, до раскладки по шардам: на этом стоит цепочка
|
распределённую таблицу, до раскладки по шардам: на этом стоит цепочка
|
||||||
STG → ODS (ADR 0005). Проверено владельцем на рабочих проектах, в документации
|
STG → ODS (ADR 0005). Проверено владельцем на рабочих проектах, в документации
|
||||||
|
|||||||
@@ -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;
|
||||||
@@ -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));
|
||||||
@@ -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;
|
||||||
Reference in New Issue
Block a user