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
  1. 所有原始数据不可变追加存储:原始日志全部原样保存,只追加、不修改,作为唯一真相源。
  2. Batch Layer 批处理层:对全部历史完整数据集做批量计算,生成批视图(Batch View)。计算慢、延迟高,但是结果 100% 准确。
  3. Speed Layer 速度层(实时层):处理新到达的增量数据,生成实时视图(Real‑time View)。只处理最近一段时间的数据,延迟低,允许存在小误差。
  4. 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 适合做服务层存储:支持随机读,方便分别读取批结果、实时结果做合并。

三、优缺点

✅优点

  1. 数据可靠性高:原始数据永久保存,一旦逻辑 bug,可以重新跑批层任务,重算全部视图;
  2. 兼顾准确性(批)+ 低延迟(速度层)
  3. 流处理出现异常,批层作为兜底,不会数据彻底错乱。

❌核心缺点

  1. 维护两套业务逻辑代码:同一指标,批写一套,流写一套;两份代码需要同步修改,极易出现逻辑不一致,维护成本高;
  2. 存储两份结果数据,存储开销更大;
  3. 架构组件多,运维复杂。

后续衍生出 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 -.需要重算时回填.-> B

Kappa 核心设计思路

  1. 所有数据以事件流形式存入消息队列 (Kafka),原始事件永久保存(设置很长的保留时间)。
  2. 只使用同一套流处理程序,既处理新来的实时数据,也可以重放历史数据完成全量离线计算。
  3. 当业务逻辑变更、程序 Bug:重新启动流作业,从最早的位置重放 Kafka 历史事件,重新计算全部结果,覆盖旧输出
  4. 计算结果输出到外部存储,对外提供查询服务。

✔关键:不再区分专门的批层、速度层;离线、实时复用同一套代码。

Kappa 架构工作流程

  1. 业务数据采集,源源不断写入 Kafka 消息队列;
  2. 流引擎消费 Kafka,执行业务计算,输出结果到数据库 / OLAP;
  3. 逻辑变更 / 修复 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 架构优缺点

✅优点

  1. 只维护一套计算代码,解决 Lambda 架构 “批、流双逻辑不一致” 的最大痛点,开发维护成本大幅下降;
  2. 架构组件更少,架构简洁;实时计算、全量重算复用同一套逻辑。

❌缺点

  1. 高度依赖消息队列:Kafka 无法无限保存数据,存储海量多年全部事件成本很高,超期数据要归档 + 回填,增加复杂度;
  2. 全量重放历史数据会消耗大量集群 CPU、IO 资源;历史数据量巨大时重算耗时很长;
  3. 如果源数据量极大,全部走流重放做全量计算,资源开销高于传统 Spark 离线批;
  4. 对引擎要求高:流引擎必须支持状态后端、大状态管理(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,2011Jay Kreps,2014(《Questioning the Lambda Architecture》)
核心思路批处理保证准确 + 流处理补延迟一切皆流,批处理是流的特例(日志重放)
处理路径两条(批处理层 + 速度层,服务层合并)一条(消息队列 + 流引擎)
代码维护两套代码,且必须保证结果一致一套代码
存储两套(批存储 + 实时存储)一套(Kafka 日志 + 结果存储)
历史重算重跑批处理作业从 Kafka 重放日志,重新跑流作业
延迟实时视图秒级,但完整视图取决于批周期端到端流式,延迟更低
准确性批处理全量兜底,结果可精确复现依赖流引擎的恰好一次语义(exactly-once)
系统复杂度高(组件多、合并逻辑复杂)低(组件少、链路单一)
故障恢复重新跑批从日志 offset 重放
典型组件Hadoop/Spark Batch + Storm/Flink + HBase/ImpalaKafka + Flink / Kafka Streams / Pulsar
典型场景历史全量 + 实时混合、复杂离线算法、强一致报表实时监控、实时推荐、风控等以实时为主的场景

怎么选

  • 选 Lambda:要求结果绝对可复现(如财务 / 审计类统计)、需要跑复杂离线算法(机器学习训练、全量关联)、批流两套已有存量系统迁移成本高。
  • 选 Kappa:业务以实时为核心(指标看板、告警、个性化推荐)、团队想统一技术栈降低运维、日志可以完整保留且重放成本可控。
  • 现实的主流做法:用 Spark/Flink 这类流批一体引擎,一套代码同时跑批和流(本质上吸收了 Kappa 的简洁,又保留了 Lambda 的批处理兜底能力);存量 Lambda 系统则逐步把速度层和批处理层统一到同一个引擎上,消除 “双写逻辑”。