Модуль 8. Продвинутые аналитические паттерны ClickHouse
(временные ряды и ASOF JOIN, последовательности и воронки, гео-аналитика, JSON/массивы/Map, bitmap/Top-K/квантили, пересечения аудиторий, сервисный слой и мини-API поверх витрин)
О чём этот модуль и когда он нужен
Вы построили слой витрин (Модуль 2), семантику (Модуль 1), наладили потоки и эксплуатацию (Модуль 3), сделали быстрые дашборды (Модуль 4) и DQ/изменения (Модуль 5–6). Следующий уровень — сложные аналитические задачи, где стандартных сумм/средних мало:
- временные задачи: gap-filling, роллинги, as-of join («цены/курсы/версия справочника на момент события»);
- последовательности и воронки (sequenceMatch/Count, windowFunnel) с сессиями;
- гео-аналитика (радиусы, полигоны, геосетки, покрытия);
- nested-данные: JSON/массивы/Map — как хранить, читать и не «убить» колонки;
- приближённые агрегаты (uniq*, topK, quantile*, bitmap) для «тяжёлых» метрик;
- пересечения/аудитории (set-операции и groupBitmap*);
- сервисный слой: как безопасно отдавать это в продуктовые мини-API (HTTP/JDBC), не проломив SLA и безопасность.
Во всех разделах — SQL, «как объяснил бы сотруднику», и риски + митигации.
Время и временные ряды: densify, роллинги, ASOF JOIN
«Денсить» календарь (убрать «дыры»)
Для корректных роллингов/ретеншна важно, чтобы у вас была плотная матрица дат × ключ. Делайте системную таблицу d_calendar и левый join:
-- дни за 90 дней и все магазины
WITH days AS (
SELECT d FROM db_marts.d_calendar
WHERE d BETWEEN today()-90 AND today()-1
),
shops AS (
SELECT DISTINCT shop_id FROM db_marts.vw_net_sales_daily
)
SELECT d AS day, s.shop_id,
coalesce(sales, 0) AS net_sales
FROM days
CROSS JOIN shops s
LEFT JOIN (
SELECT day, shop_id, sum(net_sales) AS sales
FROM db_marts.vw_net_sales_daily
WHERE day BETWEEN today()-90 AND today()-1
GROUP BY day, shop_id
) t USING (day, shop_id)
ORDER BY day, shop_id;
Риск: «дырявый» ряд ломает оконные функции и YoY.
Митигация: всегда плотните ряд; используйте один календарь (обычный или 4-5-4) на метрику.
Роллинги (скользящие окна) без самосоединений
SELECT
day, shop_id,
sum(net_sales) AS s,
sum(s) OVER (PARTITION BY shop_id ORDER BY day
ROWS BETWEEN 27 PRECEDING AND CURRENT ROW) AS roll28
FROM db_marts.vw_net_sales_daily
WHERE day >= today()-90
GROUP BY day, shop_id
ORDER BY day, shop_id;
Риск: роллинг по «дыркам» искажен.
Митигация: см. 1.1; проверяйте плотность перед публикацией.
ASOF JOIN («на момент события»)
Задача: к каждому событию прицепить «актуальное на тот момент» значение (курс валют, цена, версия атрибута).
-- Цена SKU на момент транзакции (правая таблица снапшотов: start_time сортируется) SELECT s.tx_datetime, s.sku_id, s.qty, p.price FROM sales s ASOF JOIN price_snapshots p ON p.sku_id = s.sku_id AND p.start_time <= s.tx_datetime ORDER BY s.sku_id, s.tx_datetime;
Ключевые условия успеха:
- таблица справа отсортирована по (sku_id, start_time);
- партиции/ORDER BY позволяют «подсечь» нужные окна.
Риск: ASOF JOIN без хорошего ORDER BY превращается в «полный перебор».
Митигация: держите ORDER BY (dim_id, valid_from); при больших объёмах — материализуйте результат pre-join в широкую витрину.
АргМакс/АргМин как «последнее известное»
Иногда проще использовать argMax:
-- Последнее известное значение region_name до дня X
SELECT shop_id,
argMax(region_name, valid_from) AS region_name_on_day
FROM d_shop_scd2
WHERE valid_from <= toDate('2025-07-31')
GROUP BY shop_id;
Риск: «последнее» не равно «актуальное» при наличии valid_to.
Митигация: храните SCD2 с valid_from/valid_to и по возможности используйте ASOF.
Последовательности, воронки и сессии: sequenceMatch/Count, windowFunnel
Сессии (таймаут 30 минут)
WITH ordered AS (
SELECT user_id, event_time,
lagInFrame(event_time) OVER (PARTITION BY user_id ORDER BY event_time) AS prev_t
FROM events
),
flags AS (
SELECT *, (prev_t IS NULL OR event_time - prev_t > 1800) AS is_new
FROM ordered
),
sess AS (
SELECT *,
sum(is_new) OVER (PARTITION BY user_id ORDER BY event_time) AS sess_num
FROM flags
)
SELECT user_id,
concat(toString(user_id), '_', toString(sess_num)) AS session_id,
min(event_time) AS start_t,
max(event_time) AS end_t,
count() AS events
FROM sess
GROUP BY user_id, sess_num;
Риск: разные TZ/таймауты дают разные сессии.
Митигация: фиксируйте правило в паспорте метрики (30/60 минут, локальный TZ?).
sequenceMatch / sequenceCount
Паттерны в событиях описываются регулярками:
-- Путь: view → add_to_cart → purchase в одной сессии, за 24 часа
SELECT
user_id,
sequenceMatch('(?1).*(?2).*(?3)')(
toUInt32(1) AS view,
toUInt32(event_name='add_to_cart') AS add,
toUInt32(event_name='purchase') AS buy
) AS funnel_pass
FROM events
WHERE event_time >= now() - INTERVAL 1 DAY
GROUP BY user_id;
windowFunnel — сколько дошли за окно
SELECT
toDate(event_time) AS day,
windowFunnel(3600)(
event_time,
event_name='view',
event_name='add_to_cart',
event_name='purchase'
) AS depth
FROM events
WHERE day >= today()-7
GROUP BY day;
Риски:
- «Склейка» пользователей задним числом меняет DAU/воронки.
- События дублируются после рестартов консьюмера.
Митигации:
- Identity-граф + ретро-окно пересчёта (Модуль 7).
- Идемпотентность ingestion (event_id) и дедуп на STAGE.
Гео-аналитика: расстояния, полигоны, геосетки
Фильтр по радиусу (Haversine)
-- Пользователи в радиусе 5 км от точки SELECT user_id FROM users_geo WHERE greatCircleDistance(lat, lon, 55.751244, 37.618423) <= 5000;
Попадание в полигон (зона доставки)
-- polygon: массив пар (lon, lat); порядок важен! WITH poly AS ( SELECT [(37.6,55.7),(37.7,55.7),(37.7,55.8),(37.6,55.8)] AS poly ) SELECT user_id FROM users_geo, poly WHERE pointInPolygon((lon, lat), poly);
Геосетки (geohash) для теплокарт
SELECT geohashEncode(lon, lat, 7) AS cell, count() AS c FROM pings WHERE event_time >= now()-INTERVAL 1 DAY GROUP BY cell ORDER BY c DESC;
Риски:
- Порядок координат (lon, lat) vs (lat, lon) — классический источник ошибок.
- SRID/проекции: используйте WGS-84 и функции greatCircleDistance/pointInPolygon.
- «Край» полигона: попадания на границе могут считаться по-разному.
Митигации: фиксируйте формат (lon,lat), храните гео-календарь событий (смены зон), тестируйте пограничные кейсы.
JSON, массивы и Map: хранить умно, читать быстро
Стратегия хранения
- «Шумный» JSON оставляйте как JSONEachRow на входе, но вытаскивайте нужные поля в отдельные типизированные столбцы (особенно фильтры/джойны).
- Для «ключ-значение» используйте Map(String, String/Int/Float) или Array(T) + функции arrayMap/arrayFilter.
JSONExtract* — точечное чтение
SELECT JSONExtractString(attrs, 'device.os') AS os, JSONExtractInt(attrs, 'screen.height') AS h FROM events WHERE JSONHas(attrs, 'device.os') AND os = 'iOS';
Риск: сканировать гигантскую колонку JSON ради двух полей — тяжело.
Митигация: выносите часто используемое в отдельные колонки; ставьте bloom-индексы на ключи/значения, если нужно искать по подстрокам/IN.
Массивы и arrayJoin
-- Разворачивание массива тегов товара SELECT sku_id, arrayJoin(tags) AS tag, count() AS c FROM sku_catalog GROUP BY sku_id, tag;
Map и «плоский» доступ
-- Map<String, Float64> metrics: {"cpu":0.7,"mem":0.5}
SELECT ts, metrics['cpu'] AS cpu, metrics['mem'] AS mem
FROM host_metrics
WHERE metrics['cpu'] > 0.9;
Риски:
- Чрезмерная вложенность и высококардинальные ключи → компрессия страдает, индексация не помогает.
- arrayJoin легко «взрывает» строки.
Митигации: нормализуйте до разумной глубины; для частых кейсов держите «широкие» таблицы; применяйте лимиты/фильтры до arrayJoin.
Приближённые агрегаты, heavy-hitters и квантили
Уникальные и аудитории: uniq*
- uniqExact(x) — точно, но дорого;
- uniqCombined(x) — быстрый компромисс (ошибка ~0.5–1.5% на практике);
- uniqHLL12(x) — HyperLogLog со слегка бОльшей ошибкой, но очень дешёвый.
Паттерн состояний:
-- Запись SELECT uniqCombinedState(user_id) AS dau_state INTO agg_events_daily_state GROUP BY day; -- Чтение SELECT uniqCombinedMerge(dau_state) AS DAU FROM agg_events_daily_state WHERE day BETWEEN today()-28 AND today()-1;
Top-K (heavy hitters)
-- Top-10 SKU за день SELECT day, topK(10)(sku_id) AS top10 FROM sales WHERE day BETWEEN today()-7 AND today()-1 GROUP BY day;
Для стабильности в витринах храните topKState + topKMerge во вьюхе.
Квантили и «хвосты» (p90/p95/p99)
- quantileTDigest — устойчив к «хвостам», хороший по качеству и скорости.
- Для миллисекунд/времён ответа — можно quantileTiming.
SELECT quantileTDigest(0.9)(latency_ms) AS p90, quantileTDigest(0.99)(latency_ms) AS p99 FROM http_logs WHERE ts >= now()-INTERVAL 1 DAY;
Риски:
- смешивать точные и приближённые методы в одной метрике;
- считать quantile «на лету» по сырым миллиардам рядов.
Митигации: храните состояния в агрегатах (…State) и читайте …Merge; в паспорте метрики укажите метод (TDigest).
Пересечения и аудитории: bitmap и операции над множествами
Накопление «множества пользователей»
-- Ежедневно пополняем bitmap активных пользователей
CREATE TABLE agg_active_daily_bitmap
(
day Date,
users_bm AggregateFunction(groupBitmap, UInt64)
) ENGINE = AggregatingMergeTree
ORDER BY day;
-- Запись состояния
INSERT INTO agg_active_daily_bitmap
SELECT toDate(event_time) AS day,
groupBitmapState(user_id) AS users_bm
FROM events
GROUP BY day;
Пересечения (AND/OR) и размер
-- Активные на прошлой неделе И сделавшие покупку вчера
WITH week_bm AS (
SELECT bitmapOr(users_bm) AS bm
FROM agg_active_daily_bitmap
WHERE day BETWEEN today()-7 AND today()-1
),
yest_bm AS (
SELECT users_bm FROM agg_active_daily_bitmap WHERE day = yesterday()
)
SELECT bitmapCardinality(bitmapAnd((SELECT bm FROM week_bm),
(SELECT users_bm FROM yest_bm))) AS size_intersection;
Риски: bitmap предполагает целочисленные ID и не хранит «дополнительных атрибутов».
Митигации: ID нормализовать в UInt64; для сегментации по атрибутам — вести несколько bitmap-слоёв (по регионам/каналам) или делать фильтрацию перед построением.
«Мини-API» поверх ClickHouse: как отдавать аналитику наружу
Задача: быстро отдать метрики/сегменты/топы в продукт или партнёрам, не ломая безопасность и SLA.
HTTP-интерфейс ClickHouse (базовые паттерны)
- Параметры передавайте явно: param_-переменные и FORMAT JSON/JSONEachRow.
- Не отдавайте сырые таблицы; только vw_* (семантические представления).
- Ограничьте пользователю профиль/квоты (см. Модуль 6).
-- Пример параметров в HTTP:
SELECT day, shop_id, sum(net_sales) s
FROM db_marts.vw_net_sales_daily
WHERE day BETWEEN {from:Date} AND {to:Date}
AND shop_id = {shop:UInt32}
GROUP BY day, shop_id
FORMAT JSON;
Риск: SQL-инъекции, «падает» из-за FULL SCAN.
Митигации: только параметризованные шаблоны; дефолтные периоды/лимиты; профили max_execution_time, max_rows_to_read.
Кеш-слой/гейтвей
- Простой Nginx/API-гейт с кешем на 30–300 секунд для «горячих» эндпоинтов.
- Варнировать по ключам параметров (period, shop, granularity).
- Для «утяжелённых» запросов → вынос в предагрегаты и отдача из них.
Пагинация и «пролистывание»
- Избегайте OFFSET/LIMIT на больших выборках; используйте seek-based пагинацию по (day, id) или last_value и WHERE (key > last_key).
- Для Top-листов — LIMIT BY (Top-K в группах) и хранение topKState.
Риск: «случайный» полный скан.
Митигации: оберните эндпоинты в сервис, который подставляет WHERE по партиции и валидирует параметры.
Паттерны производительности (акценты для продвинутых кейсов)
- ORDER BY = реальные WHERE/группировки: дата → разрез → id (см. Модуль 2).
- PREWHERE на широких таблицах (сначала «узкие» колонки).
- Data-skipping индексы: bloom для IN/LIKE по длинным ключам; set для маленьких доменов.
- Крупные батчи вставок; микробатчи на Kafka/MV; следите за parts.
- Никакого FINAL в продуктивных запросах; обеспечивайте консистентность на записи (Replacing с версией, Aggregating-состояния).
- Проекции — только после доказанного выигрыша.
Кейсы «от и до»
eCom: зона доставки и конверсия по «покрытию»
Вопрос: как влияет зона 30-минутной доставки на CR и AOV?
- Построили полигон «30 мин» по дорожной сети (внешний сервис), сохранили в таблицу delivery_zone(poligon).
- Отметили пользователям in_zone = pointInPolygon((lon,lat), polygon).
- Считали CR/AOV по in_zone=1/0.
Риски: редкие координаты → деанонимизация, неверный порядок (lon/lat).
Митигации: обобщать гео до сетки geohash, хранить (lon,lat), тесты границ.
FinTech: курсы валют и as-of пересчёт
Вопрос: считать NetSales в «базовой валюте на момент операции».
- dict_fx(date, from, to, rate);
- ASOF JOIN по (currency, date) → rate на дату;
- Ночной ретро-пересчёт 30 дней.
Риски: смешение «на дату операции» и «на дату отчёта» в одном отчёте.
Митигации: две отдельные вьюхи; паспорт метрики (какая логика где).
Маркетинг: аудитории «Добавили в корзину» ∧ «Покупали за 30 дней»
Паттерн: bitmap для множества user_id и пересечения.
- Ежедневно строим bitmap «add_to_cart» и «purchase»;
- Пересечение bitmapAnd → размер сегмента и список id (если нужно, bitmapToArray для выгрузки в DMP).
Риски: повторная заливка дня «накрутит» counts без состояний.
Митигации: хранить groupBitmapState, а не «готовые» битмапы; nightly rebuild окна при сомнении.
Риски в «продвинутых» задачах (сводная таблица)
|
Тема |
Риск |
Симптом |
Что делать |
|---|---|---|---|
|
ASOF JOIN |
«Не попали» в нужную версию |
Null/старая цена |
Проверьте сортировку/ORDER BY справа, перекройте pre-join |
|
Роллинги |
«Рваный» календарь |
Плавающие окна |
Плотнить ряд (join на календарь), инвариант плотности |
|
Воронки |
Дубли/склейка id |
Неожиданный рост |
Идемпотентность, ретро-окно, identity-граф |
|
Гео |
lon/lat местами |
«Мимо» зон |
Зафиксируйте формат, тесты на границах |
|
JSON |
«Жрёт» CPU/IO |
Медленные фильтры |
Вынесите ключевые поля, индексы bloom/set |
|
quantile |
«Прыгает» p99 |
Нестабильность |
TDigest, хранить состояния, достаточный объём |
|
bitmap |
Неверные пересечения |
Размер «гуляет» |
Только groupBitmapState/…Merge, rebuild окна |
|
API |
FULL SCAN/инъекция |
Пики латентности/ошибки |
Параметры, лимиты/квоты, обязательный фильтр по партиции |
Шпаргалка внедрения (по шагам)
- Календарь/время: заведите d_calendar, плотните ряды; у каждой метрики — один календарь и TZ.
- ASOF: сложите SCD2/цены/курсы с ORDER BY (id, valid_from); убедитесь, что по ключам «подсечётся».
- Последовательности: оформите правила воронки (сессия/таймаут), sequenceMatch/windowFunnel → вьюхи.
- Гео: определите формат (lon,lat), заведите полигоны/сетки, проверяйте границы.
- Nested: вынесите часто используемое из JSON; применяйте Map/массивы там, где выгодно.
- Approx: тяжёлые метрики (uniq/quantile/topK) — через состояния (…State/…Merge).
- Аудитории: храните bitmap-состояния по окнам; делайте пересечения на агрегатах.
- Сервисный слой: отдавайте только vw_*, параметризованные запросы; лимиты/квоты; кеш-гейт.
Итог
Продвинутая аналитика на ClickHouse — это не «экзотика», а набор повторяемых паттернов:
- Время: плотные ряды, роллинги, ASOF для «значений на момент».
- Последовательности: sequenceMatch/Count и windowFunnel на нормализованных событиях и сессиях.
- Гео: расстояния/полигоны/геосетки с аккуратной геометрией.
- Nested: JSON/массивы/Map — только там, где оправдано; ключи наружу.
- Приближёнка: uniq/topK/quantile/bitmap — как состояния и с обозначенной точностью.
- Аудитории и пересечения: bitmap-операции вместо «DISTINCT на лету».
- Мини-API: параметризованные запросы к vw_*, лимиты, кеш.
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. Итог — быстрый запуск витрин за недели, снижённые риски в проде и предсказуемая стоимость владения.



