Модуль 14. Стриминг и NRT-витрины на масштабе (Kafka→ClickHouse)
Это практическая «сборка под ключ» для потоковой загрузки событий (1–100k ev/s), минутных KPI и соблюдения SLA свежести. Я последовательно покажу архитектуру, DDL/конфиги, антихрупкую топологию Kafka→staging→MV→agg_minute_state→vw_*, водяные знаки (watermarks) и обработку опоздавших событий, идемпотентность/дедуп, буферизацию и backpressure, наблюдаемость и алерты, а затем — готовые runbook-и «как чинить».
Что считаем успехом (SLA и критерии приёмки)
- Throughput: стабильные 1–100k событий/сек на кластер (масштабируемо по партициям/нодам).
- Freshness SLA: минутные витрины обновляются ≤ N минут (обычно 1–5).
- Надёжность: нет «взрывов» частей (part-explosion), отставания мерджей и репликации.
- Тяжёлые метрики (uniq/quantile/top-K) считаются из состояний, а не «на лету».
- Lag Kafka и мердж-бэклог — в зелёной зоне; дубликаты погашаются идемпотентностью.
Архитектура топологии и потоки
Поставщики событий
└─→ Kafka (топик events N партиций)
└─→ ClickHouse ENGINE=Kafka (kafka_raw_json)
├─→ MV #1 → staging_raw (сырьё, дедуп, карантин)
└─→ MV #2 → fact_events_wide (нормализация полей, фильтры)
└─→ MV #3 → agg_events_minute_state (minute rollup: sum/uniq/quantile State)
└─→ vw_kpi_minute / vw_kpi_hour (семантика для BI)
Наблюдаемость:
system.query_log / parts / merges / replication_queue + kafka_exporter (lag) + алерты
Идея: разделяем этапы. Kafka-таблица + первая MV — это буферизация и карантин, вторая — нормализация (строгое зерно и идемпотентность), третья — агрегирование в состояния. BI/мини-API читают только vw_*.
Kafka Engine и первая ступень (сырьё, дедуп, карантин)
Таблица-источник (ENGINE=Kafka)
CREATE TABLE src_kafka_events ( payload String -- сырой JSON; парсим в MV, а не на самой таблице ) ENGINE = Kafka SETTINGS kafka_broker_list = 'kafka-1:9092,kafka-2:9092', kafka_topic_list = 'events', kafka_group_name = 'ch_ingest_events_v1', kafka_format = 'JSONEachRow', kafka_num_consumers = 6, -- под количество партиций/CPU kafka_max_block_size = 50000, -- размер микробатча kafka_handle_error_mode = 'stream'; -- не стопорить консьюмера из-за кривых сообщений
Замечания:
- Настройте достаточное количество партиций в Kafka (минимум = целевой параллелизм чтения).
- kafka_max_block_size управляет микробатчем в ClickHouse (и частотой коммитов оффсетов).
- Для «особо шумных» источников держите отдельный топик-DLQ (dead-letter queue).
Staging с идемпотентностью (Replacing + карантин)
Зерно staging: event_id (строго уникальный), event_time (UTC), версия (на случай переотправок/коррекций).
-- Сырые распарсенные события (сохраняем и сырой JSON) CREATE TABLE stg_events_raw ( event_id UUID, event_time DateTime64(3, 'UTC'), ingest_time DateTime64(3, 'UTC') DEFAULT now(), source LowCardinality(String), user_id UInt64, event_name LowCardinality(String), amount Decimal(18,2), attrs String, -- оригинальный JSON version UInt32 -- монотонно растёт для повторов ) ENGINE = ReplacingMergeTree(version) PARTITION BY toYYYYMM(event_time) ORDER BY (event_time, event_id); -- Карантин для «плохих» сообщений (парсинг/валидатор) CREATE TABLE stg_events_quarantine ( payload String, reason String, ingest_time DateTime64(3, 'UTC') DEFAULT now() ) ENGINE = MergeTree ORDER BY ingest_time;
MV #1: парсинг сырого JSON + карантин
CREATE MATERIALIZED VIEW mv_kafka_to_staging TO stg_events_raw AS SELECT JSON_VALUE(payload, '$.event_id')::UUID AS event_id, parseDateTime64BestEffortOrNull(JSON_VALUE(payload,'$.event_time')) AS event_time, JSON_VALUE(payload, '$.source') AS source, JSON_VALUE(payload, '$.user_id')::UInt64 AS user_id, JSON_VALUE(payload, '$.event_name') AS event_name, JSON_VALUE(payload, '$.amount')::Decimal(18,2) AS amount, payload AS attrs, coalesce(JSON_VALUE(payload, '$.version')::UInt32, 1) AS version FROM src_kafka_events WHERE isValidJSON(payload) -- грубая валидация AND JSONHas(payload, 'event_id') AND JSONHas(payload, 'event_time'); -- всё, что не прошло, складываем руками периодически: INSERT INTO stg_events_quarantine SELECT payload, 'parse_error' FROM src_kafka_events WHERE NOT isValidJSON(payload);
Почему Replacing(version): повторные доставки (retry/переигровка) «перекроются». Читать без FINAL — в fact перейдём через инкрементальные выборки.
Нормализация: факты событий и «правильное зерно»
Зерно fact: (event_id) + обязательные типизированные поля. Разрешаем только канонизированные события.
CREATE TABLE fact_events_wide ( event_id UUID, event_time DateTime64(3, 'UTC'), minute DateTime MATERIALIZED toStartOfMinute(event_time), user_id UInt64, event_name LowCardinality(String), amount Decimal(18,2), source LowCardinality(String) ) ENGINE = MergeTree PARTITION BY toYYYYMMDD(event_time) ORDER BY (event_time, user_id, event_id);
MV #2: из staging в факт (нормализация и фильтры)
CREATE MATERIALIZED VIEW mv_staging_to_fact
TO fact_events_wide AS
SELECT
event_id, event_time, user_id, event_name, amount, source
FROM stg_events_raw
-- простая фильтрация домена
WHERE event_time >= now() - INTERVAL 3 DAY
AND event_name IN ('view','add_to_cart','purchase');
На практике сюда добавляют обогащение из справочников (device/app/channel) через словари dictGet*(), чтобы избежать JOIN с большими таблицами.
Минутные витрины: агрегаты-состояния и семантика
AggregatingMergeTree (minute rollup)
CREATE TABLE agg_events_minute_state ( minute DateTime, source LowCardinality(String), view_state AggregateFunction(sum, UInt64), add_state AggregateFunction(sum, UInt64), buy_state AggregateFunction(sum, UInt64), buyers_state AggregateFunction(uniqCombined, UInt64), rev_state AggregateFunction(sum, Decimal(18,2)), p95_state AggregateFunction(quantileTDigest(0.95), Float64) ) ENGINE = AggregatingMergeTree PARTITION BY toYYYYMMDD(minute) ORDER BY (minute, source);
MV #3: накопление состояний
CREATE MATERIALIZED VIEW mv_fact_to_minute_state TO agg_events_minute_state AS SELECT toStartOfMinute(event_time) AS minute, source, sumState(event_name='view') AS view_state, sumState(event_name='add_to_cart')AS add_state, sumState(event_name='purchase') AS buy_state, uniqCombinedStateIf(user_id, event_name='purchase') AS buyers_state, sumStateIf(amount, event_name='purchase') AS rev_state, quantileTDigestStateIf(0.95)(toFloat64(1000 * rand())) AS p95_state -- пример, замените своей метрикой латентности FROM fact_events_wide GROUP BY minute, source;
Важно: храним состояния (…State) по минуте×разрезам → читаем …Merge. Это обеспечивает NRT-отчёты без «стрельбы в хранилище».
Семантическая вьюха для BI
CREATE OR REPLACE VIEW vw_kpi_minute AS SELECT minute, source, sumMerge(view_state) AS views, sumMerge(add_state) AS adds, sumMerge(buy_state) AS buys, uniqCombinedMerge(buyers_state) AS buyers, sumMerge(rev_state) AS revenue, revenue / NULLIF(buys,0) AS aov, buys / NULLIF(views,0) AS cr, quantileTDigestMerge(0.95)(p95_state) AS p95_ms FROM agg_events_minute_state GROUP BY minute, source;
Watermarks, опоздавшие события и ретро-пересборка окна
Проблема: события приходят с опозданием (минуты/часы). Если сразу «закрывать минуту», мы недосчитаем.
Решение: держим ретро-окно (например, 2–24 часа) и переагрегируем это окно периодически:
- MV (#3) даёт «первичный» минутный слой.
- Ночью (или каждый час) job пересоздаёт состояния только для окна now()-RETRO .. now() из fact_events_wide в tmp_agg_events_minute_state, затем REPLACE PARTITION в целевой:
-- пример nightly/hourly rebuild окна CREATE TABLE tmp_agg_events_minute_state AS agg_events_minute_state; INSERT INTO tmp_agg_events_minute_state SELECT ... -- тот же SELECT, что в MV #3, но WHERE minute BETWEEN now()-INTERVAL 24 HOUR AND now() ALTER TABLE agg_events_minute_state REPLACE PARTITION toYYYYMMDD(now()-INTERVAL 1 DAY) FROM tmp_agg_events_minute_state;
Для «долгоиграющих» опозданий увеличьте окно. Баланс корректности vs нагрузка фиксируем в паспорте метрик.
Backpressure, буферизация и борьба с part-explosion
Симптомы перегрузки:
- В system.parts растёт число активных частей на партицию/таблицу.
- В system.merges длинные очереди.
- Свежесть падает, Kafka lag растёт.
Меры:
- Укрупняйте микробатчи: kafka_max_block_size (30–100k), kafka_num_consumers под CPU/IO.
- Буферная таблица: если MV «бьёт» маленькими блоками, замените первую MV на вставку в Buffer/MergeTree и уже оттуда выполняйте батчевые INSERT … SELECT в факт.
- Ограничьте вычисления в MV — только лёгкие; дорогое — в периодических джобах.
- Следите за профилями пользователей BI/сервисов: max_threads, max_memory_usage, запрет FINAL.
- TTL и холодные слои для «старых» минут — уменьшают нагрузку на локальное хранилище.
Идемпотентность и дедупликация
Правило: любое событие идентифицируется (event_id, version). Повтор — либо идентичен, либо приходит с бОльшей version.
- На STAGE — ReplacingMergeTree(version);
- В FACT — не вставляем дубликаты: INSERT … SELECT с GROUP BY event_id (макс. версия) или JOIN на «последнее»;
- В rollup-слое — агрегации по состояниям устойчивы к повторной заливке одного и того же набора (при условии, что факт уникализирован).
Пример вставки «без дубля» в факт:
INSERT INTO fact_events_wide SELECT event_id, maxBy(event_time, version) AS event_time, any(user_id), any(event_name), any(amount), any(source) FROM stg_events_raw WHERE event_time >= now()-INTERVAL 3 DAY GROUP BY event_id;
async_insert vs батчи vs Kafka
|
Подход |
Когда уместен |
Плюсы |
Минусы |
|---|---|---|---|
|
Kafka→CH |
Высокая и «неровная» нагрузка, несколько продюсеров |
Дешёвая буферизация, горизонталь по партициям, надёжный ретрай |
Сложность эксплуатации, настройка lag и MV |
|
Батчи (INSERT … SELECT) |
Оффлайн/NRT-окно (секунды-минуты), контрольный ingest |
Прозрачно, повторяемо, легко ретро |
Пульсирующая свежесть, требуется планировщик |
|
async_insert |
Лёгкий поток от приложений (небольшой rps) |
Меньше накладных, быстрая запись |
Менее прозрачно под нагрузкой, сложнее гарантировать идемпотентность |
Рекомендация: для 1–100k ev/s — Kafka + MV/батчи; async_insert — для лёгких «сервисных» стримов.
Наблюдаемость ingest-а и алерты
Что мониторить в ClickHouse
- system.query_log — топ потребителей CPU/IO; отлавливаем тяжёлые MV.
- system.parts — активные части/партиции (порог на алерт).
- system.merges — длина и время мерджей.
- system.replication_queue — лаг по репликации.
- system.asynchronous_metric_log — размер файлового кэша, дисков, потоки.
Примеры запросов:
-- Части/партиции (топ) SELECT table, partition, count() parts, sum(rows) r, sum(bytes_on_disk) b FROM system.parts WHERE active GROUP BY table, partition ORDER BY parts DESC LIMIT 20; -- Мерджи (залипшие) SELECT table, elapsed, progress FROM system.merges ORDER BY elapsed DESC LIMIT 20; -- Репликация (отстающие) SELECT database, table, count() q, max(create_time) last FROM system.replication_queue GROUP BY database, table ORDER BY q DESC;
Kafka lag
Снимайте метрики lag через kafka_exporter (Prometheus) по consumer-group ch_ingest_events_v1.
Алерты:
- Lag > X (например, >100k сообщений) ≥ Y минут.
- parts/partition > порога (например, >150 активных частей в партиции).
- merges backlog > порога (elapsed > 10 мин).
- replication_queue > порога.
- доля запросов с FINAL > 0.
Дашборд ingest-здоровья: lag, parts/partition, merges, replication lag, freshness minute-витрин, топ-запросы по read_bytes/duration.
Тест-нагрузка и прогон на стенде
Генерация событий (пример псевдо-Python/kcat)
# kcat (ранее kafkacat):
seq 1000000 | kcat -b kafka-1:9092 -t events -P -l -K: <<EOF
{"event_id":"$(uuidgen)","event_time":"2025-08-05T12:34:56.000Z","event_name":"view","user_id":123,"amount":0,"source":"web"}
EOF
Совет: варьируйте размер сообщений/частоту, добавляйте 1–5% «плохих» событий для проверки карантина.
Перформанс-замер
- Включите ingestion на стенде, запустите генерацию 10–50k ev/s на 10–15 минут.
- Снимите: lag, parts/partition, merges elapsed, freshness vw_kpi_minute.
- Если лаг ≈ 0 и свежесть в SLA — профилируйте «дорогие» места (MV/обогащения).
Runbook-и (коротко)
A) «Лаг растёт, свежесть упала»
- Проверить system.merges/system.parts — part-explosion?
- Увеличить kafka_max_block_size, уменьшить число консьюмеров (баланс), временно упростить MV (выключить тяжёлые функции).
- Включить батчевую вставку: Kafka→staging (только запись), а из staging — периодический INSERT … SELECT в факт/агрегаты.
- При необходимости — остановить ретро-пересборку до стабилизации.
B) «Залипла MV, поток встал»
- system.query_log по этой MV (ошибка/таймаут).
- Проверить DLQ/карантин — нет ли «ядовитых» сообщений (JSON/тип).
- Временно отвести поток в stg_events_quarantine, снять нагрузку, починить парсер.
C) «Дубликаты в витринах/дашбордах»
- Проверить idempotency: есть ли event_id, version и логика maxBy(version).
- Проверить, не пишется ли факт параллельно двумя путями (двойная MV).
- Перебилдить ретро-окно из staging (чистый источник) → REPLACE PARTITION.
Риски и как их избегать
|
Риск |
Симптом |
Профилактика / Фикс |
|---|---|---|
|
Part-explosion |
Мерджи «захлебнулись», свежесть падает |
Укрупняйте микробатчи (kafka_max_block_size), уменьшайте kafka_num_consumers, используйте Buffer/батчи; следите за parts/partition |
|
Залипшие MV |
Lag растёт при «нулевой» записи в факт |
Упростите MV (без тяжёлых функций), вынесите обогащение в периодические джобы, включите карантин |
|
Дубликаты |
Суммы/клики «накручены» |
event_id+version, Replacing(version) на STAGE, maxBy(version) при вставке в факт |
|
Опоздавшие события |
Недосчёт минутных/часовых метрик |
Ретро-окно N часов/дней и REPLACE PARTITION; паспорт метрики фиксирует окно |
|
Тяжёлые метрики «на лету» |
P95/DAU делают запросы медленными |
Храните только …State и читайте …Merge, отдельные minute/hour rollup’ы |
|
JOIN больших на большие |
Память/таймаут |
Обогащение через словари dictGet*(); pre-join на записи |
|
Репликация отстаёт |
Разные цифры на репликах, таймауты |
Следите за replication_queue; снижайте ingest, увеличивайте фоновые пулы; не смешивайте тяжёлые мутации с горячим ingest |
|
Kafka lag «пилит» |
Всплески лагов |
Равномерный продьюсинг, достаточные партиции, backpressure на продьюсере; не «резать» блок до 1–5k |
Чек-лист NRT-контура (быстрый аудит)
- Kafka: партиций ≥ целевого параллелизма; group-id стабилен; lag < порога.
- ENGINE=Kafka настроен: consumers/блок/формат/ошибки.
- STAGE: Replacing(version), карантин для «плохих».
- FACT: уникализация по event_id; никаких дублей.
- ROLLUP: AggregatingMergeTree; только …State в записи, …Merge в чтении.
- VIEW: vw_* без FINAL и SELECT *; фильтры по партиции.
- RETRO: зафиксировано окно N часов/дней; есть job на REPLACE PARTITION.
- OBSERVABILITY: дашборд lag/parts/merges/replication/freshness; алерты.
- PERF: read_bytes и latency плиток в норме на 24–48 ч окне.
- SECURITY: BI видит только vw_*; профили/квоты/ограничители.
Что вы отдаёте на выходе (артефакты модуля)
- Конфиги ENGINE=Kafka (топики, группы, блоки, консьюмеры).
- DDL: stg_events_raw (Replacing), stg_events_quarantine, fact_events_wide, agg_events_minute_state.
- MV: mv_kafka_to_staging, mv_staging_to_fact, mv_fact_to_minute_state.
- VIEW: vw_kpi_minute/vw_kpi_hour.
- Скрипты нагрузки (kcat/генераторы), дашборд ingest-здоровья (описание метрик/алертов).
- Runbooks: лаг/залипшие MV/дубликаты/ретро-пересборка.
Итог
Стриминг в ClickHouse «летает», если соблюдать три дисциплины:
- Правильная физика: микробатчи, Replacing на STAGE, уникальный факт, rollup в Aggregating-состояния, никакого FINAL.
- Семантика минут: водяные знаки и ретро-окно для опоздавших событий; тяжёлые метрики — только через …State/…Merge.
- Наблюдаемость и гвардейлы: lag/parts/merges/replication под мониторингом, алерты и runbook-и заранее.
Arenadata QuickMarts (ADQM) — корпоративная платформа на базе ClickHouse для быстрого слоя витрин и near-real-time аналитики. Решает задачи «быстрых» дашбордов и API с низкой латентностью и высокой конкуррентностью, работает поверх вашего DWH/лейкхауса как serving-уровень. Даёт предсказуемую производительность на терабайтно-петабайтных объёмах за счёт колоночного хранения, компрессии и предагрегатов (Materialized Views, AggregatingMergeTree), подключается к Kafka/S3 и стандартным BI-инструментам по SQL/HTTP. Для корпоративных ИТ ADQM предлагает поддержку и SLA, отказоустойчивые кластеры (HA/DR), безопасность (RBAC, LDAP/OIDC, шифрование трафика и данных), мониторинг и резервное копирование. Платформа хорошо ложится на методологию курса: семантика vw_*, роллап-слои, NRT-ингест, SLO/наблюдаемость и «гвардейки» для BI/API. Итог — быстрый запуск витрин за недели, снижённые риски в проде и предсказуемая стоимость владения.



