DDL-бутстрап, топик hits и приём сырья в STG обеими нодами #37
Notifications
Due Date
No due date set.
Blocks
Depends on
#43 Типизированный ODS: ods.event, строгий приём и таблица ошибок
ddmitry/clickstream-data-platform
#36 Каркас генератора и контракт схемы события
ddmitry/clickstream-data-platform
Reference: ddmitry/clickstream-data-platform#37
Reference in New Issue
Block a user
Part of #4.
Цель
Дать стенду руки для DDL и довести сырьё до STG: механизм применения DDL
при
make up, топикhitsс двумя партициями и приём сырых событийобеими нодами. Типизированный ODS со строгим приёмом — следующий тикет
(#43): этот закладывает фундамент, на котором тот строится.
Развилки закрыты грилингом и холодным ревью до начала работы — см. «Решено
до начала».
Первым делом: четыре опыта
Дока хранилища в разделе «Что проверено» держит группу «сказано по памяти».
Четыре из пяти утверждений закрывает этот тикет — они дешёвые, а от двух
зависит DDL, который иначе придётся переписывать.
Черновой связки (чтец + матвью + таблица-цель) для этого хватит; «до написания
DDL» здесь означает «до боевых файлов в
sql/ddl/», а не «вообще без SQL».RawBLOBдаёт ровно одну строку на сообщение. Три сообщения в топик,SELECT count(). Документация описывает формат как «читает вход в однозначение» — про файл понятно, про пачку сообщений из топика нет. Склеит —
переделывать придётся решение целиком, запасной формат
LineAsString(режет по переводу строки, события у нас однострочные).
_timestamp— обнуляемость и разрядность. Отнеё зависит, чем объявлять
kafka_timestamp.DEFAULT hostName()вычисляется на шарде-получателе, а не навставляющей ноде. Ради этого служебные колонки и заполняются выражением в
SELECTматвью; если опыт покажет обратное, довод в доке надо переписать.DROP/REPLACE PARTITIONне работает поDistributed. Прямого запрета вдокументации нет, все примеры даны для семейства MergeTree.
Пятое утверждение группы — что упавшая матвью останавливает потребление —
относится к #43, там и проверяется.
Результаты записать в доку хранилища: каждое утверждение либо переезжает в
группу «сверено», либо получает запись, что опыт показал иначе.
Что войдёт
make up: два одноразовых сервиса по образцуairflow-initиsuperset-init. Сначалаkafka-initсоздаёт топик, затемclickhouse-initприменяетsql/ddl/*.sqlпо порядку имён с ноды 1,ON CLUSTER. Идемпотентность — черезCREATE ... IF NOT EXISTS.clickhouse-initобязанждать готовности обеих нод:
ON CLUSTERждёт все хосты и бросает потаймауту, а
airflow-initждёт только первую ноду,superset-init— тольковторую. И оба образца переживают
make up --waitлишь потому, что от нихзависят долгоживущие сервисы; у новой пары таких зависимых нет — проверить
поведение и решить.
hitsс 2 партициями — явно, не автосозданием:KAFKA_NUM_PARTITIONS: "1"в compose иначе молча даст одну партицию, и урок«обе ноды читают топик» умрёт.
stgиodsON CLUSTER. Имя кластера —clickstream_cluster(
infra/clickhouse/config.d/cluster.xml).stg.hits_raw_kafka— чтец на обеих нодах, группа потребителейclickstream_hits, форматRawBLOB. Колонка ровно одна,raw: форматрассчитан на таблицу с единственным полем
String, метаданные берутся толькоиз виртуальных колонок, добавить к чтецу своё нельзя.
stg.hits_raw_rep/stg.hits_raw_dist— сырьё. Колонки:raw;kafka_topicиconsumer_host—LowCardinality(String);kafka_partitionи
kafka_offset—UInt64;kafka_timestamp— по итогу опыта 2;_load_ts—DateTime64(3).Движок —
ReplicatedMergeTree, именно неReplacing: повтор доставки всырье обязан быть виден.
ORDER BY (kafka_partition, kafka_offset),PARTITION BY toDate(_load_ts), TTL трое суток,ttl_only_drop_parts = 1,шардирование
cityHash64(raw), путь в keeper/clickhouse/tables/{shard}/{database}/{table}.stg.hits_raw_mv— наполняет сырьё из чтеца, пишет в_dist, а не влокальную таблицу. Служебные колонки заполняются выражением в
SELECT, а неDEFAULTв таблице.00-databases.sql,10-stg-tables.sql,40-stg-views.sql. Пропуск в нумерации намеренный:20-ods-tables.sqlи30-ods-views.sqlприносит #43, а матвью приёма обязана создаваться послематвью разбора. Kafka-движок начинает читать топик в момент её создания, и
всё, что доедет до появления разбора, ляжет в сырьё и не попадёт в ODS
никуда — ни в событие, ни в ошибки.
Синхронной вставки на пути приёма нет:
distributed_foreground_insertнедостижим для фонового потока Kafka-движка и связывает шарды. Довод целиком —
в доке хранилища, раздел «Приём».
Новые подстановки
${...}вcompose.yamlзаписывать в.env.example:их сверяет
check_env_consistencyвscripts/stand-smoke.sh, иначе краситmake smoke.Решено до начала
снято прежнее ограничение тикета «ошибки разбора должны рождаться на шаге
Kafka-движка»: критерий #43 достигается средствами матвью.
файлов DDL, свойства приёма — docs/architecture/storage.md.
Критерии приёмки
хранилища: раздел «Что проверено» перестроен по факту, а не переписан
на глаз.
make clean && make upс нуля поднимает стенд и применяет DDL;повторный
make upповерх живого тома тоже зелёный.hitsсуществует с 2 партициями (kafka-topics.sh --describeиз контейнера Kafka).
hits, появляется вstg.hits_raw_distс метаданными доставки.stg.hits_raw_dist. Отправлятьс ключами: без ключа sticky-партиционер сложит всю пачку в одну
партицию, и критерий станет лотереей. Подобрать пару ключей, дающих
разные партиции, и приложить их к проверке.
как есть.
Ожидание доезда — опросом с таймаутом по образцу
wait_for_local_tableизtests/smoke-guards.sh, а неsleepнаугад: умолчаниеstream_flush_interval_ms— 7,5 с, и фиксированная пауза даёт шаткую проверку.Границы
ods.event, строгий приём и contract-тест — #43.orders) — этап 3.Сначала прочитать
раздел «Что проверено».
после этапа 1.
Проверка
make clean && make up && make up— второй запуск проверяет идемпотентностьhits(console producer изконтейнера Kafka), SELECT из
stg.hits_raw_distDDL событий ON CLUSTER и приём hits обеими нодамиto DDL-бутстрап, топик hits и приём сырья в STG обеими нодамиПостановка переписана после грилинга и двух холодных ревью (2026-08-03/04).
Что изменилось по существу:
— форма приёма решена: чтец читает топик формата RawBLOB, разбора на входе нет (ADR 0005). Прежнее ограничение тикета «ошибки разбора должны рождаться на шаге Kafka-движка» снято: критерий #43 достигается средствами матвью;
— имена объектов переехали на суффикс вида (ADR 0006), поэтому в критериях теперь stg.hits_raw_dist и consumer_host;
— путь в keeper решён: /clickhouse/tables/{shard}/{database}/{table}, без {uuid};
— извлечённый event_date из сырья убран: модельный день — свойство содержимого, он живёт в ODS ключом партиции;
— добавлены нарезка по дню загрузки, TTL трое суток и синхронная вставка в Distributed;
— два новых критерия: не-JSON не останавливает приём; RawBLOB даёт ровно одну строку на сообщение.
Конвенции и доводы — docs/architecture/storage.md.
Переписано после холодного ревью в три линзы и сверки утверждений о ClickHouse
с документацией. Против прежней постановки изменилось вот что.
Раскладка файлов DDL другая. Было чередование по слоям (
10-stg-tables,11-stg-views,20-ods-tables,21-ods-views). Оно включало чтение топикараньше, чем появлялся разбор: всё доехавшее в зазоре легло бы в сырьё и молча
миновало ODS. Стало — сначала все таблицы, матвью приёма последней,
40-stg-views.sql.Синхронная вставка снята.
distributed_foreground_insert = 1недостижимдля фонового потока Kafka-движка (только профилем пользователя в конфигурации
ноды) и вдобавок связывает шарды: упала вторая нода — приём встаёт целиком, а
не наполовину. Окно потери принято осознанно, довод целиком — в доке хранилища.
Появился шаг «первым делом» — до написания DDL проверить, что
RawBLOBдаёт ровно одну строку на сообщение. Документация этот случай не описывает, а
если сообщения склеятся, переделывать придётся решение, а не запрос.
Названы типы служебных колонок, движок и ключ сортировки сырья. Движок
именно
ReplicatedMergeTree:Replacingотменил бы свойство слоя, радикоторого он заведён.
Оговорка про образцы.
clickhouse-initобязан ждать обе ноды, тогда какairflow-initждёт первую, аsuperset-init— вторую; и поведение--waitуодноразового сервиса без зависимых остаётся открытым вопросом.
Вторая правка за день, по итогам холодного ревью тикетов как наряда на работу.
Вопрос ревьюеру ставился один: сможет ли исполнитель сделать задачу, имея
только тикет и названные в нём документы, ни разу не спросив автора.
Четыре опыта поручены явно. Дока хранилища обещает, что утверждения из
группы «сказано по памяти» проверяются заодно с этим тикетом, а тикет о них
молчал. Теперь они перечислены с доводом, зачем каждый: от двух зависит DDL.
Пятое утверждение группы — про упавшую матвью — переехало в #43, где ему место.
Критерий про обе ноды перестал быть лотереей. Отправка без ключей кладёт
всю пачку в одну партицию sticky-партиционером: в сырье была бы одна партиция и
один
consumer_host, и критерий зеленел или краснел по удаче. Теперь требуютсяключи, дающие разные партиции.
Ожидание доезда — опросом по образцу
wait_for_local_tableизtests/smoke-guards.sh, а неsleepнаугад: умолчаниеstream_flush_interval_ms— 7,5 секунды.«До написания DDL» уточнено: опыты требуют черновой связки чтец + матвью +
цель, речь про боевые файлы в
sql/ddl/, а не про полный отказ от SQL.Плюс напоминание про
.env.example: новые подстановки вcompose.yamlбеззаписи туда красят
make smoke.Работа сделана, коммит
daf1338в веткеfeat/37-ddl-bootstrap-stg-ingest.Все шесть критериев приёмки прогнаны на стенде с чистого тома 6 августа
2026 года; чекбоксы проставлены по факту прогона.
Четыре опыта — все подтвердили ожидание. Результаты переехали в
docs/architecture/storage.md, раздел «Что проверено», новая группа«Проверено на стенде». Из группы «сказано по памяти» осталось два
утверждения, оба ждут матвью разбора из #43.
Три решения сверх буквы тикета — приняты владельцем по ходу:
kafka_timestamp—Nullable(DateTime64(3))из виртуальной_timestamp_ms. Тикет оставлял форму на «итог опыта 2»; опыт показал,что миллисекунды доступны, и терять их незачем.
hostname:у обеих нод ClickHouse вcompose.yaml— сверх составатикета. Без него
hostName()отдаёт идентификатор контейнера, иколонка
consumer_hostтеряет учебный смысл.make up --waitс одноразовыми службами решено замером, а нерассуждением: на Compose 2.40.3 служба без зависимых считается упавшей
даже при выходе с нулём. Отсюда цепочка
kafka-init→clickhouse-init→airflow-init. Замер с датой иверсией — в доке хранилища.
Найденная граница приёма.
RawBLOBмолча теряет запись с пустымзначением и запись-надгробие: офсет потребителя двигается, строки нет ни в
сырье, ни в ошибках, приём не останавливается. Запасной
LineAsStringтеряет их точно так же — формат тут не рычаг. Принято как известное
свойство, отдельного тикета не заводим; практический признак — дыра в
kafka_offset. Замер и довод — в доке хранилища, оговорка продублированав комментарии
sql/ddl/10-stg-tables.sqlи в ADR 0005.Открытый хвост. Сценарий приёмки, которым прогонялись критерии, в
репозиторий не внесён: решение «делать ли его постоянным сторожем в
tests/» намеренно оставлено владельцу.