湖仓一体架构:Iceberg、Delta Lake 与 Hudi 的统一数据底座

深入解析湖仓一体(Lakehouse)架构:数据湖与数仓的演进关系、Apache Iceberg 开放表格式(ACID/时间旅行/增量读取)、Delta Lake 与 Hudi 三方对比、Bronze/Silver/Gold 分层设计、统一存储与多引擎协作、文件布局优化与批流一体写入,附 Spark/Flink/Trino 生产配置与落地案例。

引言

过去十年,“数据湖"与"数据仓库"被当作两种对立的哲学:湖强调低成本、Schema-On-Read 与弹性存储,仓强调可靠性、ACID 与治理。湖仓一体(Lakehouse)则试图兼得两者——在低成本对象存储之上,用一层开放表格式(Open Table Format)带来事务、Schema 演化、时间旅行与增量读取能力。以 Apache Iceberg、Delta Lake、Apache Hudi 为代表的表格式,已经成为现代数据平台的标准底座。本文将系统拆解湖仓架构的选型逻辑、分层设计、多引擎协作与生产调优。

湖仓一体的本质不是"湖"也不是"仓”,而是用开放表格式在对象存储上重建了数据仓库的可靠性契约。


一、数据湖 vs 数仓 vs 湖仓

1.1 三类架构的本质差异

维度数据湖数据仓库湖仓一体
存储介质廉价对象存储专用集群存储对象存储
数据模型Schema-On-ReadSchema-On-WriteSchema-On-Write(表格式)
ACID 事务不支持完整支持表格式支持
成本低高(存储+算力绑定)低(存算分离)
治理弱强强(元数据层)
典型引擎Spark / Presto厂商专用 SQLSpark / Flink / Trino

1.2 演进路线

第一代 HDFS + Hive
   │  (缺 ACID、缺优化器、Hive 性能差)
   ▼
第二代 云数据仓库(Redshift / BigQuery / Snowflake)
   │  (性能强,但锁定厂商、成本高)
   ▼
第三代 湖仓一体(Iceberg / Delta / Hudi)
   │  (开放、低成本、可靠性三合一)
   ▼
   统一数据底座 + 存算分离 + 多引擎共享

1.3 选型决策框架

选择表格式前先回答三个问题:是否强依赖某厂商生态(Databricks 选 Delta)、是否要跨引擎自由读写(Iceberg 最开放)、是否以 CDC 流式入湖为主(Hudi 的 MOR 有优势)。


二、Apache Iceberg 表格式

2.1 Iceberg 核心特性

Iceberg 通过三层元数据(Catalog → Metadata File → Manifest List → Manifest File)管理数据文件,因此具备传统数据湖不具备的可靠语义。

特性说明生产意义
ACID 事务快照提交原子性并发写不互相污染
时间旅行读取任意历史快照审计、回溯、口径重算
增量读取基于快照差异的增量计划增量 ETL、CDC
Schema 演化增删改列带版本演进兼容下游消费
隐藏分区分区元数据化,无需维护分区列优化查询剪枝
开放中立多引擎一致读写摆脱引擎锁定

2.2 初始化 Iceberg Catalog

用 pyiceberg 通过 REST Catalog 初始化湖仓元数据。

# iceberg_catalog.py
from pyiceberg.catalog import load_catalog

catalog = load_catalog(
    "rest",
    **{
        "uri": "http://iceberg-rest:8181",
        "warehouse": "s3://data-lake/warehouse",
        "s3.endpoint": "http://minio:9000",
    },
)

# 列出命名空间并创建订单域命名空间
print(catalog.list_namespaces())
catalog.create_namespace("ods")

table = catalog.create_table(
    "ods.orders",
    schema={
        "order_id": "string",
        "user_id": "long",
        "amount": "double",
        "status": "string",
        "event_time": "timestamp",
    },
    partition_spec="day(event_time)",
    properties={"write.format.default": "parquet"},
)

2.3 时间旅行查询

时间旅行是 Iceberg 对审计与回溯最直接的价值:无需恢复备份,直接读历史快照。

-- 基于快照 ID 读取历史版本
SELECT * FROM ods.orders
  VERSION AS OF 8201468987888888888
WHERE event_time >= '2026-09-01';

-- 基于时间戳读取
SELECT * FROM ods.orders
  TIMESTAMP AS OF '2026-09-25 08:00:00'
LIMIT 100;

三、Delta Lake 与 Hudi 对比

3.1 三方特性对照

特性IcebergDelta LakeHudi
存储格式Parquet/ORC/AvroParquet(Log 记录变更)Parquet/Avro
时间旅行✅✅部分(基于 commit)
增量读取✅ 快照 diff✅ Change Data Feed✅ Incremental View
文件布局优化rewrite_data_filesOPTIMIZEclustering
批/流 UpsertMERGE INTOMERGECOW / MOR
厂商绑定中立Databricks 深度中立 + Hive 生态

3.2 Delta Lake 的 OPTIMIZE 与 ZORDER

Delta 把文件布局优化内置进 SQL,OPTIMIZE 合并小文件,ZORDER BY 建立多维本地性。

-- 合并小文件并按 dt, user_id 聚类
OPTIMIZE dws.daily_order_revenue
ZORDER BY (dt, user_id);

-- 清理超过保留期的历史版本(默认 7 天)
VACUUM dws.daily_order_revenue RETAIN 168 HOURS;

3.3 Hudi 的 COW 与 MOR

Hudi 的核心是 Copy-On-Write(写时复制)与 Merge-On-Read(读时合并)两种表类型:COW 读快写慢,MOR 写快读慢,适合 CDC 流式入湖。

-- 创建 MOR 表并做流式 Upsert
CREATE TABLE ods.orders_hudi (
  order_id STRING,
  status STRING,
  amount DOUBLE,
  ts TIMESTAMP
) USING hudi
TBLPROPERTIES (
  'hoodie.table.type' = 'MERGE_ON_READ',
  'hoodie.datasource.write.recordkey.field' = 'order_id',
  'hoodie.datasource.write.precombine.field' = 'ts'
);

-- 增量视图读取最近一次提交的变更
SELECT * FROM hudi_incremental('ods.orders_hudi', 'beginTime', '20260925120000');

四、湖仓分层设计:Bronze / Silver / Gold

4.1 Medallion 分层

湖仓的经典分层是 Bronze(原始层)→ Silver(清洗层)→ Gold(聚合/服务层),每层对应不同的质量与消费语义。

层英文内容质量要求消费者
原始层Bronze源数据原样落湖保真、可回溯数据工程师
清洗层Silver去重、标准化、Schema 统一字段级可信分析师
聚合层Gold指标、宽表、特征业务口径确定报表/模型

4.2 分层写入的 Spark 实现

用 Spark Structured Streaming 将 Kafka 数据按三层级联写入 Iceberg。

# medallion_pipeline.py
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, to_timestamp

spark = (
    SparkSession.builder.appName("medallion_lakehouse")
    .config("spark.sql.catalog.lake", "org.apache.iceberg.spark.SparkCatalog")
    .config("spark.sql.catalog.lake.type", "rest")
    .config("spark.sql.catalog.lake.uri", "http://iceberg-rest:8181")
    .getOrCreate()
)

stream = (
    spark.readStream.format("kafka")
    .option("kafka.bootstrap.servers", "kafka:9092")
    .option("subscribe", "orders.raw")
    .load()
)

# Bronze:原样落湖
bronze = (
    stream.selectExpr("CAST(value AS STRING) AS raw", "to_timestamp(timestamp) AS event_time")
    .writeStream.format("iceberg")
    .option("path", "lake.ods.orders_bronze")
    .trigger(processingTime="60 seconds")
    .checkpointLocation("s3://data-lake/checkpoints/orders_bronze")
)

# Silver:清洗后写入(并行流)
silver = (
    stream.selectExpr("CAST(value AS STRING) AS raw")
    .select(parse_json(col("raw")))
    .filter(col("status").isNotNull())
    .writeStream.format("iceberg")
    .option("path", "lake.ods.orders_silver")
    .trigger(processingTime="60 seconds")
    .checkpointLocation("s3://data-lake/checkpoints/orders_silver")
)

五、统一存储与多引擎协作

5.1 多引擎职责分工

湖仓的价值在于"一份数据、多引擎消费",各引擎各司其职而非互相替代。

引擎场景读写 Iceberg备注
Spark批处理、大规模 ETL✅主力批引擎
Flink流式入湖、CDC✅流批一体
Trino/Presto联邦 SQL 查询✅交互式分析
ImpalaHive 生态兼容部分迁移成本高

5.2 Trino 查询 Iceberg

Trino 通过 Iceberg Connector 直接读对象存储,实现"查询与写引擎解耦"。

-- Trino: 读取 S3 上的 Iceberg 表
CALL iceberg.system.register_table(
    schema_name => 'ods',
    table_name => 'orders',
    table_location => 's3://data-lake/warehouse/ods/orders'
);

SELECT dt, count(*) AS orders
FROM "lake.ods".orders
WHERE dt >= '2026-09-20'
GROUP BY dt
ORDER BY dt;

Flink 的 Iceberg 连接器支持两阶段提交,保证流式写入的精确一次语义。

# flink_iceberg.yaml
job:
  name: orders_cdc_to_iceberg
  checkpoint:
    interval: 60s
    mode: exactly_once
  source:
    connector: debezium
    database: mysql
    tables: [shop.orders]
  sink:
    connector: iceberg
    catalog-name: lake
    table: lake.ods.orders
    write:
      format: parquet
      distribution-mode: hash
      upsert-enabled: true

六、文件布局与优化

6.1 小文件问题

流式写入与频繁 Upsert 必然产生海量小文件,拖垮查询与元数据层。湖仓的优化本质是"定期把碎片合并成合理的布局"。

优化动作Iceberg 对应频率
合并小文件rewrite_data_files小时级/日级
清理过期快照expire_snapshots日级
清理孤儿文件remove_orphan_files周级
数据布局rewrite_manifests / sort与分区频率一致

6.2 Spark 调用的优化存储过程

Iceberg 提供系统存储过程(System Procedure),在 Spark 中直接调用。

-- 合并 orders_silver 中小于 32MB 的数据文件
CALL lake.system.rewrite_data_files(
    table => 'ods.orders_silver',
    strategy => 'binpack',
    min_file_size_in_bytes => 33554432,
    max_file_size_in_bytes => 134217728
);

-- 保留 7 天内快照,清理更早历史
CALL lake.system.expire_snapshots(
    table => 'ods.orders_silver',
    older_than => TIMESTAMP '2026-09-18 00:00:00',
    retain_last => 10
);

-- 清理无引用的孤儿文件
CALL lake.system.remove_orphan_files(
    table => 'ods.orders_silver',
    older_than => TIMESTAMP '2026-09-24 00:00:00'
);

6.3 排序与 Z-Order 的取舍

数据布局排序能显著提升过滤类查询的剪枝效率,但会引入排序写放大。

方式优势代价适用
自然分区无额外代价桶内乱序一般查询
排序写单键强本地性写放大高频点查
Z-Order多列近似局部性写放大明显多条件过滤
# zorder_write.py
from pyspark.sql.functions import col

# 用 sortWithinPartitions 实现轻量 Z-Order 效果
(
    spark.table("lake.ods.orders_silver")
    .repartitionByRange(col("dt"))
    .sortWithinPartitions(col("user_id"), col("order_id"))
    .writeTo("lake.ods.orders_silver_z")
    .using("iceberg")
    .overwritePartitions()
)

七、批流一体写入

7.1 统一语义下的双模写入

批流一体的核心是"同一张表、同一套语义,批与流都能安全写入"。Iceberg 通过快照隔离让并发读写互不阻塞。

写入模式引擎一致性保证典型场景
批量覆盖Spark batch快照提交日批重算
流式追加Flink/Spark streaming两阶段提交实时入湖
Upsert/MergeSpark MERGE INTO行级语义CDC 更新
删除+重写存储过程后台异步清理/布局优化

7.2 MERGE INTO 实现 CDC Upsert

用 Iceberg 的 MERGE INTO 把 Kafka 中的 CDC 变更(含删除)应用到银层。

MERGE INTO lake.ods.orders_silver t
USING (
  SELECT order_id, user_id, amount, status, op, event_time
  FROM lake.ods.orders_cdc
) s
ON t.order_id = s.order_id
WHEN MATCHED AND s.op = 'D' THEN DELETE
WHEN MATCHED THEN UPDATE SET
  t.amount = s.amount, t.status = s.status, t.event_time = s.event_time
WHEN NOT MATCHED AND s.op = 'I' THEN
  INSERT (order_id, user_id, amount, status, event_time)
  VALUES (s.order_id, s.user_id, s.amount, s.status, s.event_time);

7.3 批流口径对账

批流一体最容易被忽略的是对账:同一指标流式结果与日批结果必须收敛。

#!/bin/bash
# reconcile.sh
# 对比流式 Silver 表与日批 Gold 表的订单数
DIFF=$(trino --execute \
  "SELECT
     (SELECT count(*) FROM \"lake.ods\".orders_silver WHERE dt = '2026-09-25') -
     (SELECT count(*) FROM \"lake.dws\".daily_order_revenue WHERE dt = '2026-09-25')")

if [ "$DIFF" != "0" ]; then
  echo "RECONCILE FAILED: diff=$DIFF"
  exit 1
fi

八、生产案例

8.1 案例:某零售企业 3 个月湖仓迁移

某零售企业从 Hive 数仓迁移到 Iceberg 湖仓,替换了 40+ 条核心 ETL。

阶段动作结果
试点选择订单域 3 张表验证 ACID 与时间旅行口径回溯效率提升 10 倍
迁移Hive 表 CONVERT TO ICEBERG 逐步切换无停机迁移
流化Debezium CDC → Kafka → Flink → Iceberg新鲜度从 T+1 到分钟级
优化定时 compaction + z-order + 快照清理查询 P95 下降 60%
治理Trino 统一查询,Spark 统一写引擎解耦,成本降 40%

8.2 迁移语法

将存量 Hive/Delta 表原地转换为 Iceberg,降低迁移风险。

-- Hive 表原地转 Iceberg
CALL lake.system.migrate('ods.orders_legacy');

-- Delta 表转 Iceberg(Spark 3.x)
ALTER TABLE lake.ods.orders
  SET TBLPROPERTIES ('engine' = 'iceberg');

九、常见问题与最佳实践

Q1: Iceberg、Delta、Hudi 最终选哪个?

结论取决于约束:中立开放 + 多引擎自由选 Iceberg;Databricks 深度生态 + 极简体验选 Delta;强 CDC 流式入湖 + Hive 生态兼容选 Hudi。多数云厂商托管湖仓(AWS Athena、阿里云 MaxCompute)已默认支持 Iceberg,长期看 Iceberg 的生态位最稳。

Q2: 湖仓还需要数仓吗?

需要。湖仓解决"存储与可靠读取",数仓(如 ClickHouse/Doris 或云数仓)解决"低延迟高并发分析"。成熟架构是湖仓为底座、OLAP 为加速层:数据在湖仓中保持单一事实源,物化到 OLAP 引擎满足交互式查询。

Q3: 时间旅行会无限占用存储吗?

会。快照保留期内的历史文件都会占用空间,所以必须配合 expire_snapshots 与 remove_orphan_files 的定期清理策略。通常保留 7-30 天即可满足审计与回溯需求。

Q4: 小文件问题真的需要天天处理吗?

取决于写入模式。流式写入(尤其是秒级触发)会产生大量小文件,建议按小时级 compaction;纯日批场景按天 compaction 即可。关键是让"合并速度"追上"产生速度",否则查询性能会随时间退化。


总结

决策点推荐理由
表格式Apache Iceberg开放中立、多引擎一致
分层Bronze/Silver/Gold语义清晰、治理可落地
写引擎Spark(批) + Flink(流)批流一体、两阶段提交
查引擎Trino/Presto联邦查询、存算解耦
优化compaction + z-order + 清理可持续的查询性能
加速层ClickHouse/Doris低延迟 OLAP 消费

湖仓一体不是技术的终点,而是数据平台走向"开放 + 可靠 + 经济"的分水岭。落地它的关键不在于选一个"最好的表格式",而在于建立一套分层清晰、写入统一、优化自动、对账常态的运营机制。先让一个域的几张核心表跑通全链路,再把范式复制到全公司,是成功率最高的路径。


参考与延伸阅读

  • Apache Iceberg 官方文档:Spec、Spark/Trino 集成与维护存储过程
  • Delta Lake 官方文档:OPTIMIZE、ZORDER 与 Change Data Feed
  • Apache Hudi 官方文档:COW/MOR 表类型与增量视图
  • The Data Lakehouse: The Key to the Future Data Platform(Big Data 界文章)

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「data-engineering」更多文章

  1. 特征存储(Feature Store)架构:从一致性到在线检索的完整实践
  2. 反向 ETL 与数据激活:让数据仓库的价值回到业务系统
  3. 数据可观测性:从管道监控到数据宕机的全方位保障