数据智能

存算分离 + 开放表格式,是现在最稳的数据架构选择

数据只存一份在 S3,Athena、Redshift、EMR、Spark 都能直接读。不再为了换查询引擎搬一次数据。

湖仓分层架构

消费层

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 表创建与增量合并

sql iceberg_merge.sql
-- 在 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

下一步

把云的复杂度交给我们,你只管做业务

留下需求,沐杉云的解决方案架构师会在一个工作日内联系你,提供免费的现状评估、迁移方案草案与 TCO 测算表。