AWS 大数据深度剖析(第二部分):S3、Parquet 与 Apache Iceberg 详解

掌握现代数据湖的存储基石 —— S3 对象存储、Parquet 列式格式,以及 Iceberg 如何为 S3 上的文件加上 ACID 事务能力。

zhuermu··14 分钟
big-dataawss3parquetapache-icebergdata-lakecolumnar-storage

本章深入讲解数据湖的两个基础主题:

  1. 数据存在哪里(S3)
  2. 数据采用什么格式(Parquet 列式存储 + Iceberg 表格式)

这两项决策是所有上层服务的基石 —— Athena、Glue、SageMaker 等等,无不建立于此。


Amazon S3:数据湖的基石

它是什么

S3(Simple Storage Service)是 AWS 最早、也最重要的服务(2006 年发布)。其核心是一个全球分布式对象存储:你上传任意大小的文件(单个对象最大 5 TB),为它分配一个 key(路径),之后就能凭这个 key 取回它。

s3://my-bucket/warehouse/ods/orders/dt=2026-05-10/part-0001.parquet
   |           |                            |             |
   bucket      path prefix                  partition dir  file

关键特性(数据湖为何选择 S3)

特性说明对数据湖意味着什么
11 个 9 的持久性每年丢失一个对象的概率约为 ~0.000000001%数据不会丢失
近乎无限的扩展能力单个 bucket 可容纳 EB 级数据并自动分片PB 级数据无需手动分片
按用量付费只为存储的 GB 付费归档历史数据成本极低
强一致性(2020 年起)PUT 之后立即 GET 总能返回最新版本不会出现脏读意外
多种存储类别(Standard / IA / Glacier)冷热分层旧数据自动流转到更便宜的层级
API 友好HTTP REST + AWS SDK每种引擎都能读写

存储类别与冷热分层

S3 并非单一层级 —— 它提供多种存储类别,价格和取回延迟相差几个数量级:

存储类别价格(us-east-1)取回延迟适用场景
S3 Standard~$0.023/GB/月毫秒级当前热数据
S3 Intelligent-Tiering~$0.023/GB/月(自动降层)毫秒到分钟级访问模式未知时的默认推荐
S3 Standard-IA(低频访问)~$0.0125/GB/月毫秒级偶尔访问
S3 Glacier Instant Retrieval~$0.004/GB/月毫秒级每月访问
S3 Glacier Flexible / Deep Archive$0.0036 / $0.00099/GB/月分钟到小时级合规归档

实践建议:把**智能分层(Intelligent-Tiering)**设为所有数仓 bucket 的默认存储类别,让 S3 根据访问频率自动搬移对象。一次配置改动,每年可节省 30-50% 的成本。

S3 是「文件系统」还是「对象存储」?

很多人本能地把 S3 当作文件系统来用。S3 不是文件系统 —— 它是一个只支持整对象读写的键值存储。你无法就地追加或修改其中的任何一个字节。

这个约束驱动了许多上层设计决策:

  • 数据湖文件遵循**一次写入、多次读取(WORM)**的模式
  • 想 UPDATE 某一行?你必须重写整个文件 —— 这恰恰是你需要 Iceberg 来管理这一切的原因(见下文 Iceberg 一节)

关于 S3 Tables(2024 年 12 月 GA,2025-2026 年持续演进)

  • 标准 S3:你能在控制台看到所有 .parquet 文件
  • S3 Tables:AWS 原生的「Iceberg 表存储」—— 控制台展示的是,而不是底层文件。AWS 会自动处理合并(compaction)、过期快照清理和元数据管理

2025-2026 新增能力

  • 跨区域复制(灾备)
  • 支持智能分层(自动冷热分层)
  • Bedrock Knowledge Bases 可直接读取 S3 Tables 进行结构化检索
  • 与 Glue Data Catalog 双向同步

在生产架构中,Zero-ETL to SageMaker Lakehouse 这条路径最终落在 S3 Tables 上。


文件格式:CSV vs JSON vs Parquet vs ORC

S3 是文件存储 —— 但这些文件内部采用什么格式至关重要。

候选格式

格式类型存储效率查询性能可读性
CSV行式文本极佳(人类可读)
JSON / JSON Lines行式文本良好
Avro行式二进制
Parquet列式二进制极高极高
ORC列式二进制极高

结论:对于数据仓库和数据湖,默认选 Parquet

为什么?因为分析型查询主要是选取少数几列、扫描海量行、进行聚合 —— 这正是列式存储的强项。


行存 vs 列存:为什么 OLAP 离不开列式

行存 vs 列存

一个具体的例子

设想一张 events 表,1 亿行、50 列。你执行:

SELECT city, SUM(amount) FROM events WHERE dt='2026-05-01' GROUP BY city;

行存(CSV / MySQL)

  • 为了取出 cityamount,必须读取每一行全部 50 列
  • 尽管只需要 2 列,却读了全部 50 列的数据
  • I/O 浪费:96%

列存(Parquet)

  • city 列连续存放,amount 列同样如此
  • 只读 city + amount,节省 96% 的 I/O
  • 由于一列内所有值类型相同,压缩比极高(例如 city 有大量重复值 —— gzip/snappy 可压缩到 1/10)
  • 现代 CPU 可利用 SIMD 向量化计算(每条指令处理 8 个值)

真实基准测试:同一份数据,CSV 100 GB 用 Snappy 压缩为 Parquet 后约 15 GB;同样的查询在 Athena 上快 5-20 倍,成本降低 80% 以上(按扫描字节数计费)。

Parquet 内部结构(简化版)

+------------------------------------------+
|  File Header (PAR1)                      |
+------------------------------------------+
|  Row Group 1 (~128 MB of rows)           |
|    Column Chunk: user_id  [encoded data] |
|    Column Chunk: city     [encoded data] |
|    Column Chunk: amount   [encoded data] |
|  ...                                     |
+------------------------------------------+
|  Row Group 2                             |
|  ...                                     |
+------------------------------------------+
|  File Footer:                            |
|    schema                                |
|    min/max for each column chunk         | <-- predicate pushdown relies on this
|    compression and encoding info         |
+------------------------------------------+

关键设计决策

  1. Row Group(行组):行被切分成 ~128 MB 的组;每组内部按列存储(在扫描吞吐与随机访问之间取得平衡)
  2. 按列压缩/编码:每一列都可以针对其数据特征选用最优算法(字典编码、RLE、位打包)
  3. 列统计信息:每个 Column Chunk 记录 min/max/null 计数,使查询引擎能够跳过整个 Column Chunk

谓词下推(Predicate Pushdown)

这是列式格式的杀手级特性。考虑:

SELECT * FROM events WHERE user_id = 99999;

当执行引擎读取一个 Parquet 文件时:

  1. 它读取 footer,发现该文件的 user_id 范围是 [100000, 200000]
  2. 整个文件被跳过 —— 一个字节的行数据都不用读

同理:

  • 文件级 min/max 可跳过整个文件
  • Row Group 级 min/max 可跳过整个行组
  • Page 级 min/max 可跳过页

经过层层跳过,引擎实际可能只需物理读取 1% 的数据。

Parquet 最佳实践(必做)

  1. 目标文件大小:128 MB 到 512 MB
    • 太小:文件数量过多,元数据开销大,查询慢
    • 太大:并行度差
  2. dt(日期)分区
    • 路径:.../events/dt=2026-05-10/part-001.parquet
    • WHERE dt='2026-05-10' 直接命中目录,跳过所有其他日期
  3. 按业务维度做二级分区(例如 event_typeapp_id),但要保持基数低(少于几千)—— 否则会引发小文件爆炸
  4. 定期执行合并(compaction),把小文件合并成更大的文件(Glue 内置的作业或 Iceberg 的 OPTIMIZE 命令)

表格式:把 S3 文件夹变成数据库

至此,我们已经让数据湖变得便宜快速。但还剩一个关键缺口:S3 文件不支持 UPDATE 或 DELETE

这为什么是个大问题?看看这些真实场景:

  • ODS 层接收 MySQL CDC 事件,需要 UPSERT(同一个 user_id 到来意味着更新记录)
  • 业务要求「删除某个用户的所有数据」(GDPR / 数据隐私法规)
  • DWD 层需要回补数据,或对特定分区做修 bug 的重写

只用 S3 + Parquet,这些操作全都需要重写整个分区 —— 成本和复杂度直线上升。

解决方案:在 S3 文件之上加一层表格式(table format)。有三个候选:

表格式创建者AWS 集成关键特点
Apache IcebergNetflix,后捐给 ApacheAWS 原生一等公民支持设计严谨,Schema 演进无痛
Apache HudiUber,后捐给 Apache良好写友好(Merge-on-Read)
Delta LakeDatabricks有限在 Databricks 生态最强;开源版本功能较少

结论:在 AWS 上,选 Iceberg。Athena、Glue、EMR、Redshift Spectrum 和 SageMaker 都原生支持它。


Apache Iceberg 深度剖析

Iceberg 内部机制

Iceberg 的分层元数据架构

Iceberg 的关键设计:在数据文件之上,增加一个 Catalog 指针 + 三层元数据文件,每一层都作为对象存储在 S3 上。

Catalog (Glue)                       <-- Layer 0: mutable pointer
    |
    +--points to-->  metadata.json (version v3)       <-- Layer 1: table metadata (schema, partition spec, snapshot list)
                |
                +--points to-->  manifest list      <-- Layer 2: which manifests compose the snapshot
                              |
                              +--points to-->  manifest        <-- Layer 3: min/max + path for each data file
                                          |
                                          +--points to-->  data.parquet (actual data)

每一次写操作:

  1. 写入新的 Parquet 数据文件
  2. 写入一个新的 manifest 来登记这些文件
  3. 写入一个新的 snapshot 来引用这批 manifest
  4. 写入一个新的 metadata.json,把当前快照指向新版本
  5. 更新 Catalog(Glue)指针,指向新的 metadata.json

整个过程是原子的(最后一步是一次单独的 KV 写入)—— 这正是 Iceberg 实现 ACID 保证的方式。

UPDATE / DELETE 如何工作

Iceberg 提供两种策略:

写时复制(Copy on Write,COW) —— 默认策略:

  • 更新一行意味着:读取包含该行的整个 Parquet 文件,修改后写入一个新文件,并把旧文件标记为作废
  • 写慢,读快

读时合并(Merge on Read,MOR)

  • 更新一行意味着:写入一个 delete 文件(「这一行已删除」)外加一个新的数据文件(包含更新后的行)
  • 写快,读时需要合并
  • 适合高频更新场景,但需要定期合并

时间旅行(一项关键能力)

每次写入都会生成一个快照。旧快照引用的文件不会立即删除(保留期可配置,默认 5 天)。

-- Query the table at a specific point in time
SELECT * FROM ads_user_features 
FOR TIMESTAMP AS OF '2026-05-01 00:00:00';

-- Query a specific snapshot version
SELECT * FROM ads_user_features 
FOR VERSION AS OF 2934856;

-- Roll back after an accidental delete
ALTER TABLE ads_user_features 
EXECUTE rollback_to_snapshot(2934856);

为什么推荐系统离不开 Iceberg:时点正确性

**时点正确性(Point-in-Time,PIT)**是推荐系统特征工程中最常见的坑。

问题所在:训练样本必须使用事件发生那一刻的特征值,而不是最新的值。

举例

  • 用户 A 在 5 月 1 日点击了一个视频(正样本)
  • 5 月 1 日,用户 A 的兴趣标签是「美食」
  • 5 月 5 日,用户 A 的兴趣标签被更新为「旅游」(由在线学习更新)
  • 5 月 6 日,你训练模型,从 ads_user_features 取用户 A 的标签 —— 得到的是「旅游」
  • 用「旅游」作为特征去训练「点击了美食视频」这个样本,这是一种特征泄漏(feature leakage) —— 模型学到了错误的模式

关于 Iceberg 时间旅行的常见误解

许多文章会建议:

-- This syntax is NOT valid (Athena/Spark/Trino all reject it)
SELECT ... 
FROM ads_user_features FOR TIMESTAMP AS OF s.event_ts
JOIN ads_sample s ON ...

真相:在 Iceberg / SQL:2011 规范中,FOR TIMESTAMP AS OF 只接受字面常量或绑定参数 —— 不接受列引用。Iceberg 时间旅行无法执行行级的 PIT join —— 它的设计初衷是「把整张表回退到某一个时间点」。

三种正确的 PIT 实现方式

方案 A:每日特征快照分区(推荐,最常用)

ads_user_features 设计成按天分区的表,每天保存一份全量快照:

-- Training sample table contains (user_id, item_id, event_ts, event_dt, label)
SELECT s.label, u.tag, u.age
FROM   ads_sample_follow s
JOIN   ads_user_features_daily u
       ON u.user_id = s.user_id
       AND u.dt    = s.event_dt;   -- align with the day the event occurred

保留 N 天的每日分区(用 Iceberg 的 expire_snapshots + 分区保留策略来控制成本)。

方案 B:缓慢变化维 Type 2(精确到秒)

-- ads_user_features_history(user_id, tag, age, valid_from, valid_to)
-- Each feature change inserts a new row
SELECT s.label, u.tag
FROM   ads_sample_follow s
JOIN   ads_user_features_history u
       ON s.user_id = u.user_id
       AND s.event_ts >= u.valid_from
       AND s.event_ts <  u.valid_to;

方案 C:SageMaker Feature Store

Feature Store 提供了内置的 PIT 检索 API(get_record(record_id, event_time))。其底层使用 Iceberg + 事件时间索引 —— AWS 替你处理了方案 A/B 的复杂性。

Iceberg 时间旅行真正擅长什么

尽管它无法做行级 PIT join,时间旅行在以下场景中极其有用:

  • 数据回滚:误操作 UPDATE 或 DELETE 后,用 rollback_to_snapshot 回滚
  • 表级审计:对比「昨天午夜的表」与「今天午夜的表」
  • 可复现训练:钉住一个快照,半年后你依然能产出完全一致的训练数据集
SELECT * FROM ads_user_features 
FOR TIMESTAMP AS OF TIMESTAMP '2026-05-01 00:00:00';   -- literal constant

SELECT * FROM ads_user_features 
FOR VERSION AS OF 2934856;                              -- snapshot id

Schema 演进

新增列、重命名列、调整列顺序 —— Iceberg 无需重写已有数据即可处理这一切:

ALTER TABLE events ADD COLUMN device_id STRING;             -- add column
ALTER TABLE events RENAME COLUMN ip TO client_ip;            -- rename
ALTER TABLE events ALTER COLUMN amount TYPE DECIMAL(20,4);   -- widen type (compatible direction)

旧的 Parquet 文件原封不动。新列从旧文件读取时为 NULL,被重命名的列继续正常工作。这在 Hive 表时代是不可能做到的。

实战用法

-- Create an Iceberg table in Athena
CREATE TABLE poc_social_layla.ods_event (
  event_id   STRING,
  user_id    BIGINT,
  event_type STRING,
  event_ts   TIMESTAMP,
  payload    STRING,
  dt         STRING
)
PARTITIONED BY (dt)
LOCATION 's3://my-bucket/warehouse/ods/event/'
TBLPROPERTIES (
  'table_type' = 'ICEBERG',
  'format'     = 'parquet',
  'write_compression' = 'snappy'
);

-- UPDATE / DELETE / MERGE just like a traditional database
UPDATE ods_event SET event_type = 'view' WHERE event_type = 'expo';

DELETE FROM ods_event WHERE user_id = 99 AND dt = '2026-05-10';

MERGE INTO ods_event t
USING staging_event s ON t.event_id = s.event_id
WHEN MATCHED THEN UPDATE SET payload = s.payload
WHEN NOT MATCHED THEN INSERT VALUES (s.*);

物理目录布局:一个生产架构

s3://my-bucket/warehouse/
+-- ods/
|   +-- ods_user/        <-- mirror of MySQL users table (CDC)
|   +-- ods_event/       <-- raw event stream (Firehose landing)
|   +-- ods_post/        <-- MySQL posts table
+-- dwd/
|   +-- dwd_user_action/ <-- events joined with user/IP dimensions
|   +-- dwd_post/        <-- enriched post details
+-- dws/
|   +-- dws_user_daily/  <-- daily aggregated user metrics
+-- ads/
|   +-- ads_user_features/    <-- recommendation feature wide table
|   +-- ads_sample_follow/    <-- follow-event training samples
|   +-- ads_recall_u2u_cf/    <-- collaborative filtering recall pool
+-- athena-results/      <-- Athena query result staging

每个目录都对应一张 Iceberg 表。每张表的元数据文件都注册在 Glue Data Catalog 中。


本章小结

概念一句话总结
S3数据湖的物理基石 —— 容量近乎无限,按 GB 付费,11 个 9 的持久性
行存 vs 列存OLAP 离不开列式:节省 80% 以上 I/O,压缩效果更佳
ParquetAWS 数据湖的标准文件格式
谓词下推Parquet 列统计信息让引擎跳过整个文件 —— 对性能至关重要
Iceberg位于 S3 文件之上的「表格式」层,增加了 ACID、UPDATE/DELETE 和时间旅行
时点正确性推荐系统必须使用每日快照分区或缓慢变化维 Type 2,以避免特征泄漏

下一篇:数据如何从源系统流入 S3。

参考资料

  1. Apache Parquet — Apache
  2. Apache ORC — Apache
  3. Apache Avro — Apache
  4. Apache Iceberg — Apache