引言
过去十年,“数据湖"与"数据仓库"被当作两种对立的哲学:湖强调低成本、Schema-On-Read 与弹性存储,仓强调可靠性、ACID 与治理。湖仓一体(Lakehouse)则试图兼得两者——在低成本对象存储之上,用一层开放表格式(Open Table Format)带来事务、Schema 演化、时间旅行与增量读取能力。以 Apache Iceberg、Delta Lake、Apache Hudi 为代表的表格式,已经成为现代数据平台的标准底座。本文将系统拆解湖仓架构的选型逻辑、分层设计、多引擎协作与生产调优。
湖仓一体的本质不是"湖"也不是"仓”,而是用开放表格式在对象存储上重建了数据仓库的可靠性契约。
一、数据湖 vs 数仓 vs 湖仓
1.1 三类架构的本质差异
| 维度 | 数据湖 | 数据仓库 | 湖仓一体 |
|---|---|---|---|
| 存储介质 | 廉价对象存储 | 专用集群存储 | 对象存储 |
| 数据模型 | Schema-On-Read | Schema-On-Write | Schema-On-Write(表格式) |
| ACID 事务 | 不支持 | 完整支持 | 表格式支持 |
| 成本 | 低 | 高(存储+算力绑定) | 低(存算分离) |
| 治理 | 弱 | 强 | 强(元数据层) |
| 典型引擎 | Spark / Presto | 厂商专用 SQL | Spark / 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 三方特性对照
| 特性 | Iceberg | Delta Lake | Hudi |
|---|---|---|---|
| 存储格式 | Parquet/ORC/Avro | Parquet(Log 记录变更) | Parquet/Avro |
| 时间旅行 | ✅ | ✅ | 部分(基于 commit) |
| 增量读取 | ✅ 快照 diff | ✅ Change Data Feed | ✅ Incremental View |
| 文件布局优化 | rewrite_data_files | OPTIMIZE | clustering |
| 批/流 Upsert | MERGE INTO | MERGE | COW / 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 查询 | ✅ | 交互式分析 |
| Impala | Hive 生态兼容 | 部分 | 迁移成本高 |
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;
5.3 Flink 流式写入 Iceberg
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/Merge | Spark 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 界文章)
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。