湖仓分层架构
消费层
QuickSight 看板
Athena 即席查询
Redshift 数仓
SageMaker 特征
API 服务
语义层
数据集市 Gold
维度建模
指标定义
血缘追踪
加工层
Glue ETL
EMR Spark
dbt 转换
Step Functions 编排
存储层
Bronze 原始
Silver 清洗
Gold 聚合
Iceberg 表格式
接入层
DMS CDC
Kinesis 流式
MSK
批量导入
SaaS 连接器
治理层
Glue Catalog
Lake Formation 权限
数据质量规则
成本归属
Iceberg 提供 ACID 事务、时间旅行和 schema 演进,解决了传统 Hive 表在数据湖上最痛的几个问题。
分层职责
| 层级 | 数据形态 | 保留期 | 主要用途 |
|---|---|---|---|
| Bronze 原始层 | 与源系统完全一致,不做任何转换 | 长期(合规要求) | 问题追溯、重跑加工、审计 |
| Silver 清洗层 | 去重、类型规范、字段标准化、质量过滤 | 1-3 年 | 通用分析、明细查询 |
| Gold 聚合层 | 面向业务的宽表与指标聚合 | 按需 | 看板、报表、API、模型特征 |
Bronze 层千万不要为了省钱删掉。加工逻辑出错时,有原始层就能重跑,没有就只能认损失。
Iceberg 表创建与增量合并
-- 在 Athena 中创建 Iceberg 表(Silver 层订单明细)
CREATE TABLE silver.orders (
order_id STRING,
user_id STRING,
order_status STRING,
total_amount DECIMAL(18, 2),
channel STRING,
created_at TIMESTAMP,
updated_at TIMESTAMP,
dt DATE
)
PARTITIONED BY (dt)
LOCATION 's3://mushan-lakehouse/silver/orders/'
TBLPROPERTIES (
'table_type' = 'ICEBERG',
'format' = 'PARQUET',
'write_compression' = 'ZSTD',
'optimize_rewrite_delete_file_threshold' = '10'
);
-- CDC 增量合并:来自 Bronze 层的变更记录合入 Silver
-- Iceberg 的 MERGE 是原子操作,不会出现读到半成品的情况
MERGE INTO silver.orders AS t
USING (
SELECT
order_id,
user_id,
order_status,
CAST(total_amount AS DECIMAL(18, 2)) AS total_amount,
channel,
created_at,
updated_at,
CAST(created_at AS DATE) AS dt,
op -- I / U / D
FROM (
SELECT *,
ROW_NUMBER() OVER (
PARTITION BY order_id ORDER BY updated_at DESC
) AS rn
FROM bronze.orders_cdc
WHERE ingest_date = CURRENT_DATE
)
WHERE rn = 1 -- 同一主键只取最新变更
) AS s
ON t.order_id = s.order_id
WHEN MATCHED AND s.op = 'D' THEN DELETE
WHEN MATCHED AND s.op IN ('I', 'U') AND s.updated_at > t.updated_at THEN
UPDATE SET
order_status = s.order_status,
total_amount = s.total_amount,
channel = s.channel,
updated_at = s.updated_at
WHEN NOT MATCHED AND s.op IN ('I', 'U') THEN
INSERT (order_id, user_id, order_status, total_amount, channel, created_at, updated_at, dt)
VALUES (s.order_id, s.user_id, s.order_status, s.total_amount, s.channel, s.created_at, s.updated_at, s.dt);
-- 定期维护:合并小文件、清理过期快照
OPTIMIZE silver.orders REWRITE DATA USING BIN_PACK;
VACUUM silver.orders;
小文件是数据湖性能的头号杀手。OPTIMIZE 建议每天跑一次,否则查询会越来越慢。
成本控制要点
- Parquet + ZSTD 压缩,相比 JSON 存储量能降 80% 以上,查询扫描量同步下降
- 按查询模式分区(通常是日期 + 高频过滤字段),避免全表扫描
- Athena 按扫描量计费,做好分区裁剪比升级引擎有效得多
- 冷数据走 S3 生命周期转 Glacier,Iceberg 元数据留在 Standard
- 小文件合并任务纳入日常调度,减少元数据开销与查询延迟
- Redshift 只放需要高并发低延迟的 Gold 层数据,明细查询交给 Athena