主题
第 19 章 综合实战:实时埋点分析平台
学习目标:把前 18 章的知识点 —— 列存、向量化、MergeTree 家族、物化视图、Projection、TTL、Mutation、副本分布式、Kafka 生态、性能调优 —— 全部串到一个真实项目里:电商网站的实时埋点分析平台。学完本章,你应该能独立在白板上画出「从 SDK 上报到大屏渲染」的完整链路,背后每一段引擎/表结构为什么这样选、踩过哪些坑、未来怎么扩容,都能讲清楚。
19.0 一句话总览
「这一章不是再教一个新语法,而是把前 18 章的乐高积木搭成一座『秒级响应的实时数仓』:Kafka 接入 → Kafka 引擎表 → MV 链 → AggregatingMergeTree → BI 大屏,全链路 P99 延迟 < 5 秒,扛住单机 5 万行/秒写入、亚秒级即席查询。」
19.1 业务需求与技术指标
19.1.1 业务场景
某电商网站日活千万级,每个用户从进站到下单平均产生 30~50 条埋点事件(曝光、点击、加购、下单、支付、退款……)。业务方要求:
| # | 业务需求 | 关键能力 |
|---|---|---|
| ① | 实时大屏 | PV / UV / 在线人数 / Top10 商品,≤ 3 秒刷新 |
| ② | 漏斗分析 | 浏览 → 加购 → 下单 → 支付,按自定义时间窗查转化率 |
| ③ | 留存分析 | 日留存 / 周留存 / 月留存矩阵 |
| ④ | 用户行为路径 | 任意用户的事件时间线,任意事件 N 步前/后的分布 |
| ⑤ | 即席查询(OLAP Cube) | 国家 × 渠道 × 商品分类 × 时间多维自由下钻 |
技术侧硬指标:
- 写入:单机 5 万行/秒稳定,峰值 10 万行/秒不丢数
- 查询:大屏类 P99 < 1 秒、漏斗/留存 P99 < 5 秒、即席 P99 < 30 秒
- 存储:明细保留 30 天,聚合保留 2 年;30 天后明细自动下沉冷盘,90 天删除
- 成本:单节点(32C / 128G / 4T NVMe)扛 3 亿行/天的写入
19.1.2 为什么选 ClickHouse 而不是 …?
| 方案 | 能不能干 | 为什么不选 |
|---|---|---|
| MySQL + 离线 T+1 | 不能 | 无法秒级响应;亿级聚合 10 分钟不出结果 |
| Elasticsearch | 部分能 | 明细检索强,但 Cube / 漏斗 / 状态聚合弱,成本高 |
| Druid | 能 | 聚合强,但 SQL 方言窄、JOIN 弱、运维复杂 |
| Doris / StarRocks | 能 | 同赛道竞品,生态略晚;团队已经在用 CK,不换 |
| ClickHouse | ✅ | 列存向量化写入快、SQL 完整、MV 实时聚合、成本低 |
📌 与传统方案一句话对比:MySQL 是「OLTP 下班回家」,ClickHouse 是「OLAP 永不下班」。想要秒级分析 + 亿级数据 + SQL 友好,当前开源生态里 ClickHouse 依然是首选。
19.2 架构总览(全链路)
分层职责:
| 层 | 作用 | 引擎 / 工具 |
|---|---|---|
| 采集层 | SDK → 网关 → Kafka | Nginx+Lua / Kafka |
| 接入层 | ClickHouse 消费 Kafka | Kafka 引擎 + MV |
| 明细层(ODS) | 全量明细,保留 30 天 | MergeTree + Projection + TTL |
| 聚合层(DWS) | 各业务口径聚合状态 | AggregatingMergeTree + 链式 MV |
| 维度层 | 商品 / 用户 / 渠道维度 | Dictionary(源:MySQL/HTTP) |
| 查询层 | 大屏 / BI / 即席 | HTTP + xxxMerge / uniqMerge |
19.3 完整建表 DDL
完整可执行脚本在
19_project/init.sql,下文按层级逐段讲「为什么这么建」。
19.3.1 建库与设置
sql
CREATE DATABASE IF NOT EXISTS learn_ck;
-- 本章默认所有表都在 learn_ck 库下
SET allow_experimental_projection_optimization = 1;19.3.2 明细表 events_raw —— 整栋楼的地基
sql
CREATE TABLE learn_ck.events_raw
(
event_date Date DEFAULT toDate(event_time),
event_time DateTime64(3, 'Asia/Shanghai'),
user_id UInt64,
session_id String,
event_type LowCardinality(String), -- view / add_cart / order / pay
page_url String CODEC(ZSTD(3)),
referrer String CODEC(ZSTD(3)),
product_id UInt32,
price Decimal(12, 2),
qty UInt16 DEFAULT 1,
channel LowCardinality(String), -- app / h5 / mini / pc
country LowCardinality(String),
city LowCardinality(String),
ua String CODEC(ZSTD(3)),
ip IPv4,
ext Map(String, String), -- 预留扩展字段
ingest_time DateTime DEFAULT now() -- 入库时间,排查延迟用
)
ENGINE = MergeTree
PARTITION BY toYYYYMMDD(event_date)
ORDER BY (event_date, event_type, user_id, event_time)
PRIMARY KEY (event_date, event_type, user_id)
SETTINGS index_granularity = 8192;关键设计解释:
- 分区按天(
toYYYYMMDD):每天一个分区,TTL / 冷热分层都按分区搬; - 排序键
(event_date, event_type, user_id, event_time):- 大屏类按
event_type+ 日期过滤,命中稀疏索引; - 用户路径查询按
user_id,命中稀疏索引中段;
- 大屏类按
- 主键只取排序键前 3 列:稀疏主键保留在内存,省内存;
LowCardinality用在基数 ≤ 1 万的枚举列(事件类型、渠道、国家);String CODEC(ZSTD(3))给文本类(URL、UA、Referrer)单列压缩,节省 60% 空间;Map(String, String)预留扩展字段,业务早期不用改表结构;ingest_time用于排查「埋点时间 ≠ 入库时间」的延迟问题。
19.3.3 跳数索引与 Projection
sql
-- Skip Index:按 product_id 查明细
ALTER TABLE learn_ck.events_raw
ADD INDEX idx_product product_id TYPE minmax GRANULARITY 4;
-- Projection:大屏场景预物化一份「按小时 × 商品」的列
ALTER TABLE learn_ck.events_raw
ADD PROJECTION proj_hour_product
(
SELECT
toStartOfHour(event_time) AS hour,
product_id,
count(),
uniqExact(user_id)
GROUP BY hour, product_id
);
ALTER TABLE learn_ck.events_raw MATERIALIZE PROJECTION proj_hour_product;📌 Projection 与 MV 如何选?查询路径不变用 Projection(自动选择),查询路径变(例如要 uniq 状态再 Merge 给跨天 UV)用 MV。两者可以共存。
19.3.4 Kafka 引擎表 + MV 链路(ODS → ODS)
sql
-- Kafka 引擎表(实际环境把 broker / topic 替换)
CREATE TABLE learn_ck.events_kafka
(
event_time DateTime64(3),
user_id UInt64,
session_id String,
event_type String,
page_url String,
referrer String,
product_id UInt32,
price Decimal(12, 2),
qty UInt16,
channel String,
country String,
city String,
ua String,
ip IPv4,
ext Map(String, String)
)
ENGINE = Kafka
SETTINGS
kafka_broker_list = '127.0.0.1:9092', -- TODO: 替换为生产集群
kafka_topic_list = 'events',
kafka_group_name = 'ck_ingest_events',
kafka_format = 'JSONEachRow',
kafka_num_consumers = 4,
kafka_max_block_size = 65536,
kafka_skip_broken_messages = 10;
-- Kafka → events_raw 的 MV:做字段清洗 + 时间转换
CREATE MATERIALIZED VIEW learn_ck.mv_kafka_to_raw
TO learn_ck.events_raw AS
SELECT
toDate(event_time) AS event_date,
event_time,
user_id,
session_id,
event_type,
page_url,
referrer,
product_id,
price,
qty,
channel,
country,
city,
ua,
ip,
ext,
now() AS ingest_time
FROM learn_ck.events_kafka;19.3.5 PV / UV 聚合表(DWS 层)
sql
CREATE TABLE learn_ck.events_agg_pv_uv
(
event_date Date,
hour UInt8,
event_type LowCardinality(String),
channel LowCardinality(String),
country LowCardinality(String),
product_id UInt32,
pv AggregateFunction(sum, UInt64),
uv AggregateFunction(uniqExact, UInt64)
)
ENGINE = AggregatingMergeTree
PARTITION BY toYYYYMM(event_date)
ORDER BY (event_date, hour, event_type, channel, country, product_id);
CREATE MATERIALIZED VIEW learn_ck.mv_raw_to_pv_uv
TO learn_ck.events_agg_pv_uv AS
SELECT
event_date,
toHour(event_time) AS hour,
event_type,
channel,
country,
product_id,
sumState(toUInt64(1)) AS pv,
uniqExactState(user_id) AS uv
FROM learn_ck.events_raw
GROUP BY event_date, hour, event_type, channel, country, product_id;💡
uniqExactState在 Merge 时能精确去重;数据量再大换uniqState(HyperLogLog,约 0.5% 误差)换 10x 空间节省。
19.3.6 漏斗聚合表
sql
CREATE TABLE learn_ck.events_agg_funnel
(
event_date Date,
channel LowCardinality(String),
funnel_state AggregateFunction(
windowFunnel(3600),
DateTime64(3), UInt8, UInt8, UInt8, UInt8
)
)
ENGINE = AggregatingMergeTree
PARTITION BY toYYYYMM(event_date)
ORDER BY (event_date, channel);
CREATE MATERIALIZED VIEW learn_ck.mv_raw_to_funnel
TO learn_ck.events_agg_funnel AS
SELECT
event_date,
channel,
windowFunnelState(3600)(
event_time,
event_type = 'view',
event_type = 'add_cart',
event_type = 'order',
event_type = 'pay'
) AS funnel_state
FROM learn_ck.events_raw
GROUP BY event_date, channel, user_id; -- 按用户聚合漏斗状态19.3.7 留存聚合表
sql
CREATE TABLE learn_ck.events_agg_retention
(
cohort_date Date,
day_offset UInt16, -- 0 / 1 / 2 / 7 / 30
channel LowCardinality(String),
retained_users AggregateFunction(uniqExact, UInt64)
)
ENGINE = AggregatingMergeTree
PARTITION BY toYYYYMM(cohort_date)
ORDER BY (cohort_date, day_offset, channel);
-- 留存的口径相对复杂,这里用「首日」作为基准:先把每个用户的首日落到 user_cohort
-- 生产上常见的写法是把 user_first_day 维护在一张 ReplacingMergeTree 里
CREATE TABLE learn_ck.user_cohort
(
user_id UInt64,
first_date Date,
channel LowCardinality(String),
created_at DateTime DEFAULT now()
)
ENGINE = ReplacingMergeTree(created_at)
ORDER BY user_id;
CREATE MATERIALIZED VIEW learn_ck.mv_raw_to_cohort
TO learn_ck.user_cohort AS
SELECT
user_id,
min(event_date) AS first_date,
any(channel) AS channel,
now() AS created_at
FROM learn_ck.events_raw
GROUP BY user_id;19.3.8 商品维度字典
sql
-- 源表假设在 MySQL
CREATE DICTIONARY learn_ck.dict_products
(
product_id UInt32,
name String,
category_id UInt32,
category String,
brand String,
price Decimal(12, 2),
shelve_time DateTime
)
PRIMARY KEY product_id
SOURCE(MYSQL(
host '127.0.0.1' port 3306 user 'ck_reader' password 'xxx'
db 'shop' table 'products'
))
LIFETIME(MIN 300 MAX 600)
LAYOUT(HASHED());查询中用 dictGet:
sql
SELECT
dictGet('learn_ck.dict_products', 'name', product_id) AS pname,
count() AS pv
FROM learn_ck.events_raw
WHERE event_date = today() AND event_type = 'view'
GROUP BY product_id
ORDER BY pv DESC LIMIT 10;19.3.9 TTL:冷热分层 + 自动删除
sql
ALTER TABLE learn_ck.events_raw
MODIFY TTL
event_date + INTERVAL 30 DAY TO VOLUME 'cold', -- 30 天后下 HDD
event_date + INTERVAL 90 DAY DELETE; -- 90 天后删
-- 对应 storage_policy 需要在 config.xml 配置 hot(ssd) / cold(hdd) 两个 volume19.4 典型业务查询 SQL + EXPLAIN
19.4.1 实时大屏:今日 PV / UV / Top10
sql
-- ① 今日实时 PV
SELECT sumMerge(pv) AS pv
FROM learn_ck.events_agg_pv_uv
WHERE event_date = today() AND event_type = 'view';
-- ② 今日实时 UV
SELECT uniqExactMerge(uv) AS uv
FROM learn_ck.events_agg_pv_uv
WHERE event_date = today() AND event_type = 'view';
-- ③ Top10 商品
SELECT
dictGet('learn_ck.dict_products', 'name', product_id) AS pname,
sumMerge(pv) AS pv
FROM learn_ck.events_agg_pv_uv
WHERE event_date = today() AND event_type = 'view'
GROUP BY product_id
ORDER BY pv DESC LIMIT 10;sql
EXPLAIN PIPELINE
SELECT sumMerge(pv) FROM learn_ck.events_agg_pv_uv
WHERE event_date = today() AND event_type = 'view';
-- 期望看到 PartialSorting -> AggregatingInOrder(索引有序,小数据量,毫秒级)预期:1 亿明细 → 聚合表仅 ~百万行 → 查询毫秒级。
19.4.2 漏斗:今日转化
sql
WITH
maxMerge(funnel_state) AS max_step -- 每个用户最深走到第几步
SELECT
countIf(max_step >= 1) AS step_view,
countIf(max_step >= 2) AS step_cart,
countIf(max_step >= 3) AS step_order,
countIf(max_step >= 4) AS step_pay,
step_cart / step_view AS r1,
step_order / step_cart AS r2,
step_pay / step_order AS r3
FROM (
SELECT
user_id,
windowFunnelMerge(3600)(funnel_state) AS max_step
FROM learn_ck.events_agg_funnel
WHERE event_date = today()
GROUP BY user_id
);19.4.3 留存矩阵
sql
SELECT
c.first_date AS cohort,
toInt32(e.event_date - c.first_date) AS offset,
uniqExact(c.user_id) AS retained
FROM learn_ck.user_cohort c
INNER JOIN learn_ck.events_raw e ON c.user_id = e.user_id
WHERE c.first_date >= today() - 30
AND e.event_date BETWEEN c.first_date AND c.first_date + 30
GROUP BY cohort, offset
ORDER BY cohort, offset;生产实战常把这段改成字典 JOIN + 按
cohort_date分桶,避免大表 JOIN;或者用retention函数族 把矩阵直接一次性算出来。
19.4.4 用户路径
sql
SELECT
user_id,
groupArray(event_type) AS path,
groupArray(toString(event_time)) AS t
FROM (
SELECT user_id, event_type, event_time
FROM learn_ck.events_raw
WHERE user_id = 10086 AND event_date = today()
ORDER BY event_time
)
GROUP BY user_id;19.4.5 多维 Cube 即席查询
sql
SELECT
country,
channel,
dictGet('learn_ck.dict_products', 'category', product_id) AS category,
sumMerge(pv) AS pv,
uniqExactMerge(uv) AS uv
FROM learn_ck.events_agg_pv_uv
WHERE event_date BETWEEN today() - 7 AND today()
AND event_type = 'view'
GROUP BY country, channel, category
WITH ROLLUP
ORDER BY pv DESC LIMIT 100;19.5 性能预期(基于同配机器经验值)
环境:单节点 32C / 128G / 4T NVMe,ClickHouse 24.x。数据:每天 3 亿条埋点,共 30 天 ~ 90 亿明细。
| 指标 | 期望值 | 说明 |
|---|---|---|
| 写入吞吐 | ~5 万行/秒稳定、峰值 10 万行/秒 | 由 kafka_num_consumers=4 + kafka_max_block_size=65536 决定 |
| 明细表单分区大小 | ~60GB/天(ZSTD 压缩后 ~18GB) | Map/String 列占比大 |
| 聚合表大小 | 聚合 ≈ 明细的 1/1000 | 仅保留每 hour × event_type × channel × product 粒度 |
| 大屏查询 P99 | < 500 ms | 仅扫聚合表,单分区几百万行 |
| 漏斗查询 P99 | 2 ~ 5 s | windowFunnelMerge + groupByUser |
| 留存查询 P99 | 3 ~ 8 s | JOIN user_cohort,仍在合理范围 |
| Cube 查询 P99 | 10 ~ 30 s | 7 天聚合 + ROLLUP |
监控关键指标:
sql
-- 写入情况
SELECT table, sum(rows) AS rows, count() AS parts
FROM system.parts
WHERE active AND database = 'learn_ck'
GROUP BY table ORDER BY rows DESC;
-- MV 是否卡住
SELECT database, table, last_exception
FROM system.materialized_views
WHERE database = 'learn_ck' AND last_exception != '';
-- Kafka 消费情况
SELECT * FROM system.kafka_consumers
WHERE database = 'learn_ck';19.6 踩坑回顾(本项目实际会踩到的 5 大坑)
坑 1:单条小批写入导致 Part 爆炸
现象:Kafka 消费过快、每秒 1000 次 INSERT,几小时后 system.parts 里单表 Part 数涨到 5 万,查询开始飘到几秒。
根因:ClickHouse 建议每次 INSERT ≥ 1 万行,Kafka 引擎默认 kafka_max_block_size=65536 已经够用。如果你用 Python 自己写 INSERT 循环,一定要攒批。
解决:开 async_insert = 1 或改 Buffer 引擎;Kafka 路径检查 kafka_flush_interval_ms。
坑 2:Distributed 表 internal_replication 写错导致重复
现象:多副本集群,分布式表写入后每行出现 2 ~ 3 次。
根因:internal_replication = false 时,Distributed 表会把同一行往每个副本各写一份,副本再互相同步 → 数据 2~3 倍。
解决:集群是 ReplicatedMergeTree 时必须 internal_replication = true;只有非 Replicated 时才设 false。
坑 3:MV 链路某一层报错导致全链断流
现象:白天突然发现大屏不更新了,查 events_raw 明细还在涨,但 events_agg_pv_uv 不涨。
根因:MV 某次 INSERT 遇到类型不兼容(例如 ext Map 里混入非 UTF-8 字节)抛异常,ClickHouse 默认 materialized_views_ignore_errors = 0,整条 INSERT 回滚。
解决:
sql
SET materialized_views_ignore_errors = 1; -- 某一层 MV 报错不影响其他
SET insert_deduplicate = 1; -- 防止重试导致重复并在 system.query_log + system.text_log 打点告警。
坑 4:大表 JOIN 爆内存
现象:留存查询 events_raw JOIN user_cohort 原始写法直接 OOM。
根因:ClickHouse 默认 Hash JOIN,右表(user_cohort)所有行塞进内存。一旦 user_cohort 过亿,必 OOM。
解决:user_cohort 改成 Dictionary(LAYOUT(HASHED) 只占几 GB)或 Join 表引擎;或把大表放左、小表放右并用 ANY LEFT JOIN;或用 SETTINGS join_algorithm = 'grace_hash' 让 24.x 的 Grace Hash 帮你分批。
坑 5:Mutation 阻塞 Merge
现象:某次误发 ALTER TABLE events_raw DELETE WHERE event_date = '2026-03-01',几分钟后所有查询变慢,SHOW MERGES 全堆着。
根因:Mutation 会重写分区,占用 Merge 线程;只要 Mutation 未完成,后续 Merge 全排队等。
解决:
sql
-- ① 能用 DROP PARTITION 就不要 DELETE
ALTER TABLE events_raw DROP PARTITION '20260301';
-- ② 已经发了 Mutation 且影响业务,立刻 KILL
KILL MUTATION WHERE database='learn_ck' AND table='events_raw' AND mutation_id='mutation_42.txt';
-- ③ 大范围 DELETE 优先用 lightweight DELETE
DELETE FROM events_raw WHERE event_date = '2026-03-01';19.7 运维方案
19.7.1 备份
bash
# 基于 clickhouse-backup(推荐)
clickhouse-backup create full_$(date +%F)
clickhouse-backup upload full_$(date +%F) # 同步到 S3
# 或逻辑备份
clickhouse-client --query="BACKUP TABLE learn_ck.events_raw TO Disk('backups','events_raw_$(date +%F).zip')"19.7.2 监控告警
Prometheus + Grafana,关键指标:
指标(来自 system.* 视图) | 告警阈值 |
|---|---|
system.parts 单表 active parts 数 | > 10 000 报警 |
system.mutations is_done=0 且时间 > 30 min | 报警 |
system.replication_queue 积压 queue_size | > 1000 或 last_exception 非空 |
system.asynchronous_metrics MaxPartCountForPartition | > 300 |
Kafka lag(system.kafka_consumers) | > 10w 报警 |
查询 P99(query_log) | > 5s 报警 |
19.7.3 扩容预案
| 场景 | 方案 |
|---|---|
| 写入吃紧 | 先加 Kafka 消费者(kafka_num_consumers)→ 再 scale out 到多节点,分布式表 rand() 分片 |
| 查询吃紧 | 再加副本(ReplicatedMergeTree),读走 Distributed 自动 round-robin |
| 存储吃紧 | TTL 缩窗;或者冷数据下 S3(storage_policy 配置 s3 磁盘) |
| 热点 SQL | 加 Projection / 新建更细粒度 MV |
19.8 本项目用到了前 18 章的哪些知识点
| 章节 | 本项目用到的点 |
|---|---|
| 01 列存 & 向量化 | events_raw 的列存布局 + ZSTD 压缩让查询只读需要的列 |
| 02 客户端 / HTTP | Python + clickhouse-connect 跑所有查询 |
| 03 数据类型 | LowCardinality / Map / Decimal / IPv4 / DateTime64 全部登场 |
| 04 表引擎 | MergeTree / Kafka / AggregatingMergeTree / ReplacingMergeTree / Dictionary |
| 05 MergeTree 核心 | 稀疏索引 + Skip Index + Projection |
| 06 MergeTree 家族 | 用 AggregatingMergeTree 存状态,ReplacingMergeTree 做用户首日 |
| 07 写入读取 | Kafka kafka_max_block_size 批处理、async_insert、max_threads |
| 08 聚合 / 窗口 | windowFunnel / uniqExactState / argMax / WITH ROLLUP |
| 09 JOIN / 字典 | 商品维度用 Dictionary 替代大表 JOIN |
| 10 MV / Projection | 三条链式 MV + Projection 加速 |
| 11 TTL / 分区 | 30 天冷分层、90 天删除、按天分区 |
| 12 Mutation | DROP PARTITION vs DELETE WHERE 的取舍 |
| 13 副本 / 分布式 | 生产版会把每张表换成 Replicated*MergeTree + Distributed |
| 14 集成生态 | Kafka 引擎 + MySQL 字典源 |
| 15 权限 / 配额 | BI 用户只给聚合表 SELECT、限制 max_memory_usage |
| 16 调优 | EXPLAIN PIPELINE 校验、query_log 抓慢 SQL |
| 17 运维 | BACKUP/RESTORE、Prometheus 监控、扩容预案 |
| 18 生态 | Superset / Metabase 直接连 8123 做大屏 |
一句话:前 18 章是「单招」,这一章是「组合拳」。能把本章画出来讲清楚,就具备了 ClickHouse 初级架构师的基本功。
19.9 本章小结
┌────────────────────────────────────────────────────────────┐
│ 实时埋点分析平台 要点卡 │
├────────────────────────────────────────────────────────────┤
│ 分层:Kafka → Kafka 引擎表 → MV → Agg 表 → 大屏 │
│ 引擎选型:明细 MergeTree;聚合 AggregatingMergeTree │
│ ;维度 Dictionary;接入 Kafka;首日 Replacing │
│ 关键技巧: │
│ ① Projection 加速「查询路径不变」的热点 │
│ ② MV 链做「INSERT 时计算」,状态函数 xxxState 延后 Merge │
│ ③ Dictionary 替代大表 JOIN │
│ ④ TTL 做冷热分层 + 自动删除 │
│ 性能预期:大屏 < 1s,漏斗/留存 < 5s,Cube < 30s │
│ 五大坑:小批写、internal_replication、MV 断流、JOIN OOM、 │
│ Mutation 阻塞 Merge │
└────────────────────────────────────────────────────────────┘19.10 面试高频题(综合项目类)
Q1:给你一个每天 10 亿条埋点、要求秒级大屏 + T+0 漏斗 + 30 天明细留存 的场景,你怎么选型和建表?
考察点:综合架构能力,是否真的能把 OLAP 知识串起来。
标准答案(分 4 层答):
- 采集层:SDK → 批量 HTTP → 网关 → Kafka(
eventstopic,按user_id hash分区)。 - 接入层:ClickHouse
Kafka引擎 + MV 落events_raw(MergeTree,按天分区,(event_date, event_type, user_id, event_time)排序)。不要用原始应用直接 INSERT 单条。 - 聚合层:用
AggregatingMergeTree+ 链式 MV 建三张 DWS:pv_uv(uniqExactState)、funnel(windowFunnelState)、retention(基于user_cohort)。 - 查询层:大屏只扫聚合表(毫秒);明细下钻扫
events_raw(亚秒,命中稀疏索引);字典加速商品维度;TTL 30 天转冷、90 天删除;副本分布式ReplicatedMergeTree+Distributed,写入rand()均匀分片,读走自动路由。
加分项:能提到 Projection 给热点查询加速、async_insert 应对突增、storage_policy 做冷热分层。
易错点:把明细表建成 ReplacingMergeTree 做去重 —— 埋点场景一般不需要精确去重,ReplacingMergeTree 的 FINAL 会让查询慢 10 倍。
Q2:为什么大屏查询要查聚合表而不是明细表?明细表再快不也能扫?
标准答案:
- 明细表 1 天 10 亿行,聚合表 1 天 100 万行,差 1000 倍 I/O;
- 聚合表存的是
uniqExactState状态二进制,uniqExactMerge可以跨时间窗合并,明细表算count(DISTINCT user_id)跨时间必须再扫一遍; - 聚合表可以直接放进大屏的缓存 / BI 工具里,前端压力也小;
- 明细表留着做「下钻」—— 大屏看到异常点击聚合表后再去明细表刨细节。
加分项:能解释 AggregatingMergeTree 的 xxxState 是半成品,不是最终值,SELECT 时必须配 xxxMerge 才正确。
Q3:本项目 Kafka → ClickHouse 这段延迟主要花在哪?怎么优化?
标准答案:
- Kafka Producer 攒批(数百毫秒)+ ClickHouse Kafka 引擎拉取批(默认 500ms + 65536 行)+ MV 计算写入(百毫秒) = 端到端 1~3 秒。
- 优化手段:
kafka_poll_max_batch_size/kafka_max_block_size调大,给批更充足的时间;kafka_flush_interval_ms调到 500;- MV 计算量大时把 SELECT 简化为纯映射,聚合推到下层 MV;
- 写入端开
async_insert = 1+wait_for_async_insert = 0。
Q4:项目上线后突然发现 system.parts 里单表 Part 数飙到 5 万,查询变慢,怎么定位和恢复?
标准答案:
定位:
sqlSELECT table, count() AS parts, sum(rows) FROM system.parts WHERE active AND database='learn_ck' GROUP BY table ORDER BY parts DESC;常见原因:写入端单条或小批 INSERT、MV 一层里面 GROUP BY 粒度过细、OPTIMIZE 太久没跑。
应急:
SYSTEM STOP MERGES learn_ck.events_raw→ 检查 →SYSTEM START MERGES;OPTIMIZE TABLE events_raw PARTITION '20260417' FINAL(只合当天,别全表 FINAL);- 让写入端攒批到 ≥ 10 000 行。
长期:Buffer 引擎 / async_insert / 调大
max_insert_block_size。
Q5:留存查询 SQL 在亿级数据上 OOM,有哪些优化手段?
标准答案:
- 把
user_cohort改成字典(HASHED)或 Join 表引擎,避免 Hash JOIN 右表内存; - 预聚合:建
events_agg_retentionMV,把(cohort_date, offset, channel, retained_users)存成状态; - 分区裁剪 + 子查询先过滤;
SETTINGS join_algorithm = 'grace_hash'或'partial_merge'切算法;- 极限场景用
retention(...)函数族,一次性算出多列留存数组,避免 JOIN。
Q6:项目要上云,你准备迁到 Doris / StarRocks,迁移路径怎么设计?
标准答案(迁移路径 + 能力对照):
| 对照项 | ClickHouse | Doris / StarRocks |
|---|---|---|
| 引擎 | MergeTree 家族 | Unique / Aggregate / Duplicate Key Model |
| 聚合 | AggregatingMergeTree + MV | Rollup + 物化视图 |
| JOIN | 偏弱,推荐字典 | 较强,优化器更成熟 |
| Kafka 接入 | Kafka 引擎 + MV | Routine Load |
| 高可用 | Replicated + Keeper | FE/BE 自带高可用 |
| 生态 | Superset / Metabase / Grafana | Doris 官方 BI,Iceberg 兼容性强 |
迁移步骤:双写过渡(Kafka 同时消费两边)→ 对齐 SQL 方言(windowFunnel 改 Doris 的 window_funnel)→ 聚合表重算 2 天做对比 → 切流量 → 下线 CK。
加分项:能点明「大部分时候不用迁 —— CK 性价比更高;迁的理由往往是 JOIN 强需求 + 向量/半结构化 + 与 Iceberg 湖仓集成」。
📌 收官:至此,19 章正文全部完成。面试题总索引与三份附录请移步:
interview.md/appendix_cheatsheet.md/appendix_ck_vs_mysql.md/appendix_pitfalls.md。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
python
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
第 19 章 综合实战:大屏查询 Demo
-------------------------------
跑 4 大业务查询并打印结果,验证:
① 实时 PV / UV / Top10 商品
② 渠道 × 国家 维度 Cube
③ 最近 1 小时热门商品
④ 用户行为路径(随机取 1 个用户)
"""
import argparse
import time
import clickhouse_connect
def header(title: str) -> None:
bar = '=' * 68
print(f'\n{bar}\n {title}\n{bar}')
def timed(label: str, fn):
t0 = time.time()
rs = fn()
dt = (time.time() - t0) * 1000
print(f'[{label}] {dt:7.1f} ms rows={len(rs.result_rows)}')
return rs
def main():
ap = argparse.ArgumentParser()
ap.add_argument('--host', default='127.0.0.1')
ap.add_argument('--port', default=8123, type=int)
ap.add_argument('--user', default='default')
ap.add_argument('--password', default='')
ap.add_argument('--db', default='learn_ck')
args = ap.parse_args()
cli = clickhouse_connect.get_client(
host=args.host, port=args.port,
username=args.user, password=args.password,
database=args.db,
)
header('① 今日 PV / UV(扫聚合表 events_agg_pv_uv)')
rs = timed('pv_uv', lambda: cli.query("""
SELECT
event_date,
sumMerge(pv) AS pv,
uniqExactMerge(uv) AS uv
FROM events_agg_pv_uv
WHERE event_date >= today() - 2 AND event_type = 'view'
GROUP BY event_date
ORDER BY event_date DESC
"""))
for r in rs.result_rows:
print(f' {r[0]} pv={r[1]:>10} uv={r[2]:>8}')
header('② Top10 商品(走字典映射 + 聚合表)')
rs = timed('top10', lambda: cli.query("""
SELECT
product_id,
dictGet('learn_ck.dict_products', 'name', product_id) AS pname,
dictGet('learn_ck.dict_products', 'category', product_id) AS cat,
sumMerge(pv) AS pv,
uniqExactMerge(uv) AS uv
FROM events_agg_pv_uv
WHERE event_date >= today() - 2 AND event_type = 'view'
GROUP BY product_id
ORDER BY pv DESC
LIMIT 10
"""))
print(f" {'pid':>4} {'name':<14} {'cat':<8} {'pv':>8} {'uv':>6}")
for r in rs.result_rows:
print(f' {r[0]:>4} {r[1]:<14} {r[2]:<8} {r[3]:>8} {r[4]:>6}')
header('③ 渠道 × 国家 Cube(WITH ROLLUP)')
rs = timed('cube', lambda: cli.query("""
SELECT
channel,
country,
sumMerge(pv) AS pv,
uniqExactMerge(uv) AS uv
FROM events_agg_pv_uv
WHERE event_date >= today() - 2 AND event_type = 'view'
GROUP BY channel, country WITH ROLLUP
ORDER BY pv DESC
LIMIT 20
"""))
print(f" {'channel':<8} {'country':<8} {'pv':>10} {'uv':>8}")
for r in rs.result_rows:
ch = r[0] if r[0] else '<ALL>'
co = r[1] if r[1] else '<ALL>'
print(f' {ch:<8} {co:<8} {r[2]:>10} {r[3]:>8}')
header('④ 用户行为路径(随机取 1 个用户)')
uid = cli.query("""
SELECT any(user_id)
FROM events_raw
WHERE event_date = today()
LIMIT 1
""").result_rows
if not uid or uid[0][0] is None:
# 兜底:从最近几天任意挑一个
uid = cli.query("""
SELECT any(user_id)
FROM events_raw
WHERE event_date >= today() - 2
""").result_rows
user_id = uid[0][0] if uid and uid[0][0] is not None else 1
print(f' sampled user_id = {user_id}')
rs = timed('path', lambda: cli.query(f"""
SELECT
event_time,
event_type,
page_url,
product_id
FROM events_raw
WHERE user_id = {user_id}
AND event_date >= today() - 2
ORDER BY event_time
LIMIT 30
"""))
for r in rs.result_rows:
print(f' {r[0]} {r[1]:<10} pid={r[3]:>3} {r[2]}')
header('Done.')
if __name__ == '__main__':
main()python
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
第 19 章 综合实战:漏斗查询 Demo
--------------------------------
演示 windowFunnel 的两种查法:
(A) 直接扫 events_raw(适合偶尔查)
(B) 走 events_agg_funnel 聚合表 + windowFunnelMerge(适合高频查)
"""
import argparse
import time
import clickhouse_connect
def header(title: str) -> None:
print(f'\n{"=" * 68}\n {title}\n{"=" * 68}')
def run(cli, sql: str, label: str):
t0 = time.time()
rs = cli.query(sql)
dt = (time.time() - t0) * 1000
print(f'[{label}] {dt:7.1f} ms')
return rs
def print_funnel(result, title):
if not result.result_rows:
print(f' {title}: empty')
return
r = result.result_rows[0]
v, c, o, p = r[:4]
print(f' {title}:')
print(f' view = {v:>10,}')
print(f' add_cart = {c:>10,} 转化={c/max(v,1)*100:6.2f}%')
print(f' order = {o:>10,} 转化={o/max(c,1)*100:6.2f}%')
print(f' pay = {p:>10,} 转化={p/max(o,1)*100:6.2f}%')
def main():
ap = argparse.ArgumentParser()
ap.add_argument('--host', default='127.0.0.1')
ap.add_argument('--port', default=8123, type=int)
ap.add_argument('--user', default='default')
ap.add_argument('--password', default='')
ap.add_argument('--db', default='learn_ck')
ap.add_argument('--window', default=3600, type=int, help='窗口秒数')
args = ap.parse_args()
cli = clickhouse_connect.get_client(
host=args.host, port=args.port,
username=args.user, password=args.password,
database=args.db,
)
header(f'(A) 直接扫 events_raw,windowFunnel({args.window}s)')
rs = run(cli, f"""
WITH t AS (
SELECT
user_id,
windowFunnel({args.window})(
event_time,
event_type = 'view',
event_type = 'add_cart',
event_type = 'order',
event_type = 'pay'
) AS step
FROM events_raw
WHERE event_date >= today() - 2
GROUP BY user_id
)
SELECT
countIf(step >= 1) AS step_view,
countIf(step >= 2) AS step_cart,
countIf(step >= 3) AS step_order,
countIf(step >= 4) AS step_pay
FROM t
""", 'raw_funnel')
print_funnel(rs, '最近 2 天')
header(f'(B) 走聚合表 events_agg_funnel,windowFunnelMerge({args.window}s)')
rs = run(cli, f"""
WITH t AS (
SELECT
user_id,
windowFunnelMerge({args.window})(funnel_state) AS step
FROM events_agg_funnel
WHERE event_date >= today() - 2
GROUP BY user_id
)
SELECT
countIf(step >= 1) AS step_view,
countIf(step >= 2) AS step_cart,
countIf(step >= 3) AS step_order,
countIf(step >= 4) AS step_pay
FROM t
""", 'agg_funnel')
print_funnel(rs, '最近 2 天(聚合表)')
header('(C) 按渠道拆分漏斗(聚合表)')
rs = run(cli, f"""
WITH t AS (
SELECT
channel,
user_id,
windowFunnelMerge({args.window})(funnel_state) AS step
FROM events_agg_funnel
WHERE event_date >= today() - 2
GROUP BY channel, user_id
)
SELECT
channel,
countIf(step >= 1) AS s1,
countIf(step >= 2) AS s2,
countIf(step >= 3) AS s3,
countIf(step >= 4) AS s4,
round(s2 / greatest(s1, 1) * 100, 2) AS r12,
round(s3 / greatest(s2, 1) * 100, 2) AS r23,
round(s4 / greatest(s3, 1) * 100, 2) AS r34
FROM t
GROUP BY channel
ORDER BY s1 DESC
""", 'by_channel')
print(f" {'ch':<6} {'view':>8} {'cart':>7} {'ord':>7} {'pay':>6} {'V→C':>7} {'C→O':>7} {'O→P':>7}")
for r in rs.result_rows:
print(f' {r[0]:<6} {r[1]:>8} {r[2]:>7} {r[3]:>7} {r[4]:>6} {r[5]:>6}% {r[6]:>6}% {r[7]:>6}%')
if __name__ == '__main__':
main()