大数据中Lambda和Kappa架构
Lambda 架构
一、理论概念
Lambda 架构是由 Nathan Marz 提出的大数据混合架构,核心思想:同时使用批处理层 + 速度层(实时层),再通过服务层合并两套结果对外提供查询,用来解决:海量数据下,既要全量准确离线结果,又要低延迟实时视图。
背景:早期纯批处理延迟很高,只能 T+1;纯流处理早期很难做到完美容错、数据不丢不重。Lambda 折中:批层保证最终准确,速度层补上实时增量。
Lambda 四层结构
简化版:
flowchart LR
Source[数据源] --> Batch[批处理层<br/>生成批视图]
Source --> Speed[速度层<br/>生成实时视图]
Batch --> Serving[服务层<br/>合并两套视图]
Speed --> Serving
Serving --> App[业务应用查询]细化版本:
flowchart TD
A[数据源<br/>日志/传感器/业务事件] --> B[(原始数据集<br/>HDFS 只追加存储)]
A --> C[Kafka消息队列<br/>最新增量数据流]
%% 批处理层
subgraph BatchLayer[批处理层 Batch Layer]
B --> B1[Spark/Hive/MapReduce<br/>全量历史数据计算]
B1 --> B2[(批视图 Batch View)]
end
%% 速度层
subgraph SpeedLayer[速度层 Speed Layer]
C --> S1[Flink/Spark Streaming<br/>增量实时计算]
S1 --> S2[(实时视图 Speed View)]
end
%% 服务层
subgraph ServingLayer[服务层 Serving Layer]
B2 --> D[合并查询<br/>HBase/Phoenix]
S2 --> D
end
D --> E[对外API/大屏/业务应用]
%% 标注说明
note1["批层:全量数据,结果准确,延迟高"]
note2["速度层:仅新数据,低延迟,允许近似误差"]
note3["查询公式:最终结果 = 批视图 + 实时视图"]
note1 -.-> BatchLayer
note2 -.-> SpeedLayer
note3 -.-> ServingLayer- 所有原始数据不可变追加存储:原始日志全部原样保存,只追加、不修改,作为唯一真相源。
- Batch Layer 批处理层:对全部历史完整数据集做批量计算,生成批视图(Batch View)。计算慢、延迟高,但是结果 100% 准确。
- Speed Layer 速度层(实时层):处理新到达的增量数据,生成实时视图(Real‑time View)。只处理最近一段时间的数据,延迟低,允许存在小误差。
- Serving Layer 服务层:对外查询接口,合并批视图 + 实时视图,对外返回完整结果。
查询公式:最终结果 = 全量批处理结果 + 近期增量实时计算结果
关键特点
- 原始数据只追加,绝不修改,出错可以重跑批任务重新生成视图;
- 批层负责正确性;速度层负责低延迟;
- 同一个业务指标,需要维护两套计算逻辑代码(批代码 + 流代码),这是 Lambda 最大痛点。
二、各层对应具体技术选型
| 层级 | 作用 | 常用技术组件 |
|---|---|---|
| 原始数据存储 | 保存全部原始不可变数据集 | HDFS、MinIO (S3 对象存储) |
| Batch Layer 批处理层 | 全量历史数据离线计算生成批视图 | Hadoop MapReduce、Spark Core / Spark SQL、Hive 输出结果存 HDFS / HBase |
| Speed Layer 速度层 | 增量实时数据计算实时视图 | Flink / Spark Streaming 消息队列:Kafka 保存增量数据流 实时结果存入:HBase、Redis |
| Serving Layer 服务层 | 合并批视图、实时视图,对外提供查询 | HBase、Phoenix、自定义 API 服务,Impala |
典型技术组合:
- 原始数据:Kafka 采集 → HDFS 持久化原始日志
- 批层:Spark+Hive 读取 HDFS 全量数据,计算统计指标,写入 HBase(批视图)
- 速度层:Flink 消费 Kafka 最新消息,计算增量指标,写入 HBase 另一张表(实时视图)
- 服务层:查询时同时读取两张表,做合并聚合,返回最终业务数据。
HBase 适合做服务层存储:支持随机读,方便分别读取批结果、实时结果做合并。
三、优缺点
✅优点
- 数据可靠性高:原始数据永久保存,一旦逻辑 bug,可以重新跑批层任务,重算全部视图;
- 兼顾准确性(批)+ 低延迟(速度层);
- 流处理出现异常,批层作为兜底,不会数据彻底错乱。
❌核心缺点
- 维护两套业务逻辑代码:同一指标,批写一套,流写一套;两份代码需要同步修改,极易出现逻辑不一致,维护成本高;
- 存储两份结果数据,存储开销更大;
- 架构组件多,运维复杂。
后续衍生出 Kappa 架构:放弃批层,全部用流处理,把历史数据重放 Kafka 来实现重算,只维护一套代码,解决 Lambda 双重代码的痛点。
四、实际业务应用场景
场景 1:网站 / APP 用户访问统计
业务需求:统计页面 PV、UV。既需要昨日至今完整准确统计,又需要看到当前分钟实时访问量。
- 批层:Spark 读取 HDFS 全部访问日志,计算从项目上线至今完整 PV/UV 批视图(T+1 更新)
- 速度层:Flink 消费 Kafka 实时用户访问日志,计算最近几十分钟增量 PV/UV
- 服务层:查询 = 历史全量批结果 + 最近增量实时结果,对外大屏展示实时访问统计。
场景 2:物联网工厂设备指标统计
工厂成千上万传感器不断上报温度、压力。
- 批层:每日凌晨 Spark 对全部历史传感器数据做统计,得到设备历史完整统计报表;
- 速度层:Flink 实时消费传感器数据流,实时计算最近 1 小时设备指标,用于监控告警;
- 服务层合并输出:既有完整准确历史统计,又有秒级实时监控。
场景 3:电商交易统计
统计商品总成交金额。
- 批层:离线计算全部历史订单,得到商品累计成交(每日更新);
- 速度层:实时消费订单消息,计算最近几十分钟新增成交;
- 对外查询:历史批结果叠加实时增量,页面可以看到近乎实时的商品销售额。
Kappa 架构
一、理论概念
Kappa 架构是 Jay Kreps 针对Lambda 架构双份代码维护痛点提出的改进大数据架构。
Lambda 最大问题:同一业务逻辑,需要维护批处理、流处理两套代码,逻辑容易不一致,维护成本高。
Kappa 核心思想:移除独立的批处理层,一切计算全部基于流处理引擎实现,只保留一套业务逻辑。把所有历史数据当作 “过去的流”,通过消息队列重放历史数据,完成离线全量重计算。
flowchart LR
A[业务数据源] --> B[Kafka消息队列<br/>保存全部原始事件]
B --> C[Flink流计算引擎<br/>同一套代码<br/>实时消费 / 历史重放]
C --> D[结果存储<br/>ClickHouse / StarRocks / HBase]
D --> E["业务查询、大屏、报表"]
subgraph 重算流程
F[业务逻辑变更/Bug修复] --> G[清空输出存储]
G --> H["Flink重置offset,重放Kafka全部历史数据"]
H --> C
end
B -.超期数据归档.-> I[HDFS/MinIO对象存储]
I -.需要重算时回填.-> BKappa 核心设计思路
- 所有数据以事件流形式存入消息队列 (Kafka),原始事件永久保存(设置很长的保留时间)。
- 只使用同一套流处理程序,既处理新来的实时数据,也可以重放历史数据完成全量离线计算。
- 当业务逻辑变更、程序 Bug:重新启动流作业,从最早的位置重放 Kafka 历史事件,重新计算全部结果,覆盖旧输出。
- 计算结果输出到外部存储,对外提供查询服务。
✔关键:不再区分专门的批层、速度层;离线、实时复用同一套代码。
Kappa 架构工作流程
- 业务数据采集,源源不断写入 Kafka 消息队列;
- 流引擎消费 Kafka,执行业务计算,输出结果到数据库 / OLAP;
- 逻辑变更 / 修复 bug:
- ①清空输出存储里面旧计算结果;
- ②启动同一个流作业,调整消费偏移量,从 Kafka 最早位置重放全部历史数据;
- ③全量重算完成后,新结果对外提供服务。
前提条件:Kafka 可以保存很长周期的历史事件(或者将海量历史数据回填回 Kafka)。
二、分层与技术选型
| 模块 | 作用 | 主流技术 |
|---|---|---|
| 事件源 & 消息存储 | 保存全部原始事件流,支持重放 | Kafka(核心);海量冷数据可以下沉 HDFS/MinIO 做归档,需要时回填 Kafka |
| 流计算引擎 | 唯一一套代码,同时支持实时消费、历史重放计算 | Apache Flink(首选);Spark Streaming |
| 结果输出存储 | 保存计算完成指标,对外查询 | Redis、HBase、ClickHouse、StarRocks、Doris |
| 数据归档 | Kafka 保存周期有限,超期原始事件归档 | HDFS、MinIO 对象存储 |
典型完整技术组合:
业务数据 → Kafka → Flink 流计算(同一套代码) → ClickHouse/StarRocks 对外查询。
Kafka 保存周期有限,过期数据归档到 HDFS;重算需要全量历史时,把归档数据回填进 Kafka 再重放。
三、Kappa 架构优缺点
✅优点
- 只维护一套计算代码,解决 Lambda 架构 “批、流双逻辑不一致” 的最大痛点,开发维护成本大幅下降;
- 架构组件更少,架构简洁;实时计算、全量重算复用同一套逻辑。
❌缺点
- 高度依赖消息队列:Kafka 无法无限保存数据,存储海量多年全部事件成本很高,超期数据要归档 + 回填,增加复杂度;
- 全量重放历史数据会消耗大量集群 CPU、IO 资源;历史数据量巨大时重算耗时很长;
- 如果源数据量极大,全部走流重放做全量计算,资源开销高于传统 Spark 离线批;
- 对引擎要求高:流引擎必须支持状态后端、大状态管理(Flink 状态)。
对比 Lambda:
- Lambda:批层保证准确,流层负责低延迟;两套代码。
- Kappa:流引擎一统天下,一套代码;依靠消息重放实现全量重计算。
四、实际应用场景
场景 1:电商实时交易数仓、实时大屏
需求:统计订单交易额、订单量,既要实时看到成交,业务逻辑变更后需要重算全部历史订单。
- 所有订单事件写入 Kafka;
- Flink 一套代码消费 Kafka,实时计算各维度交易指标写入 StarRocks;
- 当统计口径修改、修复 bug:清空目标表,Flink 从 Kafka 起始 offset 重放全部订单事件,全量重新计算。
场景 2:物联网设备监控分析
工厂传感器上报设备温度、振动等时序事件。
- 传感器数据全部上报 Kafka;Flink 实时计算设备告警、聚合指标,写入 ClickHouse。
- 修改指标计算规则:直接重放 Kafka 内历史传感器事件,重新生成全部统计结果,不需要单独写 Spark 批任务。
场景 3:用户行为分析
APP 埋点用户点击、浏览行为日志。
- 用户行为事件进入 Kafka,Flink 做 PV、UV、用户画像标签计算;
- 标签规则调整,直接重放历史行为事件,重新计算全量用户标签。
场景 4:实时数仓(ODS→DWD→DWS)
实时数仓建设,整套链路全部 Flink 流作业,Kafka 作为各层中间存储;口径变更直接重放消息完成全量刷新。
⚠不适合场景:
多年超大规模历史数据,无法全部存入 Kafka,重放代价极高,这种场景 Lambda 或者湖仓一体(Iceberg/Hudi)会更合适。
Lambda架构和Kappa架构的对比
Lambda 用 “批处理兜底准确 + 流处理补实时” 双路径,代价是两套代码两套存储;Kappa 用 “一切皆流 + 日志可重放” 单路径,代价是历史重算依赖流引擎能力和日志保留成本。两者本质是准确性 vs 简单性的权衡,近年流批一体引擎让 Kappa 越来越主流。
| 维度 | Lambda 架构 | Kappa 架构 |
|---|---|---|
| 提出者 / 时间 | Nathan Marz,2011 | Jay Kreps,2014(《Questioning the Lambda Architecture》) |
| 核心思路 | 批处理保证准确 + 流处理补延迟 | 一切皆流,批处理是流的特例(日志重放) |
| 处理路径 | 两条(批处理层 + 速度层,服务层合并) | 一条(消息队列 + 流引擎) |
| 代码维护 | 两套代码,且必须保证结果一致 | 一套代码 |
| 存储 | 两套(批存储 + 实时存储) | 一套(Kafka 日志 + 结果存储) |
| 历史重算 | 重跑批处理作业 | 从 Kafka 重放日志,重新跑流作业 |
| 延迟 | 实时视图秒级,但完整视图取决于批周期 | 端到端流式,延迟更低 |
| 准确性 | 批处理全量兜底,结果可精确复现 | 依赖流引擎的恰好一次语义(exactly-once) |
| 系统复杂度 | 高(组件多、合并逻辑复杂) | 低(组件少、链路单一) |
| 故障恢复 | 重新跑批 | 从日志 offset 重放 |
| 典型组件 | Hadoop/Spark Batch + Storm/Flink + HBase/Impala | Kafka + Flink / Kafka Streams / Pulsar |
| 典型场景 | 历史全量 + 实时混合、复杂离线算法、强一致报表 | 实时监控、实时推荐、风控等以实时为主的场景 |
怎么选
- 选 Lambda:要求结果绝对可复现(如财务 / 审计类统计)、需要跑复杂离线算法(机器学习训练、全量关联)、批流两套已有存量系统迁移成本高。
- 选 Kappa:业务以实时为核心(指标看板、告警、个性化推荐)、团队想统一技术栈降低运维、日志可以完整保留且重放成本可控。
- 现实的主流做法:用 Spark/Flink 这类流批一体引擎,一套代码同时跑批和流(本质上吸收了 Kappa 的简洁,又保留了 Lambda 的批处理兜底能力);存量 Lambda 系统则逐步把速度层和批处理层统一到同一个引擎上,消除 “双写逻辑”。

