ClickHouseNotes

第 30 章:项目实战

zjc 于 2026-01-30 发布

这是《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;

设计说明:

  1. event_date 来自事件时间;
  2. 高频过滤字段均普通列化;
  3. 低基数字段使用 LowCardinality;
  4. is_valid 支持坏数据标记;
  5. 保留 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. 回滚方案

本章小结

实战平台的关键不是“能写入能查询”,而是口径清楚、链路可恢复、资源可预算、数据可对账。明细表保留可追溯性,汇总表服务高频报表,投影和字典解决特定查询,监控和权限保障长期稳定。

思考题

  1. 为什么 UV 不能直接在 SummingMergeTree 中相加?
  2. 漏斗窗口如何定义?
  3. 用户明细查询为什么需要投影?
  4. 实时报表的数据质量如何验证?
  5. 上线前必须压测哪些场景?