项目实战:实时行为分析与报表
前面 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); -- 排序键贴合"按事件类型+地区"的分析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;真正留存 = "第 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 分组,漏斗是"每人一条"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 监控把管道加固。
- 用 Docker 起一个单节点 Kafka,写一个小脚本持续往
user_eventstopic 灌 JSON 事件(混合 4 种 event_type、随机 user_id)。 - 在 ClickHouse 建
events_raw(Kafka 引擎)、events(MergeTree)、pv_uv_1m(AggregatingMergeTree) 与mv_pv_uv,跑 5 分钟后查pv_uv_1m的 PV/UV。 - 在
pv_uv_1m之上建pv_uv_1h级联 MV,验证小时级聚合结果等于分钟级之和(sumMerge(pv)对比)。 - 用
system.kafka_consumers确认消费无滞后;故意停掉灌数据脚本,观察 MV 是否停止更新(验证 MV 是 INSERT 触发器语义,9 章)。 - bonus:用本章漏斗 MV 思路,算出「page_view → add_cart → order → pay」的整体转化率,定位流失最大的环节。