Модуль 3. Потоки данных, эксплуатация и надёжность ClickHouse-витрин
После того как у нас есть модель витрин (Модуль 2) и семантика (Модуль 1), самая частая причина проблем — подача данных и эксплуатация: «залипли» мерджи, Kafka прислала дубли, S3 тормозит, BI «падает» по памяти, ретро-пересчёт съел ночь. Здесь разбираем полный цикл: ingest → материализация → хранение → доступ → мониторинг → бэкапы → релизы.
Архитектуры потоков: откуда и как
Типовые источники
- БД приложений (OLTP): PostgreSQL/MySQL и т. п. — чаще CDC (change data capture) или ночные/часовые срезы.
- Очереди/шины: Kafka (стандарт для ClickHouse), реже другие брокеры.
- Файлы/объектное хранилище: S3/MinIO/HDFS/локальные drop-зоны — Parquet/CSV/TSV/JSON.
- Сервисы/логирование: веб/мобильные события, телеметрия (ClickHouse отлично «ест» события).
Паттерны загрузки
- Batch (микро-батчи): регулярные вставки «крупными блоками» → MergeTree. Простая эксплуатация, прогнозируемые ресурсы.
- Stream (near real-time): Kafka Engine → Materialized View → MergeTree. Нужны дисциплина ключей и микробатчи.
- CDC: Debezium/Maxwell → Kafka → ClickHouse. Важны «идемпотентность», обработка tombstone/updates.
Риск: смешать все три паттерна в одной витрине.
Митигация: чётко разделяйте слои (RAW/STAGE/CORE/MARTS) и «внешние» таблицы (ingest) vs «целевая витрина».
Batch-загрузка: S3/файлы/табличные функции
Разовые и регулярные заливки
- Форматы: Parquet (предпочтительно), CSV/TSV, JSONEachRow.
-
Способы:
- INSERT INTO target SELECT * FROM s3('s3://bucket/path/*.parquet', 'KEY', 'SECRET')
- INSERT INTO target SELECT * FROM url('https://...', 'Parquet')
- INSERT FROM file('*.parquet') для локальных дроп-зон.
Практика:
- Загружайте крупными блоками (сотни тысяч/миллионы строк за операцию).
- При необходимости — промежуточная «стейдж-таблица» и далее в целевую INSERT … SELECT с нормализацией (статусы/валюты/календарь).
Риски и митигации:
- Дубли из-за повторной заливки: делайте дедуп на STAGE, в целевую — AggregatingMergeTree (состояния) либо Replacing(version).
- «Мелкая дробь» файлов: объединяйте upstream (Spark/EMR) или через периодические INSERT SELECT (микро-батчи).
Streaming/CDC: Kafka Engine → MV → MergeTree
Минимальный конвейер
-- 1) Входной топик
CREATE TABLE raw_sales_kafka
(
event_time DateTime,
event_id String,
shop_id UInt32,
sku_id UInt32,
qty Int32,
amount Decimal(12,2),
currency FixedString(3),
status LowCardinality(String)
)
ENGINE = Kafka
SETTINGS
kafka_broker_list = 'kafka1:9092',
kafka_topic_list = 'sales',
kafka_group_name = 'ch-sales-consumers',
kafka_format = 'JSONEachRow',
kafka_num_consumers = 6;
-- 2) Целевая таблица (wide) под витрину
CREATE TABLE mart_sales_wide
(
tx_datetime DateTime,
day Date MATERIALIZED toDate(tx_datetime),
event_id String,
shop_id UInt32,
sku_id UInt32,
qty Int32,
amount Decimal(12,2),
currency FixedString(3),
status_canon LowCardinality(String)
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(day)
ORDER BY (day, shop_id, sku_id, event_id);
-- 3) Материализованное представление (микробатчи)
CREATE MATERIALIZED VIEW mv_sales_ingest TO mart_sales_wide AS
SELECT
event_time AS tx_datetime,
event_id,
shop_id,
sku_id,
qty,
amount,
currency,
transform(status,
['paid','captured','refund','cancelled'],
['PAID','PAID','REFUND','CANCELLED'],
'UNKNOWN') AS status_canon
FROM raw_sales_kafka;
Ключевые настройки и приёмы:
- Микро-батчи: регулируйте размер блоков на входе/частоту флашей (иначе получите тысячи мелких parts).
- Идемпотентность: event_id + дедуп на STAGE/целевой; для CDC — ключ + версия/ts.
- Не используйте FINAL поверх «грязных» апдейтов — добивайтесь чистоты на записи.
CDC (Debezium-стиль)
- События типов c/u/d (create/update/delete) + before/after.
-
В ClickHouse лучше нормализовать поток в STAGE (развернуть в стандартный upsert) и писать в целевую:
- либо в ReplacingMergeTree(version) (не читать с FINAL в BI),
- либо в «журнал событий» + ночной пересчёт в агрегаты.
Риски и митигации:
-
Повтор чтения после рестарта консьюмера → дубли.
→ Строгий ключ идемпотентности (event_id + версионность), дедуп на STAGE. -
«Томбы» (delete) ломают суммы.
→ Моделируйте -amount событием или отдельной таблицей корректировок и nightly overlay.
Трансформации внутри ClickHouse: MV, pre-join, pre-aggregate
Что считать «на лету», а что заранее
- Заранее (pre-join): статусы → канонические, календарь, валюты, SCD-атрибуты «на дату факта».
- Заранее (pre-aggregate): частые группировки; uniq/quantile/avg → состояния (…State/…Merge).
- На лету: лёгкие фильтры/срезы/презентация (через VIEW).
Паттерн «MV → AggregatingMergeTree»
CREATE TABLE agg_sales_daily_state ( day Date, shop_id UInt32, category_id UInt32, amount_state AggregateFunction(sum, Decimal(14,2)), qty_state AggregateFunction(sum, Int64) ) ENGINE = AggregatingMergeTree PARTITION BY toYYYYMM(day) ORDER BY (day, shop_id, category_id); CREATE MATERIALIZED VIEW mv_sales_to_daily TO agg_sales_daily_state AS SELECT toDate(tx_datetime) AS day, shop_id, sku.category_id, sumState(amount) AS amount_state, sumState(qty) AS qty_state FROM mart_sales_wide s LEFT JOIN d_sku sku ON sku.sku_id = s.sku_id GROUP BY day, shop_id, sku.category_id;
Риск: MV пересчитывает «всю историю» при каждом сообщении.
Митигация: агрегируйте только свежие окна (минуты/часы/дни), вход — микробатчи.
Хранение: партиции, TTL, холодные слои, кодеки
Партиционирование
- События/лог-стримы: дневные партиции (часто достаточно).
- Продажи/финансы: месячные (удобно для ретро-пересчётов).
- Делайте так, чтобы ретро-операции (пересборка/оптимизация) делались по одной партиции.
TTL и tiering
- TTL DELETE: удаляем старые строки (с осторожностью для фактов).
- TTL MOVE TO volume 'cold': «перевоз» старых партиций на дешёвое хранилище (S3/объектное).
- Политики хранения: «горячее» (NVMe) на N дней/недель, «холодное» (S3) — всё остальное.
Риск: TTL-правило унесло не то.
Митигация: обкатка в стейдже, точечное включение, мониторинг system.part_log.
Сжатие и типы
- Числа — целочисленные/Decimal (не Float, если это деньги).
- Строки — LowCardinality(String) для доменов.
- Кодеки — по умолчанию ок; точечно можно использовать ZSTD для текстов.
Репликация, шардинг, Distributed
ReplicatedMergeTree и ClickHouse Keeper
- Для HA каждой «горячей» таблице — 2 реплики на шард.
- Следите за system.replication_queue (лаг, ошибки).
Шардинг и Distributed-таблицы
- Выбор ключа шардинга под типовые фильтры: дата/разрез (регион/магазин/аккаунт).
- Распределённые запросы: делайте локальные агрегаты на шардах, затем собрать на кооринаторе.
Риск: Distributed «тормозит», локальные — быстрые.
Митигация: добейтесь пушдауна WHERE на шарды, правильного шардинга, и держите объём пересылаемых данных после локальной агрегации минимальным.
Distributed DDL
- Используйте ON CLUSTER для согласованных DDL, держите единые имена БД/таблиц и макросы ({shard}, {replica}) в конфигурации.
Бэкапы, восстановление, DR
Встроенные бэкапы
- BACKUP TABLE db.table TO Disk('backups', '2025-08-01/table')
- RESTORE TABLE db.table FROM Disk('backups', '2025-08-01/table')
Практика:
- Бэкапить метаданные и «критичные» партиции ежедневно; полный — реже.
- Хранить бэкапы в другом failure-домене (другая зона/облако/S3).
DR-сценарии
- Репликация в второй регион (async), холодное поднятие кластера из бэкапов.
- Тест восстановления раз в квартал: «через бэкап собрать витрину за X дней».
Риск: бэкап «зелёный», а восстановление не проверяли.
Митигация: регламент «DR-день»: поднять копию на стейдже из бэкапов и прогнать smoke-тесты.
Observability: что мониторить и где алертить
Полезные системные таблицы
- system.query_log — кто/что/как долго читает, сколько байт/строк.
- system.part_log — создание/мердж/удаление частей.
- system.merges — активные мерджи; system.replication_queue — лаг репликации.
- system.asynchronous_metrics — сводные счётчики; system.metrics — живые метрики.
Запросы-«детекторы»:
-- Топ тяжёлых запросов (за сутки) SELECT any(user) AS user, query, sum(read_bytes) AS bytes, sum(read_rows) AS rows, avg(query_duration_ms) AS ms FROM system.query_log WHERE event_time >= now() - INTERVAL 1 DAY AND type='QueryFinish' GROUP BY query ORDER BY bytes DESC LIMIT 50; -- Сводка по частям (активные parts по партициям) SELECT table, partition, count() AS parts, sum(rows) r, sum(bytes_on_disk) b FROM system.parts WHERE active GROUP BY table, partition ORDER BY parts DESC LIMIT 50;
Ключевые алерты
- Freshness витрин > SLA (по sem_meta.view_name/updated_at).
- Parts per partition > порога (например, > 20k).
- Merge backlog/elapsed растёт более N минут.
- Replication lag > порога.
- Доля запросов с FINAL > X% (линтер SQL в CI + дашборд).
- DQ-сигналы: расхождение с CORE > порога, новые «UNKNOWN» статусы.
Производительность: быстрые выигрыши
- ORDER BY под реальный WHERE (дата → разрез → id).
- Крупные батчи вставок; избегать «дроби» на kafka/mv.
- Агрегаты-состояния там, где BI регулярно считает uniq/quantile/avg.
- Без FINAL для отчётов: добивайтесь консистентности на записи.
- Срезы по hot-окну: NVMe для N дней/недель, остальное — холод/кэш.
Антипаттерны:
- Summing на данных с ретро-правками.
- Длинные ORDER BY «на всякий случай».
- MV, пересчитывающие всю историю «на каждое событие».
Безопасность и изоляция
- RBAC: роли semantic_reader (только vw_*), marts_dev, marts_admin.
- RLS: политики строк на базовых таблицах (регион/тенант).
- PII-маскирование: представления с обрезанными/хэшированными id, BI видит только vw_*.
- Квоты/профили: лимиты на память/время/чтение для групп пользователей.
- Аудит: логин/кто/когда/что читал — логируйте и храните.
CI/CD для ClickHouse
Репозиторий
/sql/tables/*.sql -- DDL таблиц (v2-варианты отдельно) /sql/views/*.sql -- семантические VIEW (CREATE OR REPLACE) /metrics/*.yaml -- паспорта метрик /tests/*.sql -- DQ/регрессионные тесты /ci/* -- скрипты деплоя, линтеры, проверки
Пайплайн
- PR → линтер (нет FINAL в VIEW, нет DROP в прод-скриптах).
- Деплой на стейдж ON CLUSTER, наполнение тестовым окном.
- Прогон DQ и регресс-сравнения v1 vs v2 на окне N дней.
- Нагрузочный тест топ-запросов (replay из query_log).
- Прод-деплой: side-by-side, alias/view-swap, пост-мониторинг.
Кейсы (end-to-end)
eCom real-time: события корзины и конверсии
Цель: дашборд «в реальном времени» (± минуты) — конверсия, AOV, p95 latency.
Поток:
- Источник: фронт/бэк → Kafka (cart_add, checkout, purchase).
- CH: Kafka Engine (6–12 консьюмеров) → MV в events_wide.
- Агрегаты: agg_minute_state (sumState/uniqCombinedState/quantileTDigestState).
- Семантика: vw_kpi_minute с …Merge и rolling-окнами.
Риски: «бурст» вечером → тысячи маленьких частей.
Как гасить: микробатчи, буферные таблицы, алерты parts/merges, ограничить поля в событии (только необходимые).
Финансы: CDC из PostgreSQL, «остатки и комиссии»
Поток:
- Debezium → Kafka (transactions c c/u/d).
- STAGE: нормализация в upsert-поток, дедуп по (tx_id, version).
- MARTS: ReplacingMergeTree(version) для «сырых» фактов + nightly overlay.
- Агрегат: agg_balance_daily_state/agg_revenue_daily_state (…State/…Merge).
- Семантика: vw_balance_daily/vw_revenue_daily.
Риски: «томбы»/повторы, валюты.
Гашение: ключ идемпотентности, ретро-окно 14–30 дней, snapshot-курс в факте, nightly сверка vs GL.
IoT/Telecom: CDR/телеметрия
Поток:
- Kafka → cdr_kafka → MV → cdr_wide.
- Агрегаты: agg_cdr_minute_state (ok/all/latency p95/avg).
- Семантика: vw_cell_kpi_hour (hourly rollup).
- Мониторинг: алерты success_rate<0.98, p95_latency выше порога.
Риски: дубликаты после рестартов, холодные партиции на S3.
Гашение: дедуп event_id, локальный кэш для S3, «горячее окно» локально.
Runbooks (сжатые сценарии)
A) Витрина отстаёт по свежести
- Проверить lag ingestion (Kafka offsets), состояние MV.
- Проверить parts/merges backlog.
- Временно сузить ретро-пересчёт, отключить тяжёлые overlay.
- Точечный OPTIMIZE проблемных партиций.
- Сообщить ETA бизнесу.
B) Резкий рост времени запросов
- Снять топ из query_log по байтам/строкам.
- Сверить ORDER BY vs WHERE; включить/подстроить skip-индексы.
- Выделить агрегаты-состояния; лимиты/фильтры в BI по умолчанию.
C) Репликация отстаёт
- system.replication_queue (ошибки/лаг).
- IO/сеть, фоновый пул.
- Точечное устранение «битых» задач; при необходимости — снять нагрузку (ингест, тяжёлые запросы).
Частые ошибки и как их не повторять
-
Summing на данных с ретро-правками → удвоения.
Решение: Aggregating (…State/…Merge) или rebuild окна. -
FINAL в продуктивных VIEW → деградация SLA.
Решение: обеспечивать консистентность на записи, линтер в CI. -
Смешанный календарь/валюта в одном представлении → «не бьются» отчёты.
Решение: отдельные VIEW, паспорт метрики. -
Чрезмерно длинный ORDER BY → большой overhead без пользы.
Решение: короткий префикс под реальный WHERE. -
Мелкие вставки (part explosion) → «задыхаются» мерджи.
Решение: микробатчи, буферные таблицы, контроль parts. -
MV пересчитывает «вселенную» на каждое событие.
Решение: агрегировать только «свежее окно», дробить по времени.
Итог
Надёжные и быстрые витрины в ClickHouse — это сочетание:
- правильного ingest-конвейера (batch/stream/CDC) с идемпотентностью;
- разумной материализации (pre-join и агрегаты-состояния);
- хорошей «физики» (партиции, ORDER BY, tiering/TTL);
- репликации/шардинга для отказоустойчивости и параллельной мощности;
- бэкапов и практикуемого восстановления;
- наблюдаемости и DQ (freshness, parts/merges, репликация, баланс с CORE);
- строгого CI/CD (версии, тесты, side-by-side миграции).
С такой «операционной рамкой» вы выдержите рост нагрузки, неожиданные корректировки, смену логики и при этом сохраните SLA для BI.
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. Итог — быстрый запуск витрин за недели, снижённые риски в проде и предсказуемая стоимость владения.



