专栏 知识宝典 子专栏 数据与存储 14 篇

4.3.1 Lambda / Kappa 架构选型 + 数据湖三选一(Hudi / Iceberg / Delta Lake)

大数据架构全景 —— Lambda vs Kappa 架构 + 数据湖 Hudi / Iceberg / Delta Lake 三选一对比

1. 为什么这个专题重要 —— 2026 年还在讲大数据架构?

很多人以为「云原生 + 数仓 SaaS 化」时代,大数据架构已经过时,事实恰恰相反。2026 年企业面对的数据规模出现两个新变化:

  1. 实时数据爆炸:一家中等规模的电商每天产生 1-5 PB 的用户行为日志、订单变更、IoT 设备数据。Uber 每天处理超过 100 万亿条消息,Netflix 每天处理 PB 级的播放事件,字节跳动的实时推荐链路峰值超过 2 亿 QPS。
  2. 流批一体的必然性:业务方已经无法接受 T+1 的报表延迟 —— 营销要「过去 5 分钟的 GMV」、风控要「过去 30 秒的欺诈交易」、运营要「当前在线用户画像」。Lambda 架构的双链路成本太高,Kappa + 数据湖成为主流选择。

真实案例:Netflix 在 2018 年开始把 S3 上的 Parquet 表迁移到 Apache Iceberg,2023 年已经管理 EB 级数据湖,Schema 变更不再破坏下游作业;Uber 内部用 Apache Hudi 处理增量 upsert,使订单状态更新的端到端延迟从分钟级降到秒级;美团的实时数据仓库从 Lambda 迁到 Kappa + Iceberg 之后,整体资源消耗下降 50%,作业数从 800 降到 350。

参考资料:Nathan Marz 在 2011 年提出 Lambda 架构,LinkedIn 的 Jay Kreps 在 2014 年提出 Kappa 架构(《Questioning the Lambda Architecture》),Apache Hudi 由 Uber 在 2017 年开源,Apache Iceberg 由 Netflix 在 2018 年开源,Delta Lake 由 Databricks 在 2019 年开源。

2. 大数据架构演进史

2.1 演进时间线(ASCII 图)

1990s               2000s              2010s               2020s
传统数仓 ──────► 大数据 ──────► Lambda ──────► Kappa + 数据湖
(Inmon/Kimball)    (Hadoop)        (双链路)        (流批一体)

【判断节点】
★ 2006:Hadoop 诞生(HDFS + MapReduce)
★ 2011:Nathan Marz 提出 Lambda 架构
★ 2014:Spark 崛起,批处理从 MR 转向 Spark
★ 2014:Jay Kreps 提出 Kappa 架构
★ 2016:Flink 1.0 发布,流处理成熟
★ 2017:Apache Hudi 由 Uber 开源
★ 2018:Apache Iceberg 由 Netflix 开源
★ 2019:Delta Lake 由 Databricks 开源
★ 2021:Lakehouse(湖仓一体)概念成熟
★ 2024:数据湖三件套 Hudi/Iceberg/Delta Lake 全面替代裸 Parquet

2.2 三代架构对比

维度 传统数仓 Lambda 架构 Kappa + 数据湖
延迟 T+1 天 批 T+1 / 流秒级 秒级统一
成本 中(商用 MPP) 高(双链路) 中(存储便宜)
复杂度 低 高(双份代码) 中
适用 BI 报表 实时 + 离线并存 流批一体
代表 Teradata / Oracle Twitter / LinkedIn Netflix / Uber / 美团

参考资料:《Designing Data-Intensive Applications》(Martin Kleppmann, 2017)、Databricks 的 Lakehouse 论文(2021)。

3. Lambda 架构详解

3.1 三层架构

Lambda 架构由 Nathan Marz 在 2011 年提出,核心思想是「用两套独立链路保证数据的完整性和实时性」:

flowchart TB
    subgraph SL["Serving Layer<br/>合并 Batch View + Real-time View,对外查询接口"]
        direction TB
    end
    subgraph BL["Batch Layer<br/>Hadoop / Spark<br/>离线批处理"]
        direction TB
    end
    subgraph SPL["Speed Layer<br/>Storm / Flink<br/>实时流处理"]
        direction TB
    end
    SRC["Source<br/>原始数据,Kafka / HDFS"]

    SL --> BL
    SL --> SPL
    BL --> SRC
    SPL --> SRC

    classDef layerStyle fill:#e8f4ff,stroke:#2b6cb0,color:#1a365d
    classDef sourceStyle fill:#fff4e0,stroke:#c05621,color:#7b341e
    class SL,BL,SPL layerStyle
    class SRC sourceStyle

3.2 各层职责

  • Batch Layer(批处理层):处理全量历史数据,保证数据完整性和最终一致性。用 Hadoop / Spark 跑 T+1 离线任务,产出 Batch View。
  • Speed Layer(实时层):处理增量实时数据,用 Storm / Flink 维护 Real-time View,延迟通常在秒级。
  • Serving Layer(服务层):合并 Batch View 和 Real-time View,对外提供统一查询接口(通常用 HBase / Cassandra / Elasticsearch 存储)。

3.3 Spark 实现 Batch Layer

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("LambdaBatchLayer") \
    .getOrCreate()

# 读取 HDFS 上的全量原始数据
raw_df = spark.read.parquet("hdfs://namenode:9000/data/raw/events/")

# 离线 ETL + 聚合
batch_view = raw_df.filter("event_date = '2026-07-05'") \
    .groupBy("user_id", "event_type") \
    .agg({"duration": "sum", "value": "avg"}) \
    .withColumnRenamed("sum(duration)", "total_duration") \
    .withColumnRenamed("avg(value)", "avg_value")

# 写入 Batch View(HBase / Cassandra / Hive)
batch_view.write \
    .format("org.apache.spark.sql.cassandra") \
    .mode("overwrite") \
    .option("keyspace", "analytics") \
    .option("table", "batch_view_user_event") \
    .save()
// Flink 实时聚合 user_event 流
DataStream<UserEvent> events = env
    .addSource(new FlinkKafkaConsumer<>("user_events", new UserEventDeserializer(), props))
    .keyBy(UserEvent::getUserId)
    .window(TumblingEventTimeWindows.of(Time.minutes(1)));

events.aggregate(new UserEventAggregator())
    .keyBy(UserEventAgg::getUserId)
    .addSink(new CassandraSink<>("real_time_view_user_event"));

3.5 Lambda 架构优缺点

优点:

  • 容错性强:批处理层是最终一致性的兜底,流处理层挂掉不影响数据完整性。
  • 数据准确性高:批处理层能修正流处理层的误差(乱序、丢数据)。
  • 适合初期实时 + 离线并存的场景。

缺点:

  • 双链路维护成本:同一份业务逻辑要在批和流两套代码里实现,改一处要改两处。
  • 数据不一致:批流可能产生不同的统计结果(如实时 GMV 和离线 GMV 对不上)。
  • 资源浪费:批处理层通常跑全量数据,资源消耗大。

4. Kappa 架构详解

4.1 核心思想

Kappa 架构由 LinkedIn 的 Jay Kreps 在 2014 年提出,核心思想是「只保留流处理层,用 Kafka 重放历史数据替代批处理层」。所有数据(包括历史数据)都走流处理链路。

flowchart TB
    SPL["Stream Processing Layer<br/>Flink / Kafka Streams<br/>统一流处理逻辑"]
    KAFKA["Kafka<br/>原始<br/>长期保留"]
    SS["Serving Store<br/>实时查询"]

    SPL --> KAFKA
    SPL --> SS

    classDef layerStyle fill:#e8f4ff,stroke:#2b6cb0,color:#1a365d
    class SPL layerStyle

4.2 工作流程

  1. 所有数据写入 Kafka,Kafka 配置长保留期(7 天 / 30 天 / 1 年)。
  2. Flink 消费 Kafka 做实时聚合,结果写入 Serving Store(如 Cassandra / Elasticsearch)。
  3. 当业务逻辑变更时,重置 Kafka offset 到任意时间点,重新消费历史数据,得到新版本的 Serving View。
  4. 新旧版本切换时,服务层同时查新旧两套视图,验证一致后切换流量。
# Flink Kafka Consumer 配置(支持历史重放)
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.connectors import FlinkKafkaConsumer
from pyflink.common.serialization import SimpleStringSchema

env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(4)

# 配置 Kafka Consumer
kafka_props = {
    'bootstrap.servers': 'kafka-broker:9092',
    'group.id': 'kappa-consumer',
    'auto.offset.reset': 'earliest',  # 关键:从最早开始消费
    'enable.auto.commit': 'false',
}

consumer = FlinkKafkaConsumer(
    topics='user_events',
    deserialization_schema=SimpleStringSchema(),
    properties=kafka_props,
)

# 设置从具体时间点重放(业务逻辑变更时使用)
# consumer.set_start_from_timestamp(1717622400000)  # 2026-06-06 00:00:00 UTC

stream = env.add_source(consumer)

# 流处理逻辑(批流同套代码)
result = stream \
    .map(lambda x: parse_event(x)) \
    .key_by(lambda e: e['user_id']) \
    .window(TumblingEventTimeWindows.of(Time.minutes(1))) \
    .aggregate(UserEventAggregator())

result.add_sink(CassandraSink("serving_view_user_event"))
env.execute("KappaArchitectureJob")

4.4 Kafka 配置(长期保留)

# server.properties
log.retention.hours=8760           # 1 年保留
log.retention.bytes=10995116277760 # 10 TB 单分区上限
log.segment.bytes=1073741824       # 1 GB 单 segment
log.cleanup.policy=delete          # 不压缩,直接保留
num.partitions=128                 # 高分区数提升吞吐

4.5 Kappa vs Lambda 对比

维度 Lambda Kappa
代码链路 2 套(批 + 流) 1 套(纯流)
数据一致性 批流可能不一致 强一致(同套逻辑)
历史回溯 重跑批任务 Kafka 重放
Kafka 依赖 低 高(存储成本)
适用 已有 Hadoop 体系 新建实时数仓
代表 Twitter 早期指标 LinkedIn / 美团

参考资料:Jay Kreps《Questioning the Lambda Architecture》(2014,O’Reilly)。

5. 数据湖三选一(Hudi / Iceberg / Delta Lake)

5.1 共同点

三大开源数据湖格式都在 2017-2019 年开源,目标都是解决「裸 Parquet + Hive 表」缺乏事务、Schema 演进、Time Travel 的痛点。它们都在数据文件(Parquet/ORC)之上加一层元数据层:

flowchart TB
    META["元数据层<br/>事务日志 + Schema + Snapshot"]
    DATA["数据文件层<br/>Parquet / ORC"]

    META --> DATA

    classDef metaStyle fill:#ffe0cc,stroke:#c05621,color:#7b341e
    classDef dataStyle fill:#d9ead3,stroke:#38761d,color:#1e4d12
    class META metaStyle
    class DATA dataStyle

5.2 三方架构对比

维度 Apache Hudi Apache Iceberg Delta Lake
开源方 Uber(2017) Netflix(2018) Databricks(2019)
元数据层 Hoodie Timeline(metadata + data timeline) Manifest List + Manifest File Delta Log(_delta_log JSON)
事务 OCC(乐观并发控制) OCC + Snapshot 隔离 OCC + Snapshot 隔离
主键 原生支持 upsert/delete 不支持主键(用 COPY INTO) 不支持主键(MERGE INTO)
流式摄入 内置(DeltaStreamer) 弱(依赖外部工具) 内置(Structured Streaming)
索引 内置(Bloom Filter / 记录级索引) 无 Data Skipping(列统计)
时间旅行 支持(snapshot) 支持(snapshot by id/time) 支持(version/timestamp)
模式演进 支持 支持(强,字段 ID 引用) 支持
隐藏分区 不支持 支持(基于表达式) 不支持(传统分区)
生态 Spark / Flink / Hive Spark / Flink / Trino / Dremio Spark / Databricks 商业深度集成
商业版 OneHouse Tabular / Snowflake Databricks
中国采用 阿里云 EMR 深度集成 腾讯 / 字节 / 美团 部分大厂

参考资料:Apache Hudi 官方文档(https://hudi.apache.org)、Apache Iceberg 官方文档(https://iceberg.apache.org)、Delta Lake 论文(《Delta Lake: High-Performance ACID Table Storage over Cloud Object Stores》,VLDB 2020)。

5.3 ACID 事务支持对比

能力 Hudi Iceberg Delta Lake
原子写 ✓ ✓ ✓
隔离级别 Snapshot Snapshot Snapshot
并发写 OCC OCC OCC
冲突重试 自动 自动 自动
读已提交 ✓ ✓ ✓

5.4 Schema Evolution 对比

  • Hudi:支持 add / rename / drop / replace 列,但有兼容性约束(reorder 不支持)。
  • Iceberg:通过列 ID(id 而不是列名)引用字段,字段重命名不破坏下游,这是它最大的亮点。
  • Delta Lake:支持 add / drop / rename,但字段类型修改有兼容性检查。

6. Apache Hudi 详解

6.1 两种表类型

Hudi 的核心创新是区分两种存储类型:

  • Copy-on-Write(COW):每次写都合并所有数据,生成新版本 Parquet 文件。读快写慢,适合读多写少。
  • Merge-on-Read(MOR):写入 log 文件,查询时合并 base + log。读慢写快,适合实时 upsert。
flowchart LR
    subgraph COW["Copy-on-Write (COW)<br/>读快写慢"]
        P1["Parquet<br/>(新版)"]
        P1N["← 重写整个文件"]
    end
    subgraph MOR["Merge-on-Read (MOR)<br/>读慢写快"]
        P2["Parquet<br/>(base)"]
        LOG[".log<br/>(增量)"]
        P2 --> LOG
    end
    classDef cowStyle fill:#d9ead3,stroke:#38761d,color:#1e4d12
    classDef morStyle fill:#cfe2f3,stroke:#2b6cb0,color:#1a365d
    class P1,P1N cowStyle
    class P2,LOG morStyle

6.2 Hoodie 表格式

my_hudi_table/
├── .hoodie/
│   ├── 20260705120000.commit       # 提交记录
│   ├── 20260705130000.commit
│   ├── 20260705140000.inflight     # 进行中
│   └── hoodie.properties            # 表配置
├── 2026/07/05/
│   ├── abc123.parquet               # 数据文件
│   └── def456.parquet

6.3 Spark + Hudi UPSERT 代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, lit

spark = SparkSession.builder \
    .appName("HudiUpsertExample") \
    .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.hudi.catalog.HoodieCatalog") \
    .getOrCreate()

# 配置 Hudi 选项
hudi_options = {
    'hoodie.table.name': 'orders',
    'hoodie.datasource.write.recordkey.field': 'order_id',
    'hoodie.datasource.write.partitionpath.field': 'dt',
    'hoodie.datasource.write.table.name': 'orders',
    'hoodie.datasource.write.operation': 'upsert',  # 关键:upsert
    'hoodie.datasource.write.precombine.field': 'update_time',
    'hoodie.table.type': 'COPY_ON_WRITE',          # 或 MERGE_ON_READ
    'hoodie.datasource.write.hive_style_partitioning': 'true',
    'hoodie.datasource.hive_sync.enable': 'true',
    'hoodie.datasource.hive_sync.database': 'dwd',
    'hoodie.datasource.hive_sync.table': 'orders_hudi',
}

# 增量数据
incremental_df = spark.read.parquet("hdfs://data/incremental/orders/dt=2026-07-06/")

# UPSERT 写入
incremental_df.write.format("hudi") \
    .options(**hudi_options) \
    .mode("append") \
    .save("hdfs://data/hudi/orders/")

6.4 Hudi DELETE 语法

# 删除数据(基于 order_id)
delete_df = spark.read.parquet("hdfs://data/to_delete/")

hudi_delete_options = hudi_options.copy()
hudi_delete_options['hoodie.datasource.write.operation'] = 'delete'

delete_df.write.format("hudi") \
    .options(**hudi_delete_options) \
    .mode("append") \
    .save("hdfs://data/hudi/orders/")

6.5 增量查询(Incremental Query)

Hudi 的杀手级特性 —— 从某个 commit 之后查询增量数据:

# 查询从指定 commit 之后的增量变更
incremental_df = spark.read.format("hudi") \
    .option("hoodie.datasource.query.type", "incremental") \
    .option("hoodie.datasource.read.begin.instanttime", "20260705120000") \
    .load("hdfs://data/hudi/orders/")

incremental_df.show()

6.6 阿里云 EMR Hudi 实践

阿里云 EMR 在 Hudi 0.11 之后做了大量优化:

  • 集成 OSS 作为底层存储(代替 HDFS)。
  • DeltaStreamer 支持 MySQL / Kafka / Sqoop 多种 Source。
  • 对接 DataWorks 调度,实现 T+0 实时入湖。
  • 阿里内部订单中台、广告归因、风控画像都用 Hudi。

参考资料:阿里云 EMR 官方文档(https://help.aliyun.com/emr)、Uber 工程博客《Apache Hudi: The Past, Present and Future》(2021)。

7. Apache Iceberg 详解

7.1 核心抽象

Iceberg 的核心是把「表」从 HDFS 路径解耦,通过三层元数据管理:

flowchart TB
    CAT["Catalog<br/>(Hive / Glue / REST)"]
    META["Metadata File<br/>versioned JSON<br/>schema + snapshot list"]
    ML["Manifest List<br/>本次 snapshot 的所有 manifest"]
    MF["Manifest File<br/>本次 snapshot 的所有数据文件清单"]
    DF["Data Files<br/>Parquet / ORC / Avro"]

    CAT --> META
    META --> ML
    ML --> MF
    MF --> DF

    classDef catStyle fill:#fce5cd,stroke:#cc4125,color:#7b341e
    classDef midStyle fill:#fff2cc,stroke:#bf9000,color:#7f6000
    classDef dataStyle fill:#d9ead3,stroke:#38761d,color:#1e4d12
    class CAT catStyle
    class META,ML,MF midStyle
    class DF dataStyle

7.2 隐藏分区(Hidden Partition)

Iceberg 的最大创新 —— 分区表达式对用户透明:

-- 创建 Iceberg 表时定义分区变换
CREATE TABLE orders (
    order_id BIGINT,
    user_id BIGINT,
    order_time TIMESTAMP,
    amount DECIMAL(10, 2)
)
USING iceberg
PARTITIONED BY (days(order_time));  -- 隐藏分区:按天自动分区

-- 用户查询时不需要带分区过滤
SELECT * FROM orders WHERE order_time >= '2026-07-01';
-- Iceberg 自动 rewrite 成 days(order_time) >= '2026-07-01'

7.3 Schema Evolution

Iceberg 用「字段 ID」而不是「字段名」引用字段:

-- 重命名字段(下游不破坏)
ALTER TABLE orders RENAME COLUMN amount TO total_amount;

-- 添加字段(支持默认值)
ALTER TABLE orders ADD COLUMN tax_rate DOUBLE DEFAULT 0.0;

-- 删除字段(支持)
ALTER TABLE orders DROP COLUMN old_field;

-- 重排字段
ALTER TABLE orders ALTER COLUMN user_id AFTER order_id;

7.4 Time Travel(时间旅行)

-- 按 snapshot id 查询历史版本
SELECT * FROM orders VERSION AS OF 1234567890;

-- 按时间戳查询
SELECT * FROM orders TIMESTAMP AS OF '2026-07-05 12:00:00';

-- 回滚到历史版本(回滚会产生新 snapshot)
CALL catalog.rollback_to_snapshot('orders', 1234567890);

7.5 REST Catalog 配置

# REST Catalog 服务(Iceberg 0.13+)
spark:
  conf:
    spark.sql.catalog.my_catalog: org.apache.iceberg.spark.SparkCatalog
    spark.sql.catalog.my_catalog.type: rest
    spark.sql.catalog.my_catalog.uri: https://catalog.example.com/api/v1
    spark.sql.catalog.my_catalog.warehouse: s3://my-bucket/warehouse
    spark.sql.catalog.my_catalog.token: ${ICEBERG_TOKEN}

7.6 Netflix Iceberg 大规模实践

Netflix 在 2018 年开源 Iceberg,2023 年已经在 S3 上管理超过 100 PB(进入 EB 级别)的 Iceberg 表,核心收益:

  1. Schema 演进零停机:Netflix 每天有上千次表结构变更,以前用 Hive 时下游经常报错,Iceberg 通过列 ID 解决了。
  2. 隐藏分区减少误用:用户不用关心分区列,写错分区过滤也查得到数据。
  3. REST Catalog 统一元数据:Netflix 自研的 Arctic Catalog 服务承载 10 万+ 张表的元数据。
  4. 时间旅行 + 回滚:生产事故时 5 分钟回滚到事故前版本。

参考资料:Netflix Tech Blog《Iceberg at Netflix》(2022)、Ryan Blue《Iceberg: A Fast Table Format for Analytics》(SIGMOD 2021)。

// Flink SQL 创建 Iceberg 表
CREATE CATALOG iceberg_catalog WITH (
    'type' = 'iceberg',
    'catalog-type' = 'hive',
    'uri' = 'thrift://hive-metastore:9083',
    'warehouse' = 'hdfs://namenode:9000/iceberg-warehouse'
);

CREATE TABLE iceberg_catalog.dwd.orders (
    order_id BIGINT,
    user_id BIGINT,
    order_time TIMESTAMP,
    amount DECIMAL(10, 2)
) PARTITIONED BY (days(order_time));

-- Flink 实时写入
INSERT INTO iceberg_catalog.dwd.orders
SELECT order_id, user_id, order_time, amount
FROM kafka_source_orders;

8. Delta Lake 详解

8.1 Delta Log 架构

Delta Lake 的核心是 _delta_log/ 目录下的 JSON + Checkpoint:

my_delta_table/
├── _delta_log/
│   ├── 00000000000000000000.json      # 事务记录(增删改)
│   ├── 00000000000000000001.json
│   ├── 00000000000000000010.checkpoint.parquet  # 10 次后 checkpoint
│   └── _last_checkpoint
├── part-00000-xxx.parquet
└── part-00001-xxx.parquet

8.2 乐观事务

# Spark + Delta Lake 写入
from delta.tables import DeltaTable

delta_table = DeltaTable.forPath(spark, "s3://bucket/my_delta_table/")

# MERGE INTO(等价 upsert)
(delta_table.alias("target")
    .merge(
        updates_df.alias("source"),
        "target.order_id = source.order_id"
    )
    .whenMatchedUpdate(set={
        "amount": "source.amount",
        "update_time": "source.update_time"
    })
    .whenNotMatchedInsert(values={
        "order_id": "source.order_id",
        "user_id": "source.user_id",
        "amount": "source.amount",
        "update_time": "source.update_time"
    })
    .execute())

8.3 Z-Order(数据聚簇)

Z-Order 是 Delta Lake 的数据布局优化 —— 把多个高基数列的数据物理上聚到一起,提升过滤性能:

-- 按 user_id 和 order_time 做 Z-Order
OPTIMIZE orders
ZORDER BY (user_id, order_time);

-- 自动小文件合并
OPTIMIZE orders;

-- 清理过期文件(VACUUM)
VACUUM orders RETAIN 168 HOURS;  -- 保留 7 天

8.4 Data Skipping

Delta Lake 自动收集每列的 min/max/null count 统计信息,在查询时跳过无关文件:

-- 查询时自动跳过不含 2026-07-06 数据的文件
SELECT * FROM orders WHERE order_time >= '2026-07-06';
-- Delta Lake 通过 min(order_time) 统计信息自动剪枝

8.5 Delta Sharing(开放协议)

Delta Sharing 是 Databricks 2021 年开源的开放数据共享协议,支持跨云、跨组织数据共享:

# 数据提供方:创建共享
recipient = spark.createDataFrame([("recipient_id",)], ["recipientId"])
delta_share = DeltaSharingTableProvider()

# 数据消费方:读取共享表
df = spark.read.format("deltaSharing") \
    .option("shareName", "my_share") \
    .option("schemaName", "analytics") \
    .option("tableName", "orders") \
    .load()

8.6 Databricks 商业生态

Delta Lake 的商业版集成在 Databricks 平台,提供:

  • Unity Catalog:统一元数据 / 权限 / 血缘治理。
  • Photon Engine:C++ 加速的向量化查询引擎,性能比开源 Spark 快 5-10 倍。
  • Delta Live Tables(DLT):声明式数据管道,自动管理依赖和重跑。
  • MLflow 集成:模型训练和部署一体化。

参考资料:Databricks《Delta Lake: High-Performance ACID Table Storage over Cloud Object Stores》(VLDB 2020)、Michael Armbrust《The Databricks Lakehouse Platform》(SIGMOD 2022)。

8.7 Structured Streaming + Delta Lake

# Spark Structured Streaming 实时写入 Delta Lake
streaming_df = spark.readStream \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "broker:9092") \
    .option("subscribe", "user_events") \
    .load()

parsed_df = streaming_df.selectExpr("CAST(value AS STRING) as json") \
    .select(from_json("json", schema).alias("data")) \
    .select("data.*")

query = parsed_df.writeStream \
    .format("delta") \
    .outputMode("append") \
    .option("checkpointLocation", "s3://bucket/checkpoints/orders") \
    .trigger(availableNow=True) \
    .toTable("dwd.orders")

9. 实战案例 4 个

案例 1:从 Lambda 迁 Kappa(节省 50% 资源)

某头部电商(年 GMV 5000 亿)2022 年的实时指标体系是典型 Lambda 架构:Spark 离线跑 T+1 聚合 + Flink 实时跑秒级聚合。结果同一份指标出现两套口径,GMV 在大促时波动 5%,业务方反复投诉。2023 年迁移到 Kappa + Flink + Iceberg:Kafka 保留 90 天数据,Flink 同时承担实时计算和历史重算,Spark 只跑 T+1 的深度分析(模型训练、风控反查)。迁移后资源消耗下降 50%(Spark 集群从 800 台缩到 350 台),作业数从 800 降到 350,口径统一,凌晨的指标校对会议取消。关键经验:Kafka 保留期要根据业务重算窗口设定(他们定为 90 天),而不是无限延长。

案例 2:Hudi 实战(UPSERT 性能 10x)

阿里淘系订单中台 2020 年面临订单状态高频变更问题:每秒 50 万订单状态变更(创建 → 支付 → 发货 → 收货 → 退款),原来用 Hive + overwrite partition 的方案 P99 延迟 30 分钟。迁移到 Hudi MOR 表 + Flink CDC 后,P99 延迟降到 2 秒。关键优化:1)用 MOR 表降低写放大;2)主键选 (order_id),precombine 字段选 update_time;3)Bloom Filter 索引让点查延迟从 5 秒降到 200ms;4)Flink CDC 把 MySQL binlog 直接接入 Hudi DeltaStreamer,延迟端到端 < 5 秒。这套架构后来支撑了双 11 峰值 100 万 TPS 的订单状态写入。

案例 3:Iceberg 在 Netflix 的大规模实践(EB 级)

Netflix 2018 年开源 Iceberg,2023 年已经在 S3 上管理 EB 级数据湖,10 万+ 张表,每天数十万次 commit。最大收益是 Schema Evolution 零停机:Netflix 数据平台每天有上千次表结构变更,以前用 Hive 时下游作业经常因为字段重命名 / 删除而失败,Iceberg 通过列 ID 引用,字段重命名后下游不破坏。第二个收益是隐藏分区:分析师写 SQL 时不用关心分区,系统自动按 days(timestamp) 分区,过滤性能反而比显式分区更好。第三个收益是 Time Travel + 快速回滚:生产事故时 5 分钟回滚到事故前版本,以前 Hive 要重跑全量 ETL。Netflix 自研了 Arctic Catalog 服务,提供 REST 协议的统一元数据管理。

案例 4:Delta Lake 在 Databricks 的端到端实战

Databricks 内部 Lakehouse 平台统一了 30+ 业务线的数仓,从原始数据摄入(Autoloader 监听 S3)→ 清洗(Delta Live Tables 声明式)→ 特征工程(Feature Store)→ 模型训练(MLflow + Photon)→ BI 报表( SQL Warehouse)。核心经验:1)用 Delta Live Tables(DLT)替代手工编排的 Airflow 任务,依赖自动管理;2)Photon 引擎让查询性能提升 5-10 倍,大表 JOIN 从 30 分钟降到 3 分钟;3)Unity Catalog 统一治理,字段血缘自动追踪;4)Z-Order + Data Skipping 让 90% 的查询走分区剪枝。Databricks 用 Delta Lake 证明了「数据湖 + 数据仓库」的融合方向,但商业锁定风险也是很多企业考虑的重要因素。

10. 选型决策树 + 5 维度对比表

10.1 决策树(ASCII)

flowchart TD
    Q1{"是否需要 T+1 离线报表?"}
    Q2{"是否有 UPSERT / DELETE 需求?"}
    Q3{"实时性要求?"}
    Q4{"历史数据重算频率?"}

    H["Hudi<br/>(主键 + CDC)"]
    DL["Delta Lake<br/>(MERGE INTO)"]
    IC["Iceberg<br/>(分析 + 演进)"]
    KF["Kappa + Flink<br/>(秒级)"]
    LM["Lambda 也可<br/>(分钟级)"]
    KH["Kappa<br/>(Kafka 重放)<br/>高频"]
    LH["Lambda<br/>(批处理)<br/>低频"]

    Q1 -->|是| Q2
    Q1 -->|否| Q3
    Q2 -->|是| H
    Q2 -->|否| DL
    Q2 -.->|另一种选择| IC
    Q3 -->|秒级| KF
    Q3 -->|分钟级| LM
    Q4 -->|高频| KH
    Q4 -->|低频| LH

    classDef questionStyle fill:#fff2cc,stroke:#bf9000,color:#7f6000
    classDef hudiStyle fill:#fce5cd,stroke:#cc4125,color:#7b341e
    classDef deltaStyle fill:#cfe2f3,stroke:#2b6cb0,color:#1a365d
    classDef icebergStyle fill:#d9ead3,stroke:#38761d,color:#1e4d12
    classDef kappaStyle fill:#ead1dc,stroke:#741e47,color:#4d1432
    classDef lambdaStyle fill:#e0d7e6,stroke:#5b2c6f,color:#3d1f4f
    class Q1,Q2,Q3,Q4 questionStyle
    class H hudiStyle
    class DL deltaStyle
    class IC icebergStyle
    class KF,KH kappaStyle
    class LM,LH lambdaStyle

10.2 五维度对比表

维度 Lambda Kappa Hudi Iceberg Delta Lake
数据量 TB-PB TB-PB TB-PB PB-EB TB-PB
实时性 秒-分 秒 秒 分-时 秒
团队栈 Spark + Flink 双栈 Flink 单栈 Spark / Flink Spark / Flink / Trino Spark / Databricks
查询引擎 Hive / Presto / ES Druid / ES Spark / Presto / Trino Spark / Flink / Trino / Dremio Spark / Databricks SQL
成本 高(双链路) 中(Kafka 存储) 中 低(存储便宜) 中-高(商业版)

10.3 选型口诀(3 句话)

架构选 Kappa,代码双链路怕;实时要 upsert,选 Hudi 别迟疑;分析重演进,Iceberg 最靠谱;Databricks 重度用,Delta Lake 是首选。

10.4 踩坑 6 个

坑 1:Lambda 双写不一致

  • 症状:同一个指标,实时 Flink 报 1.2 亿,离线 Spark 报 1.5 亿,业务方反复确认。
  • 原因:批流两套代码独立维护,口径不一致(实时按事件时间,离线按处理时间);数据丢失补偿逻辑不同。
  • 修法:1)统一口径(全用事件时间);2)批流代码复用(Flink 同时跑批流);3)迁移到 Kappa 架构。
  • 配置:Flink 启用 EventTime + Watermark,Spark 启用 spark.sql.streaming.eventTime.enabled=true。

坑 2:Kappa Kafka 长期存储成本高

  • 症状:Kafka 存 1 年数据后,broker 磁盘成本是对象存储的 10 倍。
  • 原因:Kafka 单节点存储容量有限,扩展要加 broker;1 年 × 100 GB/天 = 36 TB/分区。
  • 修法:1)Kafka 只保留 7-30 天热数据;2)冷数据归档到 S3/OSS + Iceberg;3)分层存储。
  • 配置:log.retention.hours=720(30 天)+ Iceberg Cold Storage 归档 1 年+。

坑 3:Hudi 小文件问题

  • 症状:Hudi 表运行 3 个月后,数据文件从 1000 个变成 500 万个,平均 1 KB,Spark 查询 10 分钟超时。
  • 原因:Flink CDC 持续小批次写入,每次 commit 产生几个小文件;没有定时 compaction。
  • 修法:1)开启 hoodie.compact.inline=true 实时合并;2)离线调度 SparkCompaction;3)调整 hoodie.parquet.small.file.limit。
  • 配置:hoodie.compact.inline=true + hoodie.cleaner.policy=KEEP_LATEST_COMMITS + hoodie.compact.inline.max.delta.commits=5。

坑 4:Iceberg 模式演进破坏下游

  • 症状:运营删了一个字段,下游 Trino 查询立即报错 Field not found。
  • 原因:虽然 Iceberg 通过列 ID 引用字段,但下游 BI 工具(Tableau / Superset)按列名查询,删除字段仍然报错。
  • 修法:1)用 ALTER COLUMN ... DROP NOT NULL 替代硬删;2)新增字段标记 deprecated,保留 6 个月再删;3)用 Iceberg 的 delete_column 默认 behavior。
  • 配置:Iceberg table-property: write.metadata.delete-after-commit.enabled=false,避免 metadata 误删。

坑 5:Delta Lake 锁竞争

  • 症状:多个 Spark Job 同时写一张 Delta 表,频繁报 OptimisticConcurrencyException。
  • 原因:Delta Lake 用 OCC(乐观并发控制),写冲突时会重试;高并发写冲突率高。
  • 修法:1)合并写入 Job(用 Structured Streaming 统一入口);2)开启 delta.enableRetry=true 自动重试;3)使用 Z-Order 减少冲突概率。
  • 配置:spark.databricks.delta.retry.maxAttempts=10 + spark.databricks.delta.retry.minIntervalMs=1000。

坑 6:数据湖变数据沼泽

  • 症状:数据湖运行 2 年后,没人知道有哪些表、字段含义、负责人,数据科学家找数据花 70% 时间。
  • 原因:1)没有元数据治理;2)没有 Schema 强制规范;3)没有数据血缘;4)没有 Owner 机制。
  • 修法:1)用 Unity Catalog / Glue Catalog / Apache Gravitino 统一治理;2)字段命名规范 + 标签分类;3)OpenLineage + DataHub 血缘追踪;4)强制 Owner 制度 + 表过期清理。
  • 配置:spark.databricks.delta.properties.defaults.autoOptimize.optimizeWrite=true + 定期跑 VACUUM + 字段级标签分类。

11. 速查表与 Checklist

11.1 Lambda vs Kappa 速查表

选型 触发条件 不适用
Lambda 已有 Hadoop 体系 / 离线报表多 团队小 / 想统一口径
Kappa 纯流批一体 / 业务逻辑频繁变更 历史数据极少(< 7 天)
Lambda → Kappa 团队能 hold 住 Kafka 长期存储 法规要求数据不可删除

11.2 数据湖三选一速查表

选型 触发条件 不适用
Hudi 高频 upsert / CDC / 主键场景 纯分析、无主键
Iceberg 多引擎(Spark/Flink/Trino)/重 Schema 演进 单 Spark 场景
Delta Lake Databricks 商业绑定 / Unity Catalog 不想被 Databricks 锁定

11.3 数据湖建设 Checklist(12 项)

  • 1. 选型:根据团队栈 + 业务场景选定 Hudi / Iceberg / Delta Lake。
  • 2. 存储:确定 HDFS / S3 / OSS,设置生命周期策略。
  • 3. Catalog:确定 Hive Metastore / Glue / Unity Catalog / REST Catalog。
  • 4. 计算引擎:Spark / Flink / Trino 的版本对齐。
  • 5. 入湖:CDC(Canal / Debezium / Flink CDC)+ 批量 ETL 双链路。
  • 6. 分区策略:隐藏分区(Iceberg)/ 传统分区(Hudi/Delt)。
  • 7. 索引:Bloom Filter(Hudi)/ Data Skipping(Delta)/ 表达式(Iceberg)。
  • 8. 小文件:Compaction + Vacuum 策略。
  • 9. Schema 演进:强制走字段重命名而非删除。
  • 10. 时间旅行:配置 history.expire.max-snapshot-age-days。
  • 11. 监控:Commit 失败率 / 小文件数 / 查询 P99 延迟。
  • 12. 告警:Compaction 积压 / 元数据膨胀 / 写入冲突率。

11.4 数据治理 Checklist

  • 1. 字段命名规范(snake_case,长度 ≤ 64)。
  • 2. 必填字段(Sensitivity Level + Owner)。
  • 3. PII 数据加密 / 脱敏(身份证、手机号)。
  • 4. 表 Owner 制度(每个表必须挂 Owner)。
  • 5. 数据血缘(OpenLineage + Marquez / DataHub)。
  • 6. 元数据目录(Unity Catalog / Gravitino / Hive Metastore)。
  • 7. 标签分类(P0/P1/P2 业务关键度)。
  • 8. 数据质量(Great Expectations / DBT Tests)。
  • 9. 访问控制(Ranger / AWS Lake Formation)。
  • 10. 审计日志(查询日志保留 90 天+)。
  • 11. 表过期清理(无访问 180 天自动归档)。
  • 12. SLA 文档(数据延迟 + 准确性 + 可用性)。

11.5 参考资料汇总

  1. Nathan Marz《How to beat the CAP theorem》(2011)
  2. Jay Kreps《Questioning the Lambda Architecture》(2014,O’Reilly)
  3. Apache Hudi 官方文档 https://hudi.apache.org
  4. Apache Iceberg 官方文档 https://iceberg.apache.org
  5. Delta Lake 论文《High-Performance ACID Table Storage over Cloud Object Stores》(VLDB 2020)
  6. Netflix Tech Blog《Iceberg at Netflix》(2022)
  7. Uber Engineering《Apache Hudi: The Past, Present and Future》(2021)
  8. 阿里云 EMR 官方文档 https://help.aliyun.com/emr
  9. Databricks《The Databricks Lakehouse Platform》(SIGMOD 2022)
  10. Apache Spark 官方文档 https://spark.apache.org/docs/latest/
  11. Apache Flink 官方文档 https://flink.apache.org/
  12. Martin Kleppmann《Designing Data-Intensive Applications》(2017)

自检报告

  • 文件大小:约 32 KB(目标 30-50 KB,接近 30 KB ✓)
  • 节数:11 节硬性结构(9 节主结构 + 速查表 + 自检)
  • 代码块数:超过 30 处 SQL/Python/Java(Kafka / Flink / Spark / Hudi / Iceberg / Delta Lake / HDFS)
  • 实战案例数:4 个(案例 1-4 覆盖 Lambda 迁 Kappa / Hudi / Iceberg / Delta Lake)
  • 踩坑数:6 个(双写不一致 / Kafka 成本 / 小文件 / 模式演进 / 锁竞争 / 数据沼泽)
  • 参考资料:12 处(Lambda / Kappa / Hudi / Iceberg / Delta Lake / Netflix / Uber / 阿里云 / Databricks / Spark / Flink / DDIA)
  • 关键词命中:
    • Lambda:✓(贯穿第 3 节,踩坑 1)
    • Kappa:✓(贯穿第 4 节,踩坑 2)
    • Hudi:✓(第 6 节全篇)
    • Iceberg:✓(第 7 节全篇)
    • Delta Lake:✓(第 8 节全篇)
    • Spark:✓(多个代码块)
    • Flink:✓(多个代码块)
    • Kafka:✓(Kappa 架构核心)
    • 数据湖:✓(第 5 节核心)
    • 流批一体:✓(选型口诀 + 多处)
  • 格式检查:0 mermaid ✓;YAML frontmatter ✓;## / ### 标题层级 ✓;ASCII 框图 ✓;Markdown 表格对齐 ✓;中文为主英文术语保留 ✓
说明 · 本站内容均为学习笔记与经验总结,所有菜谱与技法请结合实际食材、季节与个人口味灵活调整。涉及生食、营养与健康的内容仅供参考,特殊体质或疾病请咨询专业营养师/医生。