SDK → 网关 → Kafka → ClickHouse Kafka 引擎 → MV 链 → Agg 表 → BI。
flowchart LR
SDK["埋点 SDK"] --> GW["接入网关"]
GW --> K["Kafka topic events"]
K --> CKK["events_kafka (Kafka 引擎)"]
CKK --> RAW["events_raw (MergeTree)"]
RAW --> AGG1["events_agg_pv_uv"]
RAW --> AGG2["events_agg_funnel"]
RAW --> AGG3["user_cohort (Replacing)"]
D["dict_products"] -.字典.-> RAW
AGG1 --> BI["大屏 / Superset"]
AGG2 --> BI
AGG3 --> BI
| 层次 | 表 | 引擎 | 作用 |
|---|---|---|---|
| 接入 | events_kafka | Kafka | 消费 topic |
| 明细 ODS | events_raw | MergeTree + Projection | 全量明细,按天分区,TTL 90 天 |
| 聚合 DWS | events_agg_pv_uv | AggregatingMergeTree | PV/UV(uniqExactState) |
| 聚合 DWS | events_agg_funnel | AggregatingMergeTree | windowFunnelState |
| 用户首日 | user_cohort | ReplacingMergeTree | 留存分析基准 |
| 维度 | dict_products | Dictionary | 替代大表 JOIN |
CREATE TABLE events_raw (
event_date Date DEFAULT toDate(event_time),
event_time DateTime64(3, 'Asia/Shanghai'),
user_id UInt64,
event_type LowCardinality(String),
product_id UInt32,
channel LowCardinality(String),
...
) ENGINE = MergeTree
PARTITION BY toYYYYMMDD(event_date)
ORDER BY (event_date, event_type, user_id, event_time);
-- 每 2 秒刷新
SELECT dictGet('learn_ck.dict_products','name',product_id) AS name,
sumMerge(pv) AS pv,
uniqExactMerge(uv) AS uv
FROM events_agg_pv_uv
WHERE event_date = today() AND event_type = 'view'
GROUP BY product_id
ORDER BY pv DESC LIMIT 10;
拖动各步转化率,观察漏斗尾部的总人数变化。
WITH t AS (
SELECT user_id,
windowFunnelMerge(3600)(funnel_state) AS step
FROM events_agg_funnel
WHERE event_date = today()
GROUP BY user_id
)
SELECT countIf(step>=1) AS view,
countIf(step>=2) AS cart,
countIf(step>=3) AS "order",
countIf(step>=4) AS pay
FROM t;
SELECT c.first_date AS cohort,
toInt32(e.event_date - c.first_date) AS offset,
uniqExact(c.user_id) AS retained
FROM user_cohort c
JOIN events_raw e ON c.user_id = e.user_id
WHERE c.first_date >= today() - 14
GROUP BY cohort, offset
ORDER BY cohort, offset;