这是《ClickHouse 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 本章把前面内容串成一个实时用户行为分析平台,覆盖埋点接入、实时写入、明细建模、汇总报表、漏斗留存和运维治理。
30.1 业务目标
平台需要回答:
1. 今天各城市、渠道的 PV / UV / GMV 是多少?
2. 用户从浏览到支付的漏斗转化如何?
3. 某个用户最近发生了什么事件?
4. 昨天活动上线后指标是否异常?
5. 数据延迟和质量是否达标?
指标口径:
| 指标 | 定义 |
|---|---|
| PV | event_type = view 的次数 |
| UV | event_type = view 的去重用户数 |
| GMV | event_type = pay 的 sum(amount) |
| 支付率 | 支付用户数 / 访问用户数 |
| D1 留存 | 今日活跃用户在次日仍活跃 |
30.2 链路架构
Web / App
|
Collector
|
Kafka raw-events
|
Flink 清洗 / 补维 / 质量检查
|
Kafka clean-events
|
ClickHouse Kafka engine
|
dwd_events
|
dws_city_daily + ads_realtime
死信链路:
invalid events -> Kafka dead-letter-events -> 修复服务
30.3 明细表
CREATE TABLE analytics.dwd_events
(
event_date Date,
event_time DateTime,
event_id String,
user_id UInt64,
session_id String,
event_type LowCardinality(String),
platform LowCardinality(String),
city_id UInt32,
channel LowCardinality(String),
amount Decimal64(2),
is_valid UInt8
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(event_date)
ORDER BY (event_date, event_type, city_id, channel, user_id)
TTL event_date + INTERVAL 13 MONTH;
设计说明:
event_date来自事件时间;- 高频过滤字段均普通列化;
- 低基数字段使用 LowCardinality;
is_valid支持坏数据标记;- 保留 13 个月,支持跨年对比。
30.4 汇总表
CREATE TABLE analytics.dws_city_daily
(
stat_date Date,
city_id UInt32,
channel LowCardinality(String),
pv UInt64,
pay_count UInt64,
gmv Decimal64(2)
)
ENGINE = SummingMergeTree
PARTITION BY toYYYYMM(stat_date)
ORDER BY (stat_date, city_id, channel);
物化视图:
CREATE MATERIALIZED VIEW analytics.dws_city_daily_mv
TO analytics.dws_city_daily
AS
SELECT
event_date AS stat_date,
city_id,
channel,
countIf(event_type = 'view') AS pv,
countIf(event_type = 'pay') AS pay_count,
sumIf(amount, event_type = 'pay') AS gmv
FROM analytics.dwd_events
WHERE is_valid = 1
GROUP BY event_date, city_id, channel;
UV 需要精确口径时使用 uniqState(user_id) 与 AggregatingMergeTree。
30.5 实时报表
SELECT
city_id,
channel,
sum(pv) AS pv,
sum(pay_count) AS pay_count,
sum(gmv) AS gmv
FROM analytics.dws_city_daily
WHERE stat_date = today()
GROUP BY city_id, channel
ORDER BY gmv DESC
LIMIT 50;
实时分钟表可按分钟汇总,但保留期短:
minute metrics:保留 3 天
daily metrics:保留 13 个月
raw detail:保留 6 到 13 个月
30.6 漏斗分析
SELECT
level,
count() AS users
FROM (
SELECT
user_id,
windowFunnel(86400)(
toUnixTimestamp(event_time),
event_type = 'view',
event_type = 'cart',
event_type = 'pay'
) AS level
FROM analytics.dwd_events
WHERE event_date = today()
GROUP BY user_id
)
GROUP BY level
ORDER BY level;
注意窗口起点、事件顺序、跨天用户和多设备归属。
30.7 用户查询
SELECT
event_time,
event_type,
city_id,
channel,
amount
FROM analytics.dwd_events
WHERE event_date IN (today(), today() - 1)
AND user_id = 1001
ORDER BY event_time DESC
LIMIT 100;
主排序对 user_id 查询不友好,可以增加投影:
ALTER TABLE analytics.dwd_events
ADD PROJECTION proj_user
(
SELECT *
ORDER BY (user_id, event_date, event_time)
);
上线前对比写入吞吐、磁盘增长和 read_bytes。
30.8 权限
CREATE ROLE analyst_ro;
GRANT SELECT ON analytics.dwd_events TO analyst_ro;
GRANT SELECT ON analytics.dws_city_daily TO analyst_ro;
CREATE USER analyst IDENTIFIED BY 'strong-password'
SETTINGS PROFILE analyst_profile;
GRANT analyst_ro TO analyst;
分析用户只读,ETL 用户按库表授权,管理员与业务账号分离。
30.9 监控
| 指标 | 目标 |
|---|---|
| Kafka lag | 小于 1 分钟 |
| 写入失败 | 立即告警 |
| 核心表行数 | 环比异常告警 |
| 查询 P95 | 按报表分级 |
| parts | 不持续增长 |
| 磁盘 | 70% 预警 |
| 汇总对账 | 每小时执行 |
对账:
SELECT countIf(event_type = 'pay'), sumIf(amount, event_type = 'pay')
FROM analytics.dwd_events
WHERE event_date = today();
SELECT sum(pay_count), sum(gmv)
FROM analytics.dws_city_daily
WHERE stat_date = today();
30.10 上线清单
1. schema review
2. 分区和排序键 review
3. 写入压测
4. 查询压测
5. 权限最小化
6. 备份恢复演练
7. 监控和 runbook
8. 数据质量对账
9. 灰度切流
10. 回滚方案
本章小结
实战平台的关键不是“能写入能查询”,而是口径清楚、链路可恢复、资源可预算、数据可对账。明细表保留可追溯性,汇总表服务高频报表,投影和字典解决特定查询,监控和权限保障长期稳定。
思考题
- 为什么 UV 不能直接在 SummingMergeTree 中相加?
- 漏斗窗口如何定义?
- 用户明细查询为什么需要投影?
- 实时报表的数据质量如何验证?
- 上线前必须压测哪些场景?