Skip to content
Charles Shao
Go back

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

views

上一篇讲清了 列存与 ClickHouse/Doris——Serving 侧有了能秒查的引擎。但广告实时数据从竞价、曝光链路涌进来时,先落哪、洗净到什么程度、汇总到哪一层再给报表,才是数仓真正决定「能不能长期演进」的部分。本篇把 ODS→ADS 分层Lambda vs Kappa、以及 Flink + OLAP 的流批一体实时仓串成一张可落地的图纸。

本文是实时数仓与 OLAP 系列的第 2 篇。 全系列 3 篇:

  1. OLAP(开篇)· 列式存储原理与 ClickHouse/Doris
  2. 实时数仓 · 分层设计与 Lambda/Kappa 架构(本篇)
  3. 广告多维报表与漏斗分析实战

一句话定位:实时数仓不是「把批仓换成流」这么简单,而是用同一套分层语义,让流式增量与可重放回刷共享口径;Lambda 用双链路换新鲜度与正确性,Kappa / 流批一体则押注「一份逻辑 + 可重放日志」。

TL;DR

Table of contents

Open Table of contents

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

广告链路的数据源天然混杂:

若报表直接查「原始 Kafka 落地表」,三个月后一定会出现:字段含义漂、重复曝光未去重、维表变更无历史、同一个 CTR 有三种算法。分层的价值不是仪式,而是隔离变化——源换格式只动 ODS;口径变更收口在 DWD/DWS;产品改版只动 ADS。

2. ODS / DWD / DWS / ADS:各层写什么

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

四层像漏斗:越往下越「产品化」,越往上越「保真可重放」。

2.1 ODS · 操作数据层

落地介质可以是对象存储上的原始日志、Kafka 长期保留、或 OLAP/湖上的明细贴源表——关键是语义贴源,不是磁盘形态。

2.2 DWD · 明细数据层

这里才开始「数据工程」:

DWD 是全公司共享的「干净事实」——漏斗、归因、计费核对都应能回到这一层。

DWD 质量门禁(建议写进契约)

上线一份「DWD 准入检查」比事后对口径便宜得多:

没有门禁的 DWD,本质上还是 ODS——只是换了张漂亮的表名。

2.3 DWS · 汇总数据层

2.4 ADS · 应用数据层

经验法则: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?因为当时流引擎的正确性/状态容错还没今天稳,用批保证账单级正确,用流给广告主一个「先看看」的实时数。但代价也很硬:

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

Kappa 砍掉批层作为「正确性权威」的地位:

前提是:流引擎要扛得住状态与 exactly-once(Flink Checkpoint端到端 Exactly-Once),日志要真的能重放。对 AdTech,Kappa 的诱惑极大——指标定义只有一处,实时看板和事后对账终于能说同一种话。

实践中很少「纯」Kappa:仍会保留低频批核对(抽检、对账作业)或冷数据的批量回刷,但它们不再是产品查询的主路径。这更接近今天常说的流批一体

4.1 上 Kappa 的硬前提(缺少任一条都会伪 Kappa)

  1. 日志真能重放:Kafka 保留或同步进湖的窗口 ≥ 口径发版回滚所需天数;
  2. 作业可从指定位点重启:状态与 sink 幂等策略明确(Exactly-Once);
  3. 口径有版本号:重放输出写到新版本路径/表,Serving 切换原子;
  4. 成本可承受:全量重放七天明细不是免费的,要有容量规划与低峰窗口。

「我们消息队列里只有 24 小时,出事靠人肉补 CSV」——那不叫 Kappa,叫运气。

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

Flink 开篇已强调:批是有界流。落到数仓:

模式输入用途
Kafka 无界实时 DWD/DWS/ADS
历史分区 / 回放口径变更重算、补数、对账
同一 SQL/作业仅执行模式不同消灭双逻辑

运维上建议:口径变更走「发版 + 重放 + 切换 Serving 指针」,而不是「先改速度层、批层下个月再说」——那是 Lambda 漂移的复发路径。

5.1 湖仓一体怎么放进这张图

不少团队会问:Iceberg/Hudi/Paimon 跟 ODS–ADS 什么关系?简洁看法:

不要用湖替代所有 OLAP 秒查,也不要让 OLAP 存「永远不查的冷历史」。分层同样适用于存储介质:热在线、温可查、冷可重放

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

主路径管新鲜度与分层落点;维表广播管对齐;Exactly-Once 与物化管「敢写敢查」。

一条可生产的参考链路:

  1. Kafka 承载曝光/点击/转化(分区与并行度对齐见 Flink 篇);
  2. Flink 做 DWD:解析、过滤、去重、维表 join(Broadcast State)、窗口轻度聚合;
  3. 写出到 OLAP(ClickHouse / Doris)的 DWD/DWS 表,或先入湖再灌;
  4. OLAP 物化 / 投影生成 ADS;
  5. 报表 API 只查 ADS,长尾 dig 受限流与隔离队列访问 DWD。

新鲜度 SLA 要写进契约:例如「事件进入 Kafka 到 ADS 可查 P99 ≤ 30s」。超时告警应区分:消费 lag作业反压OLAP 写入/merge 积压——三者修法完全不同。

6.1 Schema 演进怎么不伤报表

广告事件字段几乎每月都在加。推荐约定:

Flink SQL / 表 schema 与 OLAP DDL 最好进同一变更评审,避免「流作业已写新字段、报表库还没加列」的深夜事故。

6.2 多租户与资源隔离

一条实时仓往往会同时服务:广告主看板、内部运营 dig、算法离线样本回流、财务对账。若不做隔离:

最低配置:查询队列分级(看板 / dig / 回刷)、作业资源池拆分对象存储冷热分层。这和列存选型同等重要。

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

8. 何时还该保留批层

不是教条消灭 Spark/批。以下场景批仍然香:

关键是:批不再偷偷定义另一套指标语言。 它是同一套语义在另一种执行模式下的运行。

参考


views
Share this post on:

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