Skip to content
Charles Shao
Go back

实时数仓 · 分层设计与 Lambda/Kappa 架构

–views

上一篇讲清了 列存与 ClickHouse/Doris——Serving 侧有了能秒查的引擎地基。然而,当广告主质问「为什么同一个 CTR 跑出三种算法结果」「看板怎么卡在十几分钟前」或是「结算 UV 为何与报表对不上」时,单靠底层引擎的扫列速度已无济于事。当实时数据从竞价与曝光链路汹涌而入,数据先落哪、洗净到什么程度、在何处固化口径,才是数仓能否长期演进的决定性命题。本篇将把 ODS 到 ADS 的分层设计、Lambda 与 Kappa 架构的取舍,以及 Flink + OLAP 的流批一体端到端管线串联成一张可落地的图纸。

本文是 实时数仓与 OLAP 系列的第 2 篇(分层与架构)。 全系列 4 篇:

  1. OLAP · 列式存储原理与 ClickHouse/Doris
  2. 实时数仓 · 分层设计与 Lambda/Kappa 架构(本篇)
  3. 广告多维报表与漏斗分析实战
  4. 时序数据库 InfluxDB 深挖 · 数据模型、TSM 与降采样

一句话定位:实时数仓绝非简单地将批处理替换为流计算,而是通过同一套分层语义让流式增量与可重放回刷共享唯一的指标口径;Lambda 依靠双链路换取新鲜度,而 Kappa 与流批一体则押注于「一份逻辑 + 可重放日志」。

TL;DR

Table of contents

Open Table of contents

1. 为什么实时仓必须分层

在广告业务中,数据源天然呈现出极度混杂的态势:实时事件(如竞价、曝光、点击与转化)以极高的吞吐直接流入 Kafka;离线 Dump 被用于延迟补传与财务对账;而广告主、计划等维表状态则依赖 CDC 同步。

如果放任前端报表直接去查询未经处理的原始 Kafka 落地表,不出三个月系统必定走向失控。不仅上游 JSON 字段的随意增删会让解析逻辑散落各处,重复曝光未被剔除、状态变更丢失历史轨迹等问题,更是会让同一份报表跑出截然不同的口径。正因如此,分层的核心价值并非为了追求架构上的仪式感,而是为了隔离变化。通过分层,当数据源格式更迭时只需调整 ODS 的解析逻辑,复杂的聚合与去重被严格收拢在 DWD 与 DWS,而前端看板的改版仅需在 ADS 层做轻量适配。

2. ODS / DWD / DWS / ADS:各层职责划分

实时数仓分层示意图。顶部数据源:Kafka 实时事件、对象存储日志/Dump、业务库 Binlog CDC。向下四层:ODS 操作数据层(贴源可重放、少清洗、按事件日分区)→ DWD 明细数据层(去噪去重对齐维表、统一口径、曝光点击转化事实表)→ DWS 汇总数据层(campaign×媒体×小时轻度预聚合、服务多下游)→ ADS 应用数据层(广告主看板、漏斗留存专题、物化视图/宽表)。底部提醒流批一体可共享分层语义,禁止报表直扫 ODS。

分层架构的语义递进:越靠近 ADS 越产品化,越靠近 ODS 越贴源可重放。

2.1 ODS · 操作数据层

作为实时仓的入口,ODS(Operational Data Store)的铁律是贴源存储。字段应尽可能原样保留最初的 JSON 或 Avro 结构,仅做最基础的时间戳提取与按日分区。这一层的根本使命在于保障可重放能力:一旦下游发生大面积的口径事故,系统必须能够直接切回 ODS(或更上游的 Kafka),从原始形态重新推演并回刷整个下游链路。因此,严禁在 ODS 层「发明」业务逻辑或硬编码归因窗口。无论是存放在对象存储、长周期保留的 Kafka,还是数据湖,只要语义贴源,它就是合格的 ODS。

2.2 DWD · 明细数据层

数据工程的脏活与重头戏,往往全部压在 DWD(Data Warehouse Detail)这一层:

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 架构给出了一个看似完美的妥协:

  1. Batch Layer(批处理层):定期运行全量作业,产出绝对正确但缺乏新鲜度的历史视图;
  2. Speed Layer(速度层):用流处理引擎消费增量数据,快速提供低延迟的近似估算;
  3. Serving Layer(服务层):在查询发生时,由引擎当场合并这两份结果。

Lambda 与 Kappa 对比。左侧 Lambda:原始事件分叉为 Speed Layer(Flink 实时、近似低延迟)与 Batch Layer(Spark/批、全量准结果),汇入 Serving Layer 合并视图,强调两套逻辑要对齐。右侧 Kappa:Kafka 可重放日志进入唯一计算引擎 Flink(实时与历史回放同逻辑,改口径可重放修数),再进 Serving/OLAP 一张结果一套口径。

Lambda 依靠双链路换取低延迟与正确性;Kappa 则押注可重放日志与单逻辑以换取可维护性。

早期的广告系统之所以对 Lambda 趋之若鹜,是因为当时的流计算在状态容错与 Exactly-Once 语义上尚不堪大用。团队只能靠批层来兜底账单级别的真相,同时用速度层勉强撑起大屏。然而这种架构的反噬极快:双重逻辑维护意味着你必须用两种不同的 API(如 Storm 与 Hadoop)写出绝对对齐的 CTR 算法,稍有偏差便会导致口径漂移;同时,Serving 层的合并逻辑异常繁琐,必须精准裁定批数据的生效边界并截断对应的流数据;随着指标日益繁杂,人力对账与排障成本很快就会失控。

4. Kappa:一份逻辑,靠重放修数

为了根除双写逻辑的梦魇,Kappa 架构果断剥夺了批层作为「唯一真相源」的特权:

Kappa 成立的前提,在于底层流引擎必须具备极其强悍的容错能力(如 Flink Checkpoint 与 端到端 Exactly-Once 机制)。这对 AdTech 团队的诱惑是致命的——它终于把复杂的指标定义收拢到了一处,让实时看板与月末结算真正实现了口径的统一。

4.1 上 Kappa 的硬前提

现实中极少有纯粹的 Kappa,很多团队跟风宣称上了 Kappa,却连基础的容灾底座都不具备。要真正发挥单逻辑重放的威力,以下条件缺一不可(缺少任一条即为伪 Kappa):

  1. 日志具备真实重放能力:底层 Kafka 的保留周期(或入湖窗口期)必须长于口径发版回滚所需的最大跨度;
  2. 作业支持精准位点重启:引擎的状态管理与下游 OLAP 的幂等写入必须严丝合缝;
  3. 指标定义版本化:重放时必须将新口径的数据写至带版本号的新表或新分区,且 Serving 层的读指针切换必须具备原子性;
  4. 重放成本可控:回刷数月海量明细并非轻而易举,必须预估好引擎在追赶期引发的计算资源抢占。

如果你的消息队列只保留 24 小时,出事全靠临时导 CSV 补数,这就不要自称为 Kappa 架构。

5. 流批一体:批是有界流

正如 Flink 系列开篇所强调的核心理念:批处理本质上不过是拥有明确终点的「有界流」。这一理念映射到分层架构中,就是当今被广泛讨论的流批一体。

模式输入用途
流计算无界的 Kafka 增量流构建毫秒或秒级可见的实时 DWD / DWS
批处理历史分区或湖中归档执行口径变更重算、批量补数与深度对账

回到刚才那条宽表聚合的痛点,无论是实时流式消费还是面向历史数据的离线批处理,现在都能够复用同一套 SQL 与算子逻辑。在运维层面,这意味着任何口径变更都可以被收敛为「发布新作业 → 历史重放补齐 → 原子切换前端读指针」的标准闭环,从而彻底消灭了「先上线速度层,以后再跟批层对齐」这种自欺欺人的开发习惯。

5.1 湖仓一体如何融入分层架构

许多团队在做引擎选型时常有疑惑:像 Iceberg、Hudi 这样的数据湖,与 ODS 到 ADS 的分层到底如何适配?最务实的定位是让系统各司其职:

切忌试图用高昂的 OLAP 内存去硬抗「永远不会被查询的冷历史」,也别指望用数据湖去包揽所有广告主报表的亚秒级响应。分层不仅是业务语义的递进,同样也是介质成本的规划:热数据保持在线,温数据随时可查,冷数据低成本可重放。

Flink 与 OLAP 实时仓端到端管线。主路径:Kafka 曝光点击 → Flink 清洗对齐聚合 → DWD/DWS 流式写或 upsert → OLAP(ClickHouse/Doris)→ 报表 API/看板。旁路 A:维表 CDC(广告主/创意)Broadcast 进 Flink 对齐。旁路 B:Exactly-Once Checkpoint 写稳,物化/预聚合命中报表。

主路径决定了数据的分层落点;维表广播保障了维度对齐;Exactly-Once 语义与预聚合则赋予了下游「敢写敢查」的底气。

将上述理论拼装,一条经受过大流量考验的生产管线通常如下运转:

  1. 事件接入:海量的竞价、曝光与点击汇入 Kafka,以分区机制保障核心链路的并发吞吐;
  2. 流式清洗:Flink 承接 ODS 数据,执行过滤、基于状态的精确去重,通过 Broadcast State 完成广告主维表 Join,最后按时间窗口执行轻度聚合,将散乱的日志转化为高度标准化的 DWD/DWS 事实;
  3. 高效落盘:将清洗对齐后的增量数据,通过追加或覆盖写(upsert)稳稳推入 ClickHouse 或 Doris 这类 OLAP 引擎;
  4. 视图加速:在 OLAP 内部,依托后台的 merge 机制与物化视图,进一步固化出专供报表使用的 ADS 层;
  5. 查询隔离:前端业务的 报表 API 仅被授权查询 ADS 层,而任何面向明细溯源的长尾 dig 查询,都必须被赶入独立的隔离队列中去慢查 DWD。

端到端可见延迟(从触发曝光到广告主大屏上的数字跳动)通常在秒级。但需要警惕的是,链路上的任何风吹草动——无论是 Flink 的算子反压、Kafka 的消费积压,还是 OLAP 内部由于写入过碎导致的 compaction 积压,都会瞬间把所谓「实时仓」降级为准实时。

6.1 Schema 演进如何做到对报表透明

上游业务字段的频繁变更是 AdTech 的家常便饭。为了不让下游频繁崩溃,必须设立强硬的演进规约:

更重要的是,流作业的字段增添与 OLAP 引擎的 DDL 变更,必须被纳入同一个发版流水线。否则,极易重演「Flink 开始往新列写数,而报表库还没来得及加列」这类深夜爆雷事故。

6.2 多租户与资源隔离的底线

一条成熟的管线绝不可能只伺候一种查询:外部广告主想要高优的秒级报表,内部运营在跑任意维度的探索分析,算法侧在疯狂拉取历史样本回灌,财务团队则在月末发起严苛的对账。如果把他们全揉在一个未隔离的池子里:运营随手敲下的一句全表大扫描冷路径 SQL,就足以让整个集群的 I/O 陷入瘫痪,进而导致大盘实时 pacing 报表完全刷新不出来。

因此,实施硬隔离是生存底线:查询队列必须分级,在物理层面拆分在线 Serving 与离线计算的资源池,并严格依靠对象存储或湖格式实施冷热数据分离。

7. 分层反模式(广告仓常见重灾区)

8. 何时仍应果断保留批层

推崇流批一体与 Kappa 理念,绝不是教条式地要求斩杀所有 Spark 或批处理作业。在某些严苛的业务边界处,批处理依然拥有不可动摇的地位:

这里的核心原则依然清晰:批层绝不允许再去重新发明一套指标语言。 它仅仅是那套唯一被认证的业务口径,在面对海量吞吐或审计要求时,切换出的另一种物理执行模式而已。

参考


–views
Share this post on:

Previous Post
特征提取与降维:吊高汤——把一大堆料浓缩成精华
Next Post
OLAP · 列式存储原理与 ClickHouse/Doris