Learn
ClickHouse/20-project-analytics

项目实战:实时行为分析与报表

前面 19 章把零件都讲完了。这一章把它们组装成一条真实可用的数据管道:埋点事件从 Kafka 实时流入 ClickHouse,通过物化视图预聚合,前端报表直接查聚合结果,秒级响应 PV/UV、留存、漏斗。这是 ClickHouse 最经典的生产场景。

1. 业务目标与架构

假设我们在做一个网站分析产品(类 Google Analytics),需要实时回答:

  • 今天各页面的 PV / UV 是多少?
  • 新用户 1/3/7 日 留存率 如何?
  • 从「浏览商品 → 加购 → 下单 → 支付」的 转化漏斗 在哪一步流失最多?
  • 各国家、各设备的流量分布?

整体架构:

   埋点 SDK                       Kafka                       ClickHouse
┌──────────┐   HTTP    ┌──────────────┐    ┌───────────────────────────────┐
│  Web/App │ ───────►  │ topic:       │ ──►│ events_raw (Kafka 引擎,接入)  │
│  客户端  │   JSON    │  user_events │    │        │                       │
└──────────┘           ┌──────────────┐    │        │ INSERT 触发            │
                       │ partition 0  │    │        ▼                       │
                       │ partition 1  │    │ 物化视图: 预聚合到              │
                       │ partition 2  │    │   - pv_uv_1m (按分钟)          │
                       └──────────────┘    │   - funnel_daily               │
                                           │   - retention_daily             │
                                           │        │                       │
                                           │        ▼                       │
                                           │   报表查询 (秒级,前端直查)      │
                                           └───────────────────────────────┘

设计要点:原始 events_raw 只作"管道入口",不参与查询;所有报表读的是物化视图写入的预聚合表。这正是第 9 章讲的"ETL 模式"和第 8 章的预聚合思想。

2. 原始表与 Kafka 接入

-- 1) 原始事件表:用 Kafka 引擎直接对接 topic
CREATE TABLE events_raw
(
    event_time  DateTime,
    user_id     UInt64,
    event_type  String,        -- 'page_view' / 'add_cart' / 'order' / 'pay'
    page        String,
    country     String,
    device      String,
    duration_ms UInt32
)
ENGINE = Kafka
SETTINGS
    kafka_broker_list = 'kafka-broker1:9092,kafka-broker2:9092',
    kafka_topic_list  = 'user_events',
    kafka_group_name  = 'ch-consumer',
    kafka_format      = 'JSONEachRow',
    kafka_num_consumers = 4;   -- 每个表实例起 4 个消费者并行拉取
-- 2) 落地的本地明细表(真正存全量原始数据,用于回溯/即席分析)
CREATE TABLE events
(
    event_time  DateTime,
    user_id     UInt64,
    event_type  LowCardinality(String),
    page        String,
    country     LowCardinality(String),
    device      LowCardinality(String),
    duration_ms UInt32
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_time)
ORDER BY (event_type, country, event_time);   -- 排序键贴合"按事件类型+地区"的分析
ℹ️Kafka 表只是「管道」,不存数据

events_raw 用 Kafka 引擎,它不持久化数据,只是消费位点 + 流式读取的接口。数据要靠物化视图"搬运"到 events 这类真存储表里。可以起多个物化视图从同一个 events_raw 扇出到不同目标(9 章的 Null 扇出思路)。用 SELECT count() FROM events_raw 永远返回 0 或近似,这是正常的。

3. 用物化视图预聚合(核心)

我们建三张聚合表,分别服务三类报表。全部用第 8 章的 AggregatingMergeTree,因为 .State 组合子能把"中间聚合状态"存下来,后台 merge 时自动合并。

3.1 PV/UV 按分钟聚合

-- 状态表:用 AggregateFunction 存 UV 的去重状态,而非直接存数字
CREATE TABLE pv_uv_1m
(
    minute    DateTime,
    page      String,
    country   LowCardinality(String),
    device    LowCardinality(String),
    pv        AggregateFunction(count),
    uv        AggregateFunction(uniq, UInt64)   -- 存去重"状态",不是去重结果
)
ENGINE = AggregatingMergeTree
PARTITION BY toYYYYMM(minute)
ORDER BY (page, country, device, minute);
 
-- 搬运视图:INSERT 进 events_raw 时自动聚合并写入 pv_uv_1m
CREATE MATERIALIZED VIEW mv_pv_uv
TO pv_uv_1m
AS SELECT
    toStartOfMinute(event_time) AS minute,
    page,
    country,
    device,
    countState()                AS pv,        -- 注意是 -State
    uniqState(user_id)          AS uv         -- 去重状态
FROM events_raw
GROUP BY minute, page, country, device;

查询时要用 -Merge 把状态"折叠"成结果:

-- 报表查询:秒级返回,因为 mv 已经把全量预聚合好了
SELECT
    minute,
    page,
    country,
    sumMerge(pv)  AS pv,        -- -Merge 合并状态
    uniqMerge(uv) AS uv
FROM pv_uv_1m
WHERE minute >= now() - INTERVAL 1 HOUR
GROUP BY minute, page, country
ORDER BY minute DESC;
没有物化视图时:扫全量 events(可能几十亿行) → 秒~分钟级
有物化视图时:扫 pv_uv_1m(按分钟聚合,行数缩减几千倍) → 毫秒~秒级

3.2 每日留存

留存需要"按用户记录首次/当日活跃",用位图最优雅(12 章的 groupBitmap):

CREATE TABLE retention_daily
(
    day        Date,
    country    LowCardinality(String),
    -- 当天活跃用户位图 与 各日后仍活跃的用户位图
    d0         AggregateFunction(groupBitmap, UInt64),
    d1         AggregateFunction(groupBitmap, UInt64),
    d3         AggregateFunction(groupBitmap, UInt64),
    d7         AggregateFunction(groupBitmap, UInt64)
)
ENGINE = AggregatingMergeTree
PARTITION BY toYYYYMM(day)
ORDER BY (country, day);
 
CREATE MATERIALIZED VIEW mv_retention
TO retention_daily
AS SELECT
    toDate(event_time)                       AS day,
    country,
    groupBitmapState(user_id)               AS d0,   -- 当天活跃
    groupBitmapStateIf(user_id,
        toDate(event_time) = toDate('1970-01-01')) AS d1,  -- 占位,真实留存需跨日关联,见下文
    groupBitmapState(user_id)               AS d3,
    groupBitmapState(user_id)               AS d7
FROM events_raw
GROUP BY day, country;
⚠️上例的 d1/d3/d7 是「当日活跃」的简化占位

真正留存 = "第 N 天活跃的用户 ∩ 第 0 天活跃的用户",需要跨日关联,单条 MV 无法直接表达。生产做法:MV 只按天存各日活跃位图 day_active(如 retention_bitmap(day, country, users Bitmap)),查询时 bitmapAnd(card(d0), card(d1)) / card(d0) 计算留存率。这里为讲清思路做了简化,务必理解"留存=位图交集"的本质。

查询留存率(正确写法,基于每日活跃位图表):

-- 假设有表 retention_bitmap(day, country, active AggregateFunction(groupBitmap,UInt64))
SELECT
    country,
    day,
    bitmapCardinality(bitmapMerge(active))                                   AS d0_cnt,
    bitmapCardinality(bitmapAnd(
        bitmapMerge(active),                                                -- 当天
        (SELECT bitmapMerge(active) FROM retention_bitmap rb
         WHERE rb.day = retention_bitmap.day + 1 AND rb.country = retention_bitmap.country)  -- 次日
    )) / bitmapCardinality(bitmapMerge(active))                              AS d1_retention
FROM retention_bitmap
GROUP BY country, day
ORDER BY day DESC;

3.3 转化漏斗

漏斗用第 14 章的 windowFunnel。预聚合阶段把每个用户的行为序列按时间排好,存进聚合表:

CREATE TABLE funnel_daily
(
    day        Date,
    country    LowCardinality(String),
    -- 用 AggregateFunction 存每个用户的"漏斗命中步骤数"分布
    steps      AggregateFunction(countIf, UInt8)
)
ENGINE = AggregatingMergeTree
PARTITION BY toYYYYMM(day)
ORDER BY (country, day);
 
CREATE MATERIALIZED VIEW mv_funnel
TO funnel_daily
AS SELECT
    toDate(event_time) AS day,
    country,
    -- 每个用户在当天漏斗里走到第几步
    countIfState(
        windowFunnel(86400)(
            event_time,
            event_type = 'page_view',
            event_type = 'add_cart',
            event_type = 'order',
            event_type = 'pay'
        ) = 4   -- 走完 4 步才算"完成支付"
    ) AS steps
FROM events_raw
GROUP BY day, country, user_id;   -- 注意按 user_id 分组,漏斗是"每人一条"
💡漏斗一定要按 user_id 分组

windowFunnel 是对"单个用户的行为序列"求漏斗。如果 GROUP BY 里没有 user_id,就会把不同用户的行为混在一起算,结果完全失真。MV 的 GROUP BY day, country, user_id 保证"每人一条漏斗结果",再往上汇总。

4. 级联物化视图:更粗的粒度

第 9 章讲过级联 MV。我们可以让 pv_uv_1m(分钟级)再被一个 pv_uv_1h(小时级)的 MV 消费,形成"明细 → 分钟 → 小时"的层级,越往上越快、越省。

CREATE TABLE pv_uv_1h
(
    hour    DateTime,
    page    String,
    country LowCardinality(String),
    device  LowCardinality(String),
    pv      AggregateFunction(count),
    uv      AggregateFunction(uniq, UInt64)
)
ENGINE = AggregatingMergeTree
PARTITION BY toYYYYMM(hour)
ORDER BY (page, country, device, hour);
 
-- 消费"分钟表"的插入,向上滚成小时级
CREATE MATERIALIZED VIEW mv_pv_uv_hour
TO pv_uv_1h
AS SELECT
    toStartOfHour(minute) AS hour,
    page, country, device,
    pv, uv                       -- 从分钟表直接透传 AggregateFunction 状态
FROM pv_uv_1m;

因为 pv/uv 是聚合状态,直接透传即可——ClickHouse 在后台 merge pv_uv_1h 时会自动把多个分钟状态 -Merge 成小时状态。这就是"状态可组合"的威力。

5. 和"不用 ClickHouse"的对比

场景:实时看板,每次打开要查"最近 24 小时 PV/UV 按页面分组"
 
MySQL 方案:
  原始表 50 亿行,每次 SELECT COUNT(DISTINCT user_id) ... GROUP BY page
  → 全表扫描 + 去重排序,几十秒~分钟级,看板直接超时
 
ClickHouse 方案:
  mv 早已把数据按分钟预聚合进 pv_uv_1m(行数约为原始的 1/5000)
  报表查 sumMerge/uniqMerge,毫秒~秒级返回
  → 同样的看板,从"打不开"变成"实时刷新"

这正是 ClickHouse 在 OLAP 场景碾压 MySQL 的根因:把计算在写入时提前做掉(预聚合),查询时只做轻量合并。

⚠️预聚合会「丢失明细灵活性」

预聚合表的维度是固定的(上例按 page/country/device 聚合),没法再回答"按 city 维度的 PV"这种没预聚合过的问题。取舍:高频固定维度 → 预聚合(快);偶发探索性查询 → 直接查 events 明细表(慢但灵活)。两类表并存,正是生产标准架构。

6. 生产加固清单

把这条管道推向生产,还需补几件事(对应前面章节):

  • 去重与更新:events_raw 可能重复投递。用 7 章 ReplacingMergeTree 或在 MV 里 uniqState 天然去重。
  • 跳数索引:events 明细表如果常按 user_id 查,加 16 章的 bloom_filter 跳数索引。
  • TTL:events 明细表设 3 个月 TTL 自动清理(6 章),聚合表可保留更久。
  • 副本与分片:events 与聚合表都用 18 章 ReplicatedMergeTree + 17 章 Distributed,保证高可用与扩容。
  • 监控:19 章的 system.kafka_consumers 看消费滞后,system.merges 看 MV 写入压力。
-- 检查 Kafka 消费是否落后(Lag 持续增大 = 处理能力不足)
SELECT
    topic, consumer_group, num_consumed, num_messages,
    num_messages - num_consumed AS lag
FROM system.kafka_consumers;

7. 小结

  • 经典管道:Kafka(Kafka 引擎接入)→ 物化视图预聚合(AggregatingMergeTree + -State)→ 报表查 -Merge。
  • 原始 events_raw 只作管道,不存数据、不参与查询;明细 events 保留全量供探索。
  • PV/UV 用 countState/uniqState;留存用 groupBitmap 位图交集;漏斗用 windowFunnel 按 user_id 分组。
  • 级联 MV 形成"明细 → 分钟 → 小时"层级,聚合状态可逐层 -Merge 组合。
  • 预聚合换来了极速,代价是丢失未预聚合维度的灵活性——两者并存是标准架构。
  • 上线前用副本/分片、跳数索引、TTL、Kafka Lag 监控把管道加固。
🎯结课综合练习
  1. 用 Docker 起一个单节点 Kafka,写一个小脚本持续往 user_events topic 灌 JSON 事件(混合 4 种 event_type、随机 user_id)。
  2. 在 ClickHouse 建 events_raw(Kafka 引擎)、events(MergeTree)、pv_uv_1m(AggregatingMergeTree) 与 mv_pv_uv,跑 5 分钟后查 pv_uv_1m 的 PV/UV。
  3. 在 pv_uv_1m 之上建 pv_uv_1h 级联 MV,验证小时级聚合结果等于分钟级之和(sumMerge(pv) 对比)。
  4. 用 system.kafka_consumers 确认消费无滞后;故意停掉灌数据脚本,观察 MV 是否停止更新(验证 MV 是 INSERT 触发器语义,9 章)。
  5. bonus:用本章漏斗 MV 思路,算出「page_view → add_cart → order → pay」的整体转化率,定位流失最大的环节。