第 19 章 · 实时埋点分析平台

综合实战:把前 18 章的 ClickHouse 能力串成一个 OLAP 大屏项目

全链路数据流

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_kafkaKafka消费 topic
明细 ODSevents_rawMergeTree + Projection全量明细,按天分区,TTL 90 天
聚合 DWSevents_agg_pv_uvAggregatingMergeTreePV/UV(uniqExactState)
聚合 DWSevents_agg_funnelAggregatingMergeTreewindowFunnelState
用户首日user_cohortReplacingMergeTree留存分析基准
维度dict_productsDictionary替代大表 JOIN

关键建表 DDL 速览

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);

实时 KPI

今日 PV
0
+0 / 秒
今日 UV
0
+0 / 秒
在线人数
0
实时
GMV (¥)
0
+¥0

Top 10 商品(实时滚动)

背后的查询

-- 每 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;

漏斗转化可调演示

拖动各步转化率,观察漏斗尾部的总人数变化。

200,000
35%
45%
80%

背后的 SQL

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;

留存矩阵(14 天 × 14 offset 热力图)

≤10% ≤30% ≤50% ≤70% >70%

SQL

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;