KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
06 · OLAP 选型与管道运维 — keel 龙骨
上一章的 DWS 是一张普通表。这一章回答两个收尾问题:要不要为分析负载换一个专门的 OLAP 引擎;以及这条管道上线之后,人看哪些数字才能睡得着。
上一章的 DWS 是一张普通表。这一章回答两个收尾问题:要不要为分析负载换一个专门的 OLAP 引擎;以及这条管道上线之后,人看哪些数字才能睡得着。
现场
报表从 300 毫秒慢到 3 秒。有人提议「上 ClickHouse」,有人说「先加个物化视图」。两种说法都对了一半。
问题不在「哪个引擎更快」,而在你的慢是哪一类慢:是扫描的数据量太大(列存能帮上),还是每次都重算一样的聚合(预聚合能帮上),还是查询模式根本不适合当前建模(换引擎也救不了)。本机没有 ClickHouse / Doris,所以这一章分两部分:用官方文档说明这些引擎的设计取向,用一台 PostgreSQL 上的行存宽表 vs 物化视图对比量出「预聚合」这一类优化的真实收益。
观测面
flowchart TD
subgraph UP["上游:产生与搬运"]
DB["业务库<br/>复制槽 cdc_slot"] --> P["管道进程<br/>延迟 / 积压 / offset"]
P --> K["Kafka topic<br/>lag / 分区位点"]
end
subgraph MID["中游:加工"]
K --> E["流处理<br/>水位线滞后 / 丢弃迟到数 / checkpoint 成功率"]
E --> W["数仓分层<br/>ODS/DWD/DWS 行数比 · 刷新延迟"]
end
subgraph DOWN["下游:服务查询"]
W --> OLAP["分析存储<br/>查询耗时 P95 / 扫描行数"]
OLAP --> BI["报表 / 接口"]
end
MON["运维指标面"] -. 看槽的 restart_lsn 落后量 .-> DB
MON -. 看 task 状态与 offset 是否推进 .-> P
MON -. 看窗口状态大小与丢弃计数 .-> E
MON -. 看刷新是否成功、视图是否过期 .-> W
MON -. 看是否走了预聚合 .-> OLAP
LIN["数据质量与血缘"] -. "列级血缘:orders.amount → dws.amount_sum" .-> W
LIN -. "口径变更影响面" .-> BI
style DB fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
style P fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style K fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style E fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style W fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style OLAP fill:#f3f0ff,stroke:#7048e8,color:#1a1a2e
style BI fill:#f3f0ff,stroke:#7048e8,color:#1a1a2e
style MON fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
style LIN fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
观测面的每一段都有它自己的失效方式:上游是「槽没人消费」,中游是「状态一直涨」,下游是「预聚合过期了但没人知道」。一条链路的可用性是这些段里最短的那一段决定的。
精确定义
列式存储:同一列的值连续存放。它的优势不在「读得快」本身,而在只读需要的列——一个只用到 3 列的分析查询,行存要读出整行,列存只读那 3 列。代价是单行写入要散到多个列文件。
向量化执行:按批(block)而不是按行处理,减少虚函数调用、提高 CPU 缓存命中率、便于用 SIMD。它和列存是一对搭档,但可以独立存在。
MPP(大规模并行处理):查询被拆到多个节点并行执行,节点间 shuffle。它决定的是「一张大表能不能被多台机器同时扫」。
物化视图:把查询结果物化下来。分两种——增量型在写入时就被触发更新(ClickHouse 的 incremental MV 官方描述是「把计算成本从查询时挪到写入时」),可刷新型按周期或手动全量重算(ClickHouse 的 refreshable MV、PostgreSQL 的 REFRESH MATERIALIZED VIEW)。本机实验中用的是后者。
一次完整运行:行存宽表 vs 物化视图
原始输出见 lab/evidence/stream-data-pipeline/07-olap.txt。一张 100 万行、88 MB 的行存明细表,同一个查询「每个状态多少单、金额多少」:
--- 直接查明细表 ---
Finalize GroupAggregate ... (actual time=267.163..312.357 rows=5)
-> Gather Merge ... Parallel Seq Scan on olap_events
Execution Time: 312.475 ms
--- 建物化视图 ---
SELECT 5 Time: 501.786 ms mv_size = 16 kB
--- 查物化视图 ---
Seq Scan on mv_olap_by_status (actual time=0.010..0.011 rows=5)
Execution Time: 0.018 ms
312.475 ms → 0.018 ms,约 17000 倍。这个倍数不是「引擎变强了」,是问题变简单了:原本要在 100 万行上做分组聚合,现在只读 5 行。把这 5 行的计算挪到了建视图的那 502 ms 里。
然后是预聚合一定会付的代价——过期。灌进 20 万新行之后:
物化视图(未刷新): mv_total = 1000000
明细表(真实): base_total = 1200000
刷新物化视图: Time: 488.259 ms
刷新后: mv_total_after_refresh = 1200000
视图在这段时间里对外给的是错的数字,而且它不会报错、不会告警,只是安静地少 20 万。刷新一次 488 ms,也就是说「刷新频率」直接决定了「报表最大误差窗口」。
最后一个对照,说明预聚合的边界:
明细点查: explain analyze select count(*) from olap_events
where customer_id = 4321 and status = 'paid';
-> Parallel Seq Scan ... Rows Removed by Filter: 400000
Execution Time: 272.264 ms
同样的问题问物化视图: 只返回 status='paid' 那一行(240000 单),答不了「这个客户的」
物化视图是按 status 聚合的,它根本没有 customer_id 这个维度。预聚合只能回答它预先算过的问题;换一个维度,回到 272 ms 的全表扫。所以真实系统里会看到「一张主表 + 好几个不同维度的物化视图」,每多一个维度就多一份写入放大和一份过期风险。
官方文档怎么规定(检索日期 2026-10-05)
以下为 ClickHouse 与 Apache Doris 官方文档陈述,不属于本机实测(本机跑不了这两个引擎):
- ClickHouse MergeTree:表由「按主键排序的 data part」组成,插入会创建新的 part,由后台进程合并;主键不指向单行而是 8192 行一个的 granule,所以主键索引可以整块装进内存;分区裁剪能跳过不需要的 part;Wide 格式下「每一列存在单独的文件里」(ClickHouse 官方 MergeTree 文档)。需要注意的是,这一页并没有直接使用「列式存储」和「向量化执行」这两个词——它的表述落在「按列存文件」「granule 是最小读取单位」「读取自动并行化」上。把「列存/向量化」当成能力标签时,最好回到这些可核验的描述。
- ClickHouse 物化视图:官方把增量型物化视图的定位写成「把计算成本从查询时挪到插入时,从而让
SELECT更快」;另有 refreshable 型,需要周期性地对全量数据执行查询并把结果写进目标表(ClickHouse 官方 Materialized views 文档)。这正是本机REFRESH MATERIALIZED VIEW那条路的对应物。 - Apache Doris:官方定位是「基于 MPP 的实时数据仓库」,强调亚秒级查询;存储侧「列式存储引擎,按列编码、压缩、读取」;查询侧「查询引擎完全向量化,所有内存结构按列布局……官方称在宽表聚合场景下比非向量化引擎快 5~10 倍」;模型侧提供明细模型、主键模型和聚合模型(同 key 的 value 列合并,用预聚合提升性能);物化视图分单表(系统自动刷新维护)和多表(周期刷新)两类(Apache Doris 官方文档)。
- Flink 的端到端恰好一次(补第 04 章):Kafka sink 的
DeliveryGuarantee.EXACTLY_ONCE把消息写进 Kafka 事务并在 checkpoint 时提交,要求transactionalIdPrefix全局唯一,代价是记录在 checkpoint 完成前对外不可见(Flink 1.20 Kafka connector 文档)。这解释了为什么上端到端事务会把结果延迟和 checkpoint 间隔绑在一起。
选型判断
把上面的东西收成几句可以直接用的判断:
- 慢在扫描的数据量(宽表、列多、只取几列、大范围过滤)→ 列存 + 分区裁剪是主要收益来源,换 OLAP 引擎值得。
- 慢在每次都重算同样的聚合(固定几个维度、固定几张报表)→ 先上物化视图/预聚合,不动引擎也能拿到数量级的提升(本机 17000 倍)。
- 慢在查询模式本身(每个查询都换维度、都要明细)→ 预聚合帮不上,列存也只能缓解扫描,真正的解法是重新建模(第 05 章的分层)。
- 预聚合的隐性成本是过期:多一个视图就多一份刷新调度、多一个「刷失败/刷慢了」的告警面。视图数量要克制。
运维指标与数据质量
上线后需要盯的(按观测面从上到下):
- 上游:每个复制槽的
pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)、active是否长期为f、pg_wal目录大小。这一条最容易被忽略,也最容易直接把主库写死。 - 搬运:Kafka 各分区 lag、Connect task 的状态(RUNNING / FAILED)、offset 提交是否推进。第 02 章实测过:task 启动失败不会自愈,只能靠状态告警发现。
- 加工:水位线相对最新事件时间的滞后(= 结果延迟)、被丢弃的迟到事件数、checkpoint 成功率、状态大小。丢弃计数从 0 变正,说明口径开始偏了。
- 落地:ODS/DWD/DWS 的行数比、DWS 的最近刷新时间与刷新耗时、物化视图是否落后于基表。本机那条「100 万 vs 120 万」的偏差,只有主动对比才会暴露。
- 查询:P95 耗时、扫描行数、是否命中预聚合。只看平均耗时会把「少数走了全表扫的查询」平均掉。
数据质量与血缘在链路里是同一件事的两面:血缘回答「这一列从哪来」,质量规则回答「它还在不在合理范围」。列级血缘(orders.amount → ods.amount → dwd.amount → dws.amount_sum)决定了口径变更时的影响面;质量规则则应挂在每一跳的出口——比如「DWS 行数不超过 DWD 状态基数」「ODS 的 DELETE 数不超过基线 3 倍」,越靠近产生点,定位越快。
生产边界
- 教学替身 vs 真实依赖:本机只有一台 PostgreSQL,行存 + 物化视图是「不加新组件能拿到的最大收益」;真实环境会引入列存 MPP 引擎和对象存储湖表,随之而来的是数据搬运、元数据同步和跨系统一致性,这些成本本机看不到。
- 失败策略:物化视图刷新失败要保留上一版可用并告警;列存引擎的导入任务失败要能重跑(幂等键落在业务主键上);预聚合与明细不一致时,以明细为准并触发刷新。
- 容量:列存 + 预聚合会让同一份数据存多份,压缩率高的列存能缓解,但「多几个视图就多几份存储」的账要提前算。
动手
- 跑
exp07_olap.sh,把group by status换成group by status, date_trunc('month', created_at),比较基表和物化视图的耗时变化,说明为什么倍数变了。 - 给
olap_events(customer_id, status)建一个索引,重跑那个 272 ms 的明细点查,比较加索引前后;再讨论「为什么物化视图不能替代这个索引」。 - 断言题:写一条 SQL 对比
mv_olap_by_status的总行数与olap_events的行数,灌一批新数据后运行,确认它能检测出「视图过期」。
自测
- 列存、向量化、MPP、物化视图各自解决慢的哪一类原因?分别对应本机实验里的哪个数字?
- 本机物化视图把查询从 312 ms 降到 0.018 ms,这个倍数里有多少来自「引擎变强」,多少来自「问题变简单」?
- 物化视图过期时为什么不报错?你会用什么机制发现它?
- 预聚合只能回答预先算过的问题——本机那个 272 ms 的明细点查说明了什么?
- 一条 CDC 链路要上线,你会按什么顺序布置监控?为什么上游的复制槽指标排在最前面?
这门课到这里结束。回头看这条链:一次 UPDATE 在 WAL 里留下物理记录(00)→ 逻辑槽把它解码成表/列/值,并用
restart_lsn决定主库能删多少 WAL(01)→ 连接器把变更搬进 topic,位点决定断了从哪续(02)→ 事件时间和水位线决定窗口算出来是多少(03)→ 快照和幂等键决定崩溃恢复后会不会多算(04)→ 分层决定这个数字存在哪、多快能查到(05)→ 选型和监控决定它上线以后能不能一直对。这条链上的每一环都会以不同的方式失效,而症状往往出现在另一环——这正是流式管道最难排查的地方。