Модуль 12. Capstone-проект: полный контур витрин (от CORE до BI)
Это практическая глава «собери и запусти всё»: архитектура, DDL, пайплайны, DQ/SLA, безопасность, дашборды, наблюдаемость и процедуры. Сделаем так, чтобы за 1–2 спринта получить рабочую систему, которую не стыдно показать бизнесу и аудиторам.
Что строим и что считаем «готово»
Цель: поднять слой витрин на ClickHouse, подключить BI, включить DQ/наблюдаемость и доказать SLA.
Готово =
- Свежесть данных ≤ SLA (оперативные 5–15 мин, дневные ≤ 60 мин).
- Баланс витрин vs CORE в допуске (по умолчанию ≤ 0.2% за окно 30 дней).
- Вьюхи без FINAL/SELECT *; тяжёлые метрики (uniq, quantile, top-K) читаются из состояний.
- Дашборды открываются в пределах целевой латентности; read_bytes под контролем.
- DR-восстановление на тестовом стенде прошло (smoke-тест «зелёный»).
Два трека практики:
- Retail: продажи+возвраты, мультивалюта «на дату операции», минутные KPI, дашборд CR/AOV/Top-K.
- SaaS: активность/подписки/NRR, когорты/ретеншн, «микро» A/B, дашборд NRR/DAU/p95.
Картина мира и артефакты спринта
Архитектура (минимум жизнеспособный слой)
CORE (DDS, внешние БД/стримы) │ batch/CDC/stream ▼ STAGE (валидатор/карантин/дедуп) │ INSERT INTO ... SELECT ... ▼ MARTS (факты wide/agg-state, SCD-измерения) │ ├─ rollup (minute/hour/day AggregatingMergeTree) ▼ SEMANTIC (vw_*: метрики как код, календарь/валюта, RLS/маскирование) │ ├─ BI (дашборды) └─ mini-API (опционально) + Observability (freshness, query_log, parts/merges, DQ)
Репозиторий (скелет)
repo/
infra/terraform | ansible | k8s
sql/
tables/ # DDL (v1/v2 отдельно)
views/ # vw_*.sql (CREATE OR REPLACE)
dicts/ # словари (валюты/статусы)
testdata/ # сиды/синтетика
metrics/ # паспорта метрик (YAML)
tests/ # DQ/регресс/перформанс (SQL)
ci/ # pipeline (линтеры, сборка, тесты, деплой)
docs/ # автодоки (md/html)
План работ на 1–2 спринта (примерный ритм)
Спринт 1 (неделя)
День 1–2:
- Поднять ClickHouse (dev/stage) и настроить макросы {cluster}.
- Завести storage policy (hot/warm/cold), роли/квоты/профили (RBAC).
- Сборка таблиц STAGE/MARTS (wide + agg-state), словари (валюты/статусы).
День 3–4:
- Наполнить историю (backfill) по партициям (30–90 дней).
- Включить инкремент (batch/CDC/stream).
- Сделать семантические VIEW (vw_*) и базовый rollup (hour/day).
День 5:
- DQ-набор: свежесть, баланс vs CORE, дубли зерна, плотность рядов, инварианты.
- Дашборд Observability (freshness, read_bytes, parts/merges).
- Мини-дашборд для предметного домена (Retail или SaaS).
Спринт 2 (неделя)
- Добавить минутные KPI и MV (если есть NRT).
- BI-подключение, дефолт-фильтры и guardrails.
- CI/CD: линтеры (FINAL/SELECT * запрет), регрессы v1↔v2, перформанс-реплей.
- DR-процедура: BACKUP/RESTORE, smoke-тест.
- Приёмка по чек-листу; flip на prod/alias (при наличии).
Выбор модели и ключей (общий каркас)
Таблицы MARTS: wide + agg-state
Широкий факт (retail-пример):
CREATE TABLE mart_sales_wide ( tx_id UInt64, tx_datetime DateTime, day Date MATERIALIZED toDate(tx_datetime), shop_id UInt32, sku_id UInt32, qty Int32, amount Decimal(14,2), currency FixedString(3), status_canon LowCardinality(String), buyer_id UInt64 ) ENGINE = MergeTree PARTITION BY toYYYYMM(day) ORDER BY (day, shop_id, sku_id, tx_id);
Агрегаты-состояния (день):
CREATE TABLE agg_sales_daily_state ( day Date, shop_id UInt32, category_id UInt32, amount_state AggregateFunction(sum, Decimal(18,2)), qty_state AggregateFunction(sum, Int64), buyers_state AggregateFunction(uniqCombined, UInt64) ) ENGINE = AggregatingMergeTree PARTITION BY toYYYYMM(day) ORDER BY (day, shop_id, category_id);
Семантическая вьюха:
CREATE OR REPLACE VIEW vw_retail_daily AS SELECT day, shop_id, category_id, sumMerge(amount_state) AS net_sales, sumMerge(qty_state) AS qty, uniqCombinedMerge(buyers_state) AS buyers, net_sales / NULLIF(qty, 0) AS aov FROM agg_sales_daily_state GROUP BY day, shop_id, category_id;
Правила:
- Партиция — время (месяц/день); ORDER BY — префикс типовых WHERE.
- Тяжёлые метрики = состояния (…State на записи, …Merge на чтении).
- Никакого FINAL в продуктивном чтении.
Трек Retail (продажи/возвраты, мультивалюта, минутные KPI)
Курсы валют и «на дату операции»
Справочник:
CREATE DICTIONARY dict_fx ( date Date, from FixedString(3), to FixedString(3), rate Decimal(18,8) ) PRIMARY KEY (date, from, to) SOURCE(CLICKHOUSE(HOST 'localhost' PORT 9000 USER 'default' DB 'core' TABLE 'fx_rates')) LAYOUT(FLAT());
ASOF-присоединение цены/курса:
-- цена на момент транзакции
SELECT s.tx_datetime, s.sku_id, s.qty, p.price
FROM mart_sales_wide s
ASOF JOIN price_snapshots p
ON p.sku_id = s.sku_id AND p.valid_from <= s.tx_datetime;
-- перевод суммы в базовую валюту
SELECT day, shop_id,
sum( amount * dictGetDecimal64('dict_fx','rate', (toDate(tx_datetime), currency, 'USD')) ) AS net_sales_usd
FROM mart_sales_wide
WHERE day BETWEEN today()-30 AND today()-1
GROUP BY day, shop_id;
Минутные KPI (CR, AOV, Top-K)
Событийный слой → минутный rollup:
CREATE TABLE agg_events_minute_state ( minute DateTime, shop_id UInt32, view_state AggregateFunction(sum, UInt64), add_state AggregateFunction(sum, UInt64), buy_state AggregateFunction(sum, UInt64), buyers_state AggregateFunction(uniqCombined, UInt64), aov_state AggregateFunction(sum, Decimal(18,2)) -- суммарный чек ) ENGINE = AggregatingMergeTree PARTITION BY toYYYYMMDD(minute) ORDER BY (minute, shop_id);
Вьюха:
CREATE OR REPLACE VIEW vw_kpi_minute AS SELECT minute, shop_id, sumMerge(view_state) AS views, sumMerge(add_state) AS adds, sumMerge(buy_state) AS buys, uniqCombinedMerge(buyers_state) AS buyers, sumMerge(aov_state) AS revenue, buys / NULLIF(views,0) AS cr FROM agg_events_minute_state GROUP BY minute, shop_id;
Дашборд (примеры запросов)
- CR/AOV за сутки по часам:
SELECT toStartOfHour(minute) AS hour, sum(buys)/NULLIF(sum(views),0) AS cr,
sum(revenue)/NULLIF(sum(buys),0) AS aov
FROM vw_kpi_minute
WHERE minute >= now()-INTERVAL 24 HOUR AND shop_id = {shop:UInt32}
GROUP BY hour ORDER BY hour;- Top-10 SKU за вчера:
SELECT shop_id, sku_id, sum(net_sales) s FROM vw_retail_daily WHERE day = yesterday() GROUP BY shop_id, sku_id ORDER BY shop_id, s DESC LIMIT 10 BY shop_id;
Риски и фиксы (Retail):
- Возвраты «задним числом» → nightly ретро-окно 14–30 дней.
- Курсы валют «на дату отчёта» vs «на дату операции» → две вьюхи, прописать в паспорте.
- Перцентили/uniq на лету → только состояния.
- Узкий ORDER BY → чтение «пол-таблицы» → перепроектировать v2-таблицу.
Трек SaaS (подписки/активность, NRR, когорты, p95)
Подписки и NRR
CREATE TABLE mart_subscriptions ( sub_id UInt64, account_id UInt64, valid_from Date, valid_to Date, mrr Decimal(14,2), change_type LowCardinality(String) ) ENGINE = MergeTree PARTITION BY toYYYYMM(valid_from) ORDER BY (account_id, valid_from);
NRR (упрощённо):
CREATE OR REPLACE VIEW vw_mrr_month AS SELECT toStartOfMonth(valid_from) AS month, sumIf(mrr, change_type='START') AS mrr_start, sumIf(mrr, change_type='END') AS mrr_end, sumIf(mrr, change_type='EXPANSION') AS expansion, sumIf(-mrr, change_type='CHURN') AS churn, (mrr_start + expansion - churn) / NULLIF(mrr_start,0) AS nrr FROM mart_subscriptions GROUP BY month;
DAU/p95 перформанса
Состояния:
CREATE TABLE agg_app_minute_state ( minute DateTime, account_id UInt64, dau_state AggregateFunction(uniqCombined, UInt64), p95_state AggregateFunction(quantileTDigest(0.95), Float64) ) ENGINE = AggregatingMergeTree PARTITION BY toYYYYMMDD(minute) ORDER BY (minute, account_id);
Чтение:
SELECT minute, account_id,
uniqCombinedMerge(dau_state) AS dau,
quantileTDigestMerge(0.95)(p95_state) AS p95_ms
FROM agg_app_minute_state
WHERE minute >= now()-INTERVAL 24 HOUR
GROUP BY minute, account_id;
Когорты/ретеншн (идея)
- Вьюха vw_user_cohort (первая активация), затем денсить матрицу «cohort_day × day», считать D7/D30 retention.
- Хранить дневные показатели когорт (не пересчитывать «на лету» на годы).
Риски и фиксы (SaaS):
- SRD (slow-rollout drift) по регионам → стратификация отчётов.
- Неправильный календарь (финансовые месяцы) → явные справочники календаря.
- p95 из сырых → только TDigest-состояния.
DQ и SLA: «здоровье» витрин
Таблицы наблюдаемости
CREATE TABLE sem_meta (view_name String, updated_at DateTime) ENGINE=MergeTree ORDER BY view_name; CREATE TABLE dq_results ( test_name String, scope String, ts DateTime, status LowCardinality(String), value Float64, threshold Float64, details String ) ENGINE=MergeTree ORDER BY (test_name, ts);
Набор тестов (SQL/идея)
- Свежесть: now() - updated_at <= SLA.
- Баланс vs CORE (вчера/окно):
WITH m AS (SELECT sum(net_sales) s FROM vw_retail_daily WHERE day=yesterday()),
c AS (SELECT sum(amount) s FROM core.sales WHERE toDate(ts)=yesterday() AND status='PAID')
SELECT abs(m.s-c.s)/NULLIF(c.s,0) <= 0.002 AS ok FROM m,c;- Дубли зерна, плотность ряда, инварианты (GM ≤ NetSales, rate ∈ [0,1]).
- Аномалии: z/MAD на скользящем окне (для ключевых KPI).
Алерты и панели
- Алерты: freshness, balance, parts/merges backlog, replication lag, доля FINAL>0.
- Графики: read_bytes/query_duration (топ-10 запросов), parts/partition, merges latency.
CI/CD и guardrails
Линтеры (пример)
- Запрет FINAL и SELECT * в views/.
- Обязательный фильтр по времени в vw_nrt_*.
- Запрет SummingMergeTree на ретро-корректируемых данных.
- Проверка, что изменение views/*.sql сопровождается bump metrics/*.yaml.
Пайплайн
- PR → линтеры.
- Сборка DDL на mini-кластер, загрузка testdata.
- Прогон tests/ (DQ/регресс/perf-replay).
- Stage-деплой ON CLUSTER.
- Prod: side-by-side (_v2), dual-run (7–30 дней), alias-flip, пост-мониторинг.
Безопасность: RBAC/RLS/маскирование
CREATE ROLE semantic_reader, marts_dev;
CREATE USER bi_ro IDENTIFIED BY '***';
GRANT semantic_reader TO bi_ro;
GRANT SELECT ON db_marts.vw_* TO semantic_reader;
CREATE ROW POLICY rp_region ON db_marts.mart_sales_wide
FOR SELECT USING region_id = currentSetting('region_id');
ALTER USER bi_ro SETTINGS region_id=77;
CREATE OR REPLACE VIEW db_marts.vw_sales_masked AS
SELECT day, shop_id, net_sales,
substring(sha256Hex(toString(buyer_id)),1,12) AS buyer_hash
FROM db_marts.vw_retail_daily;
Профили/квоты для BI: max_execution_time, max_memory_usage, max_threads, quota per minute.
DR и хранение: TTL/S3/бэкапы
- TTL MOVE: горячие 90 дней на NVMe → warm → S3 (cold).
ALTER TABLE mart_sales_wide
MODIFY TTL day + INTERVAL 90 DAY TO VOLUME 'warm',
day + INTERVAL 365 DAY TO VOLUME 'cold';- BACKUP/RESTORE и DR-день (квартально):
BACKUP DATABASE db_marts TO Disk('backups','2025-08-01/db_marts');
RESTORE DATABASE db_marts FROM Disk('backups','2025-08-01/db_marts');
BI-подключение и «красные линии»
- Подключать только к vw_*.
- Дефолтные фильтры: период (последние 28/90 дней), лимиты строк/Top-N.
- Запрет тяжёлых JOIN/перцентилей «на лету» — всё есть в VIEW.
- Для NRT-плиток — отдельные источники (vw_kpi_minute/vw_kpi_hour).
- Пагинация «по ключу», не OFFSET.
Чек-лист приёмки Capstone
- Freshness ≤ SLA (vw_* обновляются вовремя).
- Баланс vs CORE ≤ 0.2% на окне 30 дней.
- Вьюхи без FINAL/SELECT *, тяжёлые метрики — через …Merge.
- Latency плиток/отчётов ≤ целевых; read_bytes в норме.
- DQ-панель зелёная (дубли, плотность, инварианты, аномалии).
- RBAC/RLS/маскирование включены; BI → только vw_*.
- DR: backup/restore прошёл, smoke-тесты ок.
- CI/CD: PR-линтеры, stage-сборка, тесты, prod-flip с откатом.
- Документация: паспорта метрик (YAML), lineage, инструкции.
Runbooks (короткие инструкции)
A. Свежесть просела
- Проверить Kafka lag / MV / system.merges / system.replication_queue.
- Временно выключить глубокие ретро/overlay; укрупнить батчи.
- OPTIMIZE проблемных партиций (точечно).
- Коммуникация SLA-отклонения и ETA.
B. BI стало медленно
- system.query_log: топ по read_bytes/duration.
- Убрать FINAL/SELECT *; заменить uniq/quantile на …Merge.
- Скорректировать WHERE под ORDER BY; добавить skip-индексы (по факту выгоды).
- При необходимости — вынести pre-aggregates.
C. Баланс vs CORE красный
- Проверка новых статусов/валют/календаря.
- Ретро-пересчёт окна (14–30 дней) из чистого источника.
- Запись причин в changelog метрики (v2, если логика изменилась).
Риски и митигации (сводная таблица)
|
Риск |
Проявление |
Митигация |
|---|---|---|
|
Несостыковка календарей/валют |
отчёты «не бьются» |
Разные vw_* для «операция/отчёт», календарь в паспорте |
|
Part-explosion |
отставание merges, свежесть падает |
Крупные батчи, буферные таблицы, алерты по parts |
|
Тяжёлые метрики «на лету» |
минуты/таймауты |
AggregatingMergeTree и …State/…Merge, rollup |
|
FINAL в продуктивных VIEW |
провалы SLA |
Запрет в линтерах, дисциплина записи |
|
JOIN «большого с большим» |
память/таймаут |
Предагрегаты/словари/ASOF, pre-join при записи |
|
S3 без кэша |
пила по латентности |
«Горячее окно» локально, filesystem cache, TTL MOVE |
|
Утечка PII |
комплаенс-риски |
RLS/маскирование, BI → только vw_*, аудит query_log |
Что отдать на выходе (артефакты Capstone)
- Репозиторий: infra/sql/views/dicts/tests/metrics/docs/ci.
- DDL: wide/agg-state, словари, календарь.
- VIEW (vw_*): метрики и правила (календарь/валюта), NRT-вьюхи.
- DQ-набор: freshness, баланс, дубли, плотность, инварианты, аномалии.
- CI/CD: линтеры, сборка, тесты, stage→prod flip/rollback.
- RBAC/RLS и маскирование, профили/квоты.
- Observability-дашборды: свежесть, parts/merges/replication, топ-запросы.
- BI-дашборды: Retail (CR/AOV/Top-K) или SaaS (NRR/DAU/p95).
- Runbooks и DR-процедуры.
Заключение
Capstone — это не «ещё один модуль», а итоговая сборка всех практик: правильная физика данных, метрики как код, агрегаты-состояния, наблюдаемость, безопасность и дисциплина релизов. Собрав такой контур за 1–2 спринта, вы получите витринный слой, который:
а) быстро отвечает, б) показывает одни и те же цифры всем, в) переживает ретро-изменения и аварии, г) контролируем по SLA и стоимости.
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. Итог — быстрый запуск витрин за недели, снижённые риски в проде и предсказуемая стоимость владения.



