上一篇讲清了 列存与 ClickHouse/Doris——Serving 侧有了能秒查的引擎。但广告实时数据从竞价、曝光链路涌进来时,先落哪、洗净到什么程度、汇总到哪一层再给报表,才是数仓真正决定「能不能长期演进」的部分。本篇把 ODS→ADS 分层、Lambda vs Kappa、以及 Flink + OLAP 的流批一体实时仓串成一张可落地的图纸。
本文是实时数仓与 OLAP 系列的第 2 篇。 全系列 3 篇:
- OLAP(开篇)· 列式存储原理与 ClickHouse/Doris
- 实时数仓 · 分层设计与 Lambda/Kappa 架构(本篇)
- 广告多维报表与漏斗分析实战
一句话定位:实时数仓不是「把批仓换成流」这么简单,而是用同一套分层语义,让流式增量与可重放回刷共享口径;Lambda 用双链路换新鲜度与正确性,Kappa / 流批一体则押注「一份逻辑 + 可重放日志」。
TL;DR
- 四层分工:ODS 贴源可重放 → DWD 标准化事实 → DWS 轻度汇总复用 → ADS 面向产品的专题/API。
- 报表禁止直扫 ODS:口径散、重复计算、上游一抖全盘跟着抖。
- Lambda:Speed(低延迟近似)+ Batch(全量准结果)+ Serving 合并——新鲜度好,但对账与双逻辑成本高。
- Kappa:Kafka(或同等可重放日志)上一套流逻辑,改口径靠重放修数——现代 AdTech 实时仓主旋律。
- 上 Kappa 的硬前提:可重放窗口、位点重启、口径版本、重放成本——缺一则是伪 Kappa。
- 流批一体:批是有界流;Flink/Spark 等可用同一套 SQL/作业跑实时与回刷,减少「两个数对不上」。
- 端到端:Kafka → Flink(清洗/对齐/聚合)→ DWD/DWS → OLAP → 报表 API。
- SLA 现实:端到端常见秒~数十秒;倾斜、写入批次、merge/compaction 积压会把实时打成准实时。
Table of contents
Open Table of contents
1. 为什么实时仓必须分层
广告链路的数据源天然混杂:
- 实时事件:竞价、曝光、点击、转化进 Kafka;
- 离线日志 / Dump:补传、对账、冷启动全量;
- 业务维表:广告主、计划、创意、出价策略,常靠 CDC / binlog 同步。
若报表直接查「原始 Kafka 落地表」,三个月后一定会出现:字段含义漂、重复曝光未去重、维表变更无历史、同一个 CTR 有三种算法。分层的价值不是仪式,而是隔离变化——源换格式只动 ODS;口径变更收口在 DWD/DWS;产品改版只动 ADS。
2. ODS / DWD / DWS / ADS:各层写什么
四层像漏斗:越往下越「产品化」,越往上越「保真可重放」。
2.1 ODS · 操作数据层
- 贴源:字段尽量保留原始 JSON/Avro;只做最必要的解析与分区。
- 可重放:出了口径事故,能从 ODS(或上游 Kafka)重放重建下游。
- 不做业务口径:不去「发明」CTR、不在这层硬编码归因窗口。
落地介质可以是对象存储上的原始日志、Kafka 长期保留、或 OLAP/湖上的明细贴源表——关键是语义贴源,不是磁盘形态。
2.2 DWD · 明细数据层
这里才开始「数据工程」:
- 过滤无效流量、刷子(可接 布隆/黑名单 一类前置);
- 曝光/点击去重(按 impression_id 等);
- 对齐维表:补 campaign 名称、广告主行业、计费类型;
- 统一时区、货币、枚举码,输出曝光事实 / 点击事实 / 转化事实。
DWD 是全公司共享的「干净事实」——漏斗、归因、计费核对都应能回到这一层。
DWD 质量门禁(建议写进契约)
上线一份「DWD 准入检查」比事后对口径便宜得多:
- 完整性:关键字段空值率、枚举越界率;
- 唯一性:
imp_id日重复率(超出阈值告警); - 时效性:事件时间 vs 处理时间的延迟分布;
- 维表命中率:join 不上广告主/计划的比例(维表延迟的信号);
- 水量守恒:ODS 行数 ≈ DWD 有效行 + 过滤行(按原因码拆开)。
没有门禁的 DWD,本质上还是 ODS——只是换了张漂亮的表名。
2.3 DWS · 汇总数据层
- 按高频粒度轻度聚合,例如
campaign × media × hour的曝光、点击、花费; - 服务多个 ADS,避免每个专题各自扫 DWD;
- 粒度别切太碎(组合爆炸)也别太粗(报表仍要二次扫明细)。
2.4 ADS · 应用数据层
- 广告主看板、运营驾驶舱、漏斗/留存专题、对外报表 API 宽表或物化;
- 可直接建在 OLAP 的物化视图上(见 末篇);
- 只服务明确产品问法,允许适度冗余。
经验法则:DWD 求真,DWS 求复用,ADS 求体验。 三者口径变更要有版本与发布说明,否则「同一个花费三个数」会再次出现。
3. Lambda:批流双跑的经典答案
Nathan Marz 提出的 Lambda 很直白:
- Batch Layer 周期性跑全量,产出「正确但不够新鲜」的视图;
- Speed Layer 用流处理补「从上次批到现在」的增量,低延迟、可近似;
- Serving Layer 把两份结果合并给查询。
Lambda 用双链路换「又快又准」;Kappa 用可重放 + 单逻辑换「可维护」。
广告侧为何曾经爱 Lambda?因为当时流引擎的正确性/状态容错还没今天稳,用批保证账单级正确,用流给广告主一个「先看看」的实时数。但代价也很硬:
- 两套代码/两套 SQL,CTR、归因窗口稍有不一致就对账爆雷;
- Serving 合并逻辑复杂(批覆盖到哪一刻、速度层从哪切开);
- 人与告警成本随指标数量线性涨。
4. Kappa:一份逻辑,靠重放修数
Kappa 砍掉批层作为「正确性权威」的地位:
- 所有事实以 可重放的日志(通常 Kafka,配合足够保留或落入湖)为真相源;
- 同一套流作业既处理实时,也在改口径/修 bug 时从 earliest / 指定 offset 重放;
- Serving 只消费一份结果。
前提是:流引擎要扛得住状态与 exactly-once(Flink Checkpoint、端到端 Exactly-Once),日志要真的能重放。对 AdTech,Kappa 的诱惑极大——指标定义只有一处,实时看板和事后对账终于能说同一种话。
实践中很少「纯」Kappa:仍会保留低频批核对(抽检、对账作业)或冷数据的批量回刷,但它们不再是产品查询的主路径。这更接近今天常说的流批一体。
4.1 上 Kappa 的硬前提(缺少任一条都会伪 Kappa)
- 日志真能重放:Kafka 保留或同步进湖的窗口 ≥ 口径发版回滚所需天数;
- 作业可从指定位点重启:状态与 sink 幂等策略明确(Exactly-Once);
- 口径有版本号:重放输出写到新版本路径/表,Serving 切换原子;
- 成本可承受:全量重放七天明细不是免费的,要有容量规划与低峰窗口。
「我们消息队列里只有 24 小时,出事靠人肉补 CSV」——那不叫 Kappa,叫运气。
5. 流批一体:批只是有界流
Flink 开篇已强调:批是有界流。落到数仓:
| 模式 | 输入 | 用途 |
|---|---|---|
| 流 | Kafka 无界 | 实时 DWD/DWS/ADS |
| 批 | 历史分区 / 回放 | 口径变更重算、补数、对账 |
| 同一 SQL/作业 | 仅执行模式不同 | 消灭双逻辑 |
运维上建议:口径变更走「发版 + 重放 + 切换 Serving 指针」,而不是「先改速度层、批层下个月再说」——那是 Lambda 漂移的复发路径。
5.1 湖仓一体怎么放进这张图
不少团队会问:Iceberg/Hudi/Paimon 跟 ODS–ADS 什么关系?简洁看法:
- 湖表很适合做 ODS / 历史 DWD 的廉价可重放底座(列存文件 + 时间旅行);
- OLAP(CH/Doris) 继续承担 ADS / 热 DWS 的交互式查询;
- Flink 可同时写湖与 OLAP,或「湖为源、OLAP 为Serving 投影」。
不要用湖替代所有 OLAP 秒查,也不要让 OLAP 存「永远不查的冷历史」。分层同样适用于存储介质:热在线、温可查、冷可重放。
6. Flink + OLAP 的端到端实时仓
主路径管新鲜度与分层落点;维表广播管对齐;Exactly-Once 与物化管「敢写敢查」。
一条可生产的参考链路:
- Kafka 承载曝光/点击/转化(分区与并行度对齐见 Flink 篇);
- Flink 做 DWD:解析、过滤、去重、维表 join(Broadcast State)、窗口轻度聚合;
- 写出到 OLAP(ClickHouse / Doris)的 DWD/DWS 表,或先入湖再灌;
- OLAP 物化 / 投影生成 ADS;
- 报表 API 只查 ADS,长尾 dig 受限流与隔离队列访问 DWD。
新鲜度 SLA 要写进契约:例如「事件进入 Kafka 到 ADS 可查 P99 ≤ 30s」。超时告警应区分:消费 lag、作业反压、OLAP 写入/merge 积压——三者修法完全不同。
6.1 Schema 演进怎么不伤报表
广告事件字段几乎每月都在加。推荐约定:
- ODS 允许多余字段(Map/JSON 尾部),向前兼容;
- DWD 显式列 + 版本:新指标进新列,废弃列标记但不立刻删;
- 下游 ADS 只选表白名单列,禁止
SELECT *; - 破坏性变更走发版:双写过渡期 → 校验 → 切读 → 下线旧列。
Flink SQL / 表 schema 与 OLAP DDL 最好进同一变更评审,避免「流作业已写新字段、报表库还没加列」的深夜事故。
6.2 多租户与资源隔离
一条实时仓往往会同时服务:广告主看板、内部运营 dig、算法离线样本回流、财务对账。若不做隔离:
- 一条 exploration SQL 扫三个月明细,打满 OLAP;
- 算法回刷抢 Flink 集群 slot,实时 pacing 报表延迟飙升。
最低配置:查询队列分级(看板 / dig / 回刷)、作业资源池拆分、对象存储冷热分层。这和列存选型同等重要。
7. 分层反模式(广告仓常见)
- ADS 直连 ODS 原始 JSON:字段一变前端一起炸。
- 每个报表团队自建一套「准 DWD」:同一曝光三种去重规则。
- DWS 粒度按「所有维度笛卡尔积」预聚合:组合爆炸,写入比查询还贵。
- 速度层「先上线、口径以后对齐」:以后永远不对齐。
- Kafka 保留 24h 却声称 Kappa:出事无法重放,只是伪 Kappa。
- 用「实时数」给财务结账:速度层近似或未稳定窗口内的数进了出账,合规与客诉双爆。
8. 何时还该保留批层
不是教条消灭 Spark/批。以下场景批仍然香:
- 超大历史重算(年级别),流重放成本明显高于批扫湖;
- 复杂图计算 / 全站训练样本,本就不是 Serving 路径;
- 监管或财务要求「日终冻结快照」——可以是同一逻辑在日终跑一遍有界流,产出不可变分区。
关键是:批不再偷偷定义另一套指标语言。 它是同一套语义在另一种执行模式下的运行。
参考
- Nathan Marz. How to beat the CAP theorem / Lambda Architecture:Lambda 的原始论述。
- Jay Kreps. Questioning the Lambda Architecture:Kappa 视角的经典批评与主张。
- Apache Flink. Streaming Concepts:流批一体与有界/无界流。
- 本系列开篇. 列式存储原理与 ClickHouse/Doris:Serving 引擎地基。