Lakehouse 架构实战
概述
实时湖仓一体把数仓分层搬到湖上:Iceberg 提供事务与版本,Flink 负责实时采集加工,Spark 负责离线批处理,Presto 提供即席查询。本文给出一个可落地的 Iceberg + Spark + Flink 湖仓一体架构:分层设计、全链路实现、Catalog 统一与运维要点。
一、总体架构
┌────────────────────────────────────────────────┐
│ 采集层:Flink CDC → Kafka │
│ 业务库(MySQL)→ binlog → Kafka │
├────────────────────────────────────────────────┤
│ 加工层:Flink(实时) Spark(离线) │
│ Kafka → Iceberg ODS/DWD │
│ Spark 批 → DWS/ADS(聚合宽表) │
├────────────────────────────────────────────────┤
│ 存储层:Iceberg 表(HDFS/S3 + Hive Metastore) │
│ ODS → DWD → DWS → ADS(分区/快照) │
├────────────────────────────────────────────────┤
│ 服务层:Presto/Trino → BI(Superset) │
└────────────────────────────────────────────────┘二、分层设计
| 层 | 内容 | 更新频率 |
|---|---|---|
| ODS | 原始明细(CDC 落湖) | 实时追加 |
| DWD | 清洗明细(宽表) | 实时 |
| DWS | 汇总层(按维度聚合) | 分钟级 |
| ADS | 应用层(报表/指标) | 分钟/小时 |
| DIM | 维表(快照/版本) | 低频 |
分层在湖上 = 不同库/表 + 分区 + 快照
流批共用同一份表三、实时采集(Flink CDC → Kafka → Iceberg)
3.1 链路
MySQL binlog → Flink CDC → Kafka → Flink SQL → Iceberg ODS3.2 实现
sql
-- 1. CDC 表(读 MySQL)
CREATE TABLE cdc_orders (...) WITH (
'connector' = 'mysql-cdc', ...);
-- 2. 中间 Kafka
CREATE TABLE kafka_orders (...) WITH (
'connector' = 'kafka', ...);
-- 3. Iceberg ODS
CREATE TABLE lake.ods_orders (...) WITH (
'connector' = 'iceberg',
'catalog-name' = 'lake',
'catalog-type' = 'hive',
'uri' = 'thrift://metastore:9083',
'write.format.default' = 'parquet');
-- 4. 流转
INSERT INTO kafka_orders SELECT ... FROM cdc_orders;
INSERT INTO lake.ods_orders SELECT ... FROM kafka_orders;3.3 实时写注意
| 要点 | 说明 |
|---|---|
| Checkpoint | 开启保证精确一次 |
| 主键 | CDC 变更按主键 upsert |
| 分区 | 按事件时间分区 |
| 小文件 | 定期合并 |
四、离线批处理(Spark)
4.1 批任务
scala
// 读 ODS,加工到 DWD
val ods = spark.table("lake.ods_orders")
.filter("dt = '2026-08-04'")
ods.transform(cleanTransform)
.writeTo("lake.dwd_orders")
.overwritePartitions()4.2 批写语义
| 模式 | 说明 |
|---|---|
| append | 追加 |
| overwritePartitions | 覆盖分区 |
| upsert | 更新插入 |
4.3 调度
离线任务定时(如每 10 分钟/每小时):
读取增量分区 → 加工 → 写目标层
调度:Airflow/DolphinScheduler五、DWS/ADS 聚合
5.1 分钟级聚合
sql
-- 汇总:按分钟 + 维度
INSERT INTO lake.dws_order_min
SELECT dt, hour, minute, city,
SUM(amount) AS gmv, COUNT(*) AS cnt
FROM lake.dwd_orders
WHERE dt = '2026-08-04'
GROUP BY dt, hour, minute, city;5.2 指标存储
| 指标 | 存储 |
|---|---|
| 明细 | Iceberg |
| 汇总 | Iceberg / Doris / ClickHouse |
| 大屏 | Doris/ClickHouse |
建议:
明细与汇总留湖(可追溯)
高并发指标查 → OLAP 引擎六、查询服务
6.1 Presto 即席
sql
SELECT city, SUM(gmv) FROM lake.dws_order_min
WHERE dt = '2026-08-04'
GROUP BY city;6.2 与 BI 集成
Superset/Metabase 连 Presto
报表直接查湖表| 查询优化 | 说明 |
|---|---|
| 分区裁剪 | WHERE 分区列 |
| 快照查询 | 时间旅行 |
| 合并小文件 | 定期 OPTIMIZE |
| Z-Order | 多维过滤优化 |
七、元数据与权限
7.1 Catalog 统一
单一 Hive Metastore:
Flink/Spark/Presto 共用 lake catalog| 组件 | 配置 |
|---|---|
| Hive Metastore | 元数据 |
| Spark | spark.sql.catalog.lake |
| Flink | catalog-type=hive |
| Presto | iceberg catalog |
7.2 权限
Ranger 集成:
表/库级别权限
引擎统一鉴权八、运维要点
8.1 例行维护
| 任务 | 说明 |
|---|---|
| 小文件合并 | OPTIMIZE(每日) |
| 快照清理 | 过期快照(保留 N 天) |
| 孤儿文件清理 | 删除无引用文件 |
| 表健康检查 | 文件数/大小 |
8.2 监控
| 指标 | 说明 |
|---|---|
| 写入延迟 | CDC 到湖延迟 |
| 小文件数 | 合并必要性 |
| 查询耗时 | Presto 性能 |
| 任务失败 | 批任务重试 |
8.3 常见问题
| 问题 | 处理 |
|---|---|
| 流批写冲突 | 错峰 + 乐观并发 |
| 表膨胀 | 快照清理 |
| 查询慢 | 分区/合并优化 |
| 元数据漂移 | 统一 catalog 治理 |
九、架构演进
当前:Iceberg + Flink + Spark + Presto
演进方向:
- Flink 直接写 DWS(减少 Spark 批)
- 数据网格(领域化数据产品)
- 云原生(对象存储 + K8s)常见问题速查
| 问题 | 原因与处理 |
|---|---|
| CDC 到湖有延迟 | 检查 Flink 背压与 checkpoint |
| 小文件爆炸 | 定期 OPTIMIZE |
| 流批冲突 | 错峰调度 |
| 查询不裁剪 | 检查分区过滤条件 |
| 快照过大 | 清理过期快照 |