Skip to content

第 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 → 网关 → KafkaNginx+Lua / Kafka
接入层ClickHouse 消费 KafkaKafka 引擎 + 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;

关键设计解释:

  1. 分区按天toYYYYMMDD):每天一个分区,TTL / 冷热分层都按分区搬;
  2. 排序键 (event_date, event_type, user_id, event_time)
    • 大屏类按 event_type + 日期过滤,命中稀疏索引;
    • 用户路径查询按 user_id,命中稀疏索引中段;
  3. 主键只取排序键前 3 列:稀疏主键保留在内存,省内存;
  4. LowCardinality 用在基数 ≤ 1 万的枚举列(事件类型、渠道、国家);
  5. String CODEC(ZSTD(3)) 给文本类(URL、UA、Referrer)单列压缩,节省 60% 空间;
  6. Map(String, String) 预留扩展字段,业务早期不用改表结构;
  7. 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) 两个 volume

19.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仅扫聚合表,单分区几百万行
漏斗查询 P992 ~ 5 swindowFunnelMerge + groupByUser
留存查询 P993 ~ 8 sJOIN user_cohort,仍在合理范围
Cube 查询 P9910 ~ 30 s7 天聚合 + 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 改成 DictionaryLAYOUT(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 客户端 / HTTPPython + 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_insertmax_threads
08 聚合 / 窗口windowFunnel / uniqExactState / argMax / WITH ROLLUP
09 JOIN / 字典商品维度用 Dictionary 替代大表 JOIN
10 MV / Projection三条链式 MV + Projection 加速
11 TTL / 分区30 天冷分层、90 天删除、按天分区
12 MutationDROP 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 层答):

  1. 采集层:SDK → 批量 HTTP → 网关 → Kafka(events topic,按 user_id hash 分区)。
  2. 接入层:ClickHouse Kafka 引擎 + MV 落 events_rawMergeTree,按天分区,(event_date, event_type, user_id, event_time) 排序)。不要用原始应用直接 INSERT 单条
  3. 聚合层:用 AggregatingMergeTree + 链式 MV 建三张 DWS:pv_uvuniqExactState)、funnelwindowFunnelState)、retention(基于 user_cohort)。
  4. 查询层:大屏只扫聚合表(毫秒);明细下钻扫 events_raw(亚秒,命中稀疏索引);字典加速商品维度;TTL 30 天转冷、90 天删除;副本分布式 ReplicatedMergeTree + Distributed,写入 rand() 均匀分片,读走自动路由。

加分项:能提到 Projection 给热点查询加速、async_insert 应对突增、storage_policy 做冷热分层。

易错点:把明细表建成 ReplacingMergeTree 做去重 —— 埋点场景一般不需要精确去重,ReplacingMergeTreeFINAL 会让查询慢 10 倍。


Q2:为什么大屏查询要查聚合表而不是明细表?明细表再快不也能扫?

标准答案

  1. 明细表 1 天 10 亿行,聚合表 1 天 100 万行,差 1000 倍 I/O;
  2. 聚合表存的是 uniqExactState 状态二进制uniqExactMerge 可以跨时间窗合并,明细表算 count(DISTINCT user_id) 跨时间必须再扫一遍;
  3. 聚合表可以直接放进大屏的缓存 / BI 工具里,前端压力也小;
  4. 明细表留着做「下钻」—— 大屏看到异常点击聚合表后再去明细表刨细节。

加分项:能解释 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 万,查询变慢,怎么定位和恢复?

标准答案

  1. 定位

    sql
    SELECT table, count() AS parts, sum(rows)
    FROM system.parts WHERE active AND database='learn_ck' GROUP BY table ORDER BY parts DESC;
  2. 常见原因:写入端单条或小批 INSERT、MV 一层里面 GROUP BY 粒度过细、OPTIMIZE 太久没跑。

  3. 应急

    • SYSTEM STOP MERGES learn_ck.events_raw → 检查 → SYSTEM START MERGES
    • OPTIMIZE TABLE events_raw PARTITION '20260417' FINAL(只合当天,别全表 FINAL);
    • 让写入端攒批到 ≥ 10 000 行。
  4. 长期:Buffer 引擎 / async_insert / 调大 max_insert_block_size


Q5:留存查询 SQL 在亿级数据上 OOM,有哪些优化手段?

标准答案

  1. user_cohort 改成字典(HASHED)或 Join 表引擎,避免 Hash JOIN 右表内存;
  2. 预聚合:建 events_agg_retention MV,把 (cohort_date, offset, channel, retained_users) 存成状态;
  3. 分区裁剪 + 子查询先过滤;
  4. SETTINGS join_algorithm = 'grace_hash''partial_merge' 切算法;
  5. 极限场景用 retention(...) 函数族,一次性算出多列留存数组,避免 JOIN。

Q6:项目要上云,你准备迁到 Doris / StarRocks,迁移路径怎么设计?

标准答案(迁移路径 + 能力对照)

对照项ClickHouseDoris / StarRocks
引擎MergeTree 家族Unique / Aggregate / Duplicate Key Model
聚合AggregatingMergeTree + MVRollup + 物化视图
JOIN偏弱,推荐字典较强,优化器更成熟
Kafka 接入Kafka 引擎 + MVRoutine Load
高可用Replicated + KeeperFE/BE 自带高可用
生态Superset / Metabase / GrafanaDoris 官方 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()

dashboard_query.py ↗ · funnel_query.py ↗