上一篇讲清了 列存与 ClickHouse/Doris——Serving 侧有了能秒查的引擎地基。然而,当广告主质问「为什么同一个 CTR 跑出三种算法结果」「看板怎么卡在十几分钟前」或是「结算 UV 为何与报表对不上」时,单靠底层引擎的扫列速度已无济于事。当实时数据从竞价与曝光链路汹涌而入,数据先落哪、洗净到什么程度、在何处固化口径,才是数仓能否长期演进的决定性命题。本篇将把 ODS 到 ADS 的分层设计、Lambda 与 Kappa 架构的取舍,以及 Flink + OLAP 的流批一体端到端管线串联成一张可落地的图纸。
本文是 实时数仓与 OLAP 系列的第 2 篇(分层与架构)。 全系列 4 篇:
- OLAP · 列式存储原理与 ClickHouse/Doris
- 实时数仓 · 分层设计与 Lambda/Kappa 架构(本篇)
- 广告多维报表与漏斗分析实战
- 时序数据库 InfluxDB 深挖 · 数据模型、TSM 与降采样
一句话定位:实时数仓绝非简单地将批处理替换为流计算,而是通过同一套分层语义让流式增量与可重放回刷共享唯一的指标口径;Lambda 依靠双链路换取新鲜度,而 Kappa 与流批一体则押注于「一份逻辑 + 可重放日志」。
TL;DR
- 四层分工:ODS 贴源可重放 → DWD 标准化事实 → DWS 轻度汇总复用 → ADS 面向产品的专题与 API。
- 直扫反模式:严禁报表直扫 ODS,否则不仅极易引发重复计算,上游格式的任何抖动都会导致全盘指标瘫痪。
- Lambda:Speed(低延迟近似)+ Batch(全量准结果)+ Serving 合并——新鲜度极佳,但对账与双逻辑维护成本高昂。
- Kappa:基于 Kafka 等可重放日志运行一套流逻辑,修改口径依赖重放修数——这是现代 AdTech 实时仓的主旋律。
- 上 Kappa 的硬前提:可重放窗口、精准位点重启、口径版本与重放成本缺一不可,缺少任一条即为伪 Kappa。
- 流批一体:批是有界流;Flink 等引擎可用同一套 SQL 兼顾实时接入与历史回刷,彻底消除口径漂移。
- 端到端管线:Kafka → Flink(清洗/对齐/聚合)→ DWD/DWS upsert → OLAP → 报表 API。
- SLA 现实:端到端可见延迟通常在秒至数十秒量级;Flink 反压、数据倾斜或 OLAP 端的 merge / compaction 积压,都会随时将其降级为准实时甚至更糟。
Table of contents
Open Table of contents
1. 为什么实时仓必须分层
在广告业务中,数据源天然呈现出极度混杂的态势:实时事件(如竞价、曝光、点击与转化)以极高的吞吐直接流入 Kafka;离线 Dump 被用于延迟补传与财务对账;而广告主、计划等维表状态则依赖 CDC 同步。
如果放任前端报表直接去查询未经处理的原始 Kafka 落地表,不出三个月系统必定走向失控。不仅上游 JSON 字段的随意增删会让解析逻辑散落各处,重复曝光未被剔除、状态变更丢失历史轨迹等问题,更是会让同一份报表跑出截然不同的口径。正因如此,分层的核心价值并非为了追求架构上的仪式感,而是为了隔离变化。通过分层,当数据源格式更迭时只需调整 ODS 的解析逻辑,复杂的聚合与去重被严格收拢在 DWD 与 DWS,而前端看板的改版仅需在 ADS 层做轻量适配。
2. ODS / DWD / DWS / ADS:各层职责划分
分层架构的语义递进:越靠近 ADS 越产品化,越靠近 ODS 越贴源可重放。
2.1 ODS · 操作数据层
作为实时仓的入口,ODS(Operational Data Store)的铁律是贴源存储。字段应尽可能原样保留最初的 JSON 或 Avro 结构,仅做最基础的时间戳提取与按日分区。这一层的根本使命在于保障可重放能力:一旦下游发生大面积的口径事故,系统必须能够直接切回 ODS(或更上游的 Kafka),从原始形态重新推演并回刷整个下游链路。因此,严禁在 ODS 层「发明」业务逻辑或硬编码归因窗口。无论是存放在对象存储、长周期保留的 Kafka,还是数据湖,只要语义贴源,它就是合格的 ODS。
2.2 DWD · 明细数据层
数据工程的脏活与重头戏,往往全部压在 DWD(Data Warehouse Detail)这一层:
- 流量清洗:剔除作弊请求与爬虫流量,此处常前置 布隆过滤器 进行高效的穿透拦截;
- 精确去重:严格按照
impression_id等唯一业务标识对曝光与点击进行清洗; - 维表对齐:把枯燥的 ID 扩充为具体的 campaign 名称、广告主行业或计费类型;
- 标准化事实:统一时区、货币并对齐枚举值,最终产出全局通用的曝光事实表与转化事实表。
DWD 是全公司共享的「唯一事实底座」。无论是事后的漏斗留存分析,还是严苛的结算对账,都必须能无缝追溯到这里的明细记录。为了守住这道防线,强烈建议设立强制的 DWD 质量门禁:监控关键字段的空值率、追踪事件时间与处理时间之间的可见延迟(Event Time vs Processing Time)、校验 imp_id 的日重复率,并确保上游总行数与过滤后的有效水量绝对守恒。没有质量门禁的 DWD,不过是披着明细外衣的 ODS。
2.3 DWS · 汇总数据层
面对广告多维报表动辄扫数月的查询,DWS(Data Warehouse Service)通过轻度预聚合来提前消解计算压力。通过将细粒度的事件按业务高频维度(如 campaign × 媒体 × 小时)进行汇总,DWS 能作为一个高复用的中间件服务于多个前端场景。但在设计时必须权衡粒度:切得过细会引发维度组合爆炸,使得预聚合失去意义;聚合得太粗,又会导致报表层在处理下钻分析时被迫回退到全量扫描 DWD。
2.4 ADS · 应用数据层
ADS(Application Data Store)完全呈现出产品导向的特质,专为广告主看板、运营驾驶舱、报表 API 等明确场景定制宽表。为了追求极致的端到端 P99 延迟,这一层极其依赖 OLAP 引擎的物化视图(Materialized View)来实现查询加速(详情将在本系列 收官篇 深入探讨)。在此处,适度的数据冗余被视为换取体验的合理代价。
经验法则:DWD 求真,DWS 求复用,ADS 求体验。 这三层的口径变更必须配备极其严格的版本控制,否则「同一个花费在两个报表里跑出不同数字」的口径漂移灾难必将重现。
3. Lambda:批流双跑的经典答案
在流引擎发展早期,面对实时性与准确性的双重拷问,Nathan Marz 提出的 Lambda 架构给出了一个看似完美的妥协:
- Batch Layer(批处理层):定期运行全量作业,产出绝对正确但缺乏新鲜度的历史视图;
- Speed Layer(速度层):用流处理引擎消费增量数据,快速提供低延迟的近似估算;
- Serving Layer(服务层):在查询发生时,由引擎当场合并这两份结果。
Lambda 依靠双链路换取低延迟与正确性;Kappa 则押注可重放日志与单逻辑以换取可维护性。
早期的广告系统之所以对 Lambda 趋之若鹜,是因为当时的流计算在状态容错与 Exactly-Once 语义上尚不堪大用。团队只能靠批层来兜底账单级别的真相,同时用速度层勉强撑起大屏。然而这种架构的反噬极快:双重逻辑维护意味着你必须用两种不同的 API(如 Storm 与 Hadoop)写出绝对对齐的 CTR 算法,稍有偏差便会导致口径漂移;同时,Serving 层的合并逻辑异常繁琐,必须精准裁定批数据的生效边界并截断对应的流数据;随着指标日益繁杂,人力对账与排障成本很快就会失控。
4. Kappa:一份逻辑,靠重放修数
为了根除双写逻辑的梦魇,Kappa 架构果断剥夺了批层作为「唯一真相源」的特权:
- 所有的业务事实均被视作可重放的日志(配置了充足保留时长的 Kafka 或是沉淀入湖的数据);
- 依托同一套流作业不仅处理实时增量,在面临业务口径变更或逻辑修复时,直接从指定位点发起重放(Replay)来刷写历史;
- Serving 层自始至终只对着单一的聚合结果提供查询。
Kappa 成立的前提,在于底层流引擎必须具备极其强悍的容错能力(如 Flink Checkpoint 与 端到端 Exactly-Once 机制)。这对 AdTech 团队的诱惑是致命的——它终于把复杂的指标定义收拢到了一处,让实时看板与月末结算真正实现了口径的统一。
4.1 上 Kappa 的硬前提
现实中极少有纯粹的 Kappa,很多团队跟风宣称上了 Kappa,却连基础的容灾底座都不具备。要真正发挥单逻辑重放的威力,以下条件缺一不可(缺少任一条即为伪 Kappa):
- 日志具备真实重放能力:底层 Kafka 的保留周期(或入湖窗口期)必须长于口径发版回滚所需的最大跨度;
- 作业支持精准位点重启:引擎的状态管理与下游 OLAP 的幂等写入必须严丝合缝;
- 指标定义版本化:重放时必须将新口径的数据写至带版本号的新表或新分区,且 Serving 层的读指针切换必须具备原子性;
- 重放成本可控:回刷数月海量明细并非轻而易举,必须预估好引擎在追赶期引发的计算资源抢占。
如果你的消息队列只保留 24 小时,出事全靠临时导 CSV 补数,这就不要自称为 Kappa 架构。
5. 流批一体:批是有界流
正如 Flink 系列开篇所强调的核心理念:批处理本质上不过是拥有明确终点的「有界流」。这一理念映射到分层架构中,就是当今被广泛讨论的流批一体。
| 模式 | 输入 | 用途 |
|---|---|---|
| 流计算 | 无界的 Kafka 增量流 | 构建毫秒或秒级可见的实时 DWD / DWS |
| 批处理 | 历史分区或湖中归档 | 执行口径变更重算、批量补数与深度对账 |
回到刚才那条宽表聚合的痛点,无论是实时流式消费还是面向历史数据的离线批处理,现在都能够复用同一套 SQL 与算子逻辑。在运维层面,这意味着任何口径变更都可以被收敛为「发布新作业 → 历史重放补齐 → 原子切换前端读指针」的标准闭环,从而彻底消灭了「先上线速度层,以后再跟批层对齐」这种自欺欺人的开发习惯。
5.1 湖仓一体如何融入分层架构
许多团队在做引擎选型时常有疑惑:像 Iceberg、Hudi 这样的数据湖,与 ODS 到 ADS 的分层到底如何适配?最务实的定位是让系统各司其职:
- 数据湖凭借廉价的存储与良好的时间旅行机制,极其适合充当 ODS 与历史 DWD 的贴源底座;
- OLAP 引擎(如 ClickHouse、Doris) 凭借极致的向量化扫列能力,继续死磕 热 DWS 与 ADS 的交互式查询;
- 在管道层,Flink 可以采用双写策略,同时将清洗好的数据沉入数据湖并 upsert 至 OLAP 进行快速 Serving。
切忌试图用高昂的 OLAP 内存去硬抗「永远不会被查询的冷历史」,也别指望用数据湖去包揽所有广告主报表的亚秒级响应。分层不仅是业务语义的递进,同样也是介质成本的规划:热数据保持在线,温数据随时可查,冷数据低成本可重放。
6. Flink + OLAP 的端到端实时仓
主路径决定了数据的分层落点;维表广播保障了维度对齐;Exactly-Once 语义与预聚合则赋予了下游「敢写敢查」的底气。
将上述理论拼装,一条经受过大流量考验的生产管线通常如下运转:
- 事件接入:海量的竞价、曝光与点击汇入 Kafka,以分区机制保障核心链路的并发吞吐;
- 流式清洗:Flink 承接 ODS 数据,执行过滤、基于状态的精确去重,通过 Broadcast State 完成广告主维表 Join,最后按时间窗口执行轻度聚合,将散乱的日志转化为高度标准化的 DWD/DWS 事实;
- 高效落盘:将清洗对齐后的增量数据,通过追加或覆盖写(upsert)稳稳推入 ClickHouse 或 Doris 这类 OLAP 引擎;
- 视图加速:在 OLAP 内部,依托后台的 merge 机制与物化视图,进一步固化出专供报表使用的 ADS 层;
- 查询隔离:前端业务的 报表 API 仅被授权查询 ADS 层,而任何面向明细溯源的长尾 dig 查询,都必须被赶入独立的隔离队列中去慢查 DWD。
端到端可见延迟(从触发曝光到广告主大屏上的数字跳动)通常在秒级。但需要警惕的是,链路上的任何风吹草动——无论是 Flink 的算子反压、Kafka 的消费积压,还是 OLAP 内部由于写入过碎导致的 compaction 积压,都会瞬间把所谓「实时仓」降级为准实时。
6.1 Schema 演进如何做到对报表透明
上游业务字段的频繁变更是 AdTech 的家常便饭。为了不让下游频繁崩溃,必须设立强硬的演进规约:
- ODS 宽容接入:采用具备向后兼容性的 JSON 扩展机制,允许新字段在尾部自由追加;
- DWD 显式列控:严禁随意覆盖旧口径。新指标必须开辟新列(如
ctr_v2),废弃列打上标记但不做物理删除,依赖底层列存引擎的跳过机制来规避 I/O 损耗; - ADS 严格白名单:下游 API 查询必须显式声明列名,全面封杀
SELECT *; - 破坏性变更闭环:遇到不兼容的重构,必须老老实实走双写过渡期、灰度验证并原位切换读指针。
更重要的是,流作业的字段增添与 OLAP 引擎的 DDL 变更,必须被纳入同一个发版流水线。否则,极易重演「Flink 开始往新列写数,而报表库还没来得及加列」这类深夜爆雷事故。
6.2 多租户与资源隔离的底线
一条成熟的管线绝不可能只伺候一种查询:外部广告主想要高优的秒级报表,内部运营在跑任意维度的探索分析,算法侧在疯狂拉取历史样本回灌,财务团队则在月末发起严苛的对账。如果把他们全揉在一个未隔离的池子里:运营随手敲下的一句全表大扫描冷路径 SQL,就足以让整个集群的 I/O 陷入瘫痪,进而导致大盘实时 pacing 报表完全刷新不出来。
因此,实施硬隔离是生存底线:查询队列必须分级,在物理层面拆分在线 Serving 与离线计算的资源池,并严格依靠对象存储或湖格式实施冷热数据分离。
7. 分层反模式(广告仓常见重灾区)
- 报表穿透直连 ODS:产品端直接去解析带有嵌套层级的原始 JSON,上游任何一次格式重构都会让报表全局雪崩。
- 各业务线私搭准 DWD:导致同一次大促曝光,在公司内部流传着三种截然不同的去重基数(UV),甚至连结算和合同口径都无法对齐。
- DWS 盲目笛卡尔积:为了提前把所有查询算好,强行按几十个维度做预聚合,其带来的维度组合爆炸使得写入成本甚至远超现场扫明细的代价。
- 口径脱节:速度层坚信「低延迟最重要,口径以后再修」,历史经验一再证明,一旦速度层与批层结果未收敛,以后就永远不可能再对账成功。
- 伪 Kappa 自欺欺人:Kafka 仅保留 24 小时,口径写错时根本无力回天,却对外标榜公司全面拥抱了单逻辑架构。
8. 何时仍应果断保留批层
推崇流批一体与 Kappa 理念,绝不是教条式地要求斩杀所有 Spark 或批处理作业。在某些严苛的业务边界处,批处理依然拥有不可动摇的地位:
- 超大规模历史重算:当面临数月甚至跨年的数据整体回刷时,直接起批任务扫描数据湖,其吞吐与成本效益远好于让流引擎慢吞吞地消化回放队列;
- 复杂离线挖掘:涉及全站图谱推演或算法大模型训练的样本构建,其本质不属于在线分析链路,理应交给批计算;
- 强监管日终快照:当财务审计强制要求提交「某月某日零点的绝对不可变分区」时,我们完全可以复用流批一体的 SQL,在日终以有界流模式执行一遍大批处理,产出审计认可的冻结数据。
这里的核心原则依然清晰:批层绝不允许再去重新发明一套指标语言。 它仅仅是那套唯一被认证的业务口径,在面对海量吞吐或审计要求时,切换出的另一种物理执行模式而已。
参考
- 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 引擎地基。