# AWS 大数据深度剖析（第二部分）：S3、Parquet 与 Apache Iceberg 详解

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

- 作者: zhuermu
- 发布: 2025-05-11
- 网页版: https://zhuermu.com/blog/bigdata-deep-dive-part-2-storage-formats/

---
> 本章深入讲解数据湖的两个基础主题：
> 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 列存](/images/blog/bigdata-deep-dive/02-row-vs-columnar.svg)

### 一个具体的例子

设想一张 `events` 表，1 亿行、50 列。你执行：

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

**行存（CSV / MySQL）**：
- 为了取出 `city` 和 `amount`，必须读取每一行全部 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）

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

```sql
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_type`、`app_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 Iceberg** | Netflix，后捐给 Apache | **AWS 原生一等公民支持** | 设计严谨，Schema 演进无痛 |
| Apache Hudi | Uber，后捐给 Apache | 良好 | 写友好（Merge-on-Read） |
| Delta Lake | Databricks | 有限 | 在 Databricks 生态最强；开源版本功能较少 |

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

---

## Apache Iceberg 深度剖析

![Iceberg 内部机制](/images/blog/bigdata-deep-dive/02-iceberg-internals.svg)

### 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 天）。

```sql
-- 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 时间旅行的常见误解

许多文章会建议：
```sql
-- 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` 设计成**按天分区**的表，每天保存一份全量快照：

```sql
-- 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（精确到秒）**

```sql
-- 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` 回滚
- **表级审计**：对比「昨天午夜的表」与「今天午夜的表」
- **可复现训练**：钉住一个快照，半年后你依然能产出完全一致的训练数据集

```sql
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 无需重写已有数据即可处理这一切：

```sql
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 表时代是不可能做到的。

### 实战用法

```sql
-- 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，压缩效果更佳 |
| Parquet | AWS 数据湖的标准文件格式 |
| 谓词下推 | Parquet 列统计信息让引擎跳过整个文件 —— 对性能至关重要 |
| Iceberg | 位于 S3 文件之上的「表格式」层，增加了 ACID、UPDATE/DELETE 和时间旅行 |
| 时点正确性 | 推荐系统必须使用每日快照分区或缓慢变化维 Type 2，以避免特征泄漏 |

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

---

## 常见问题

### 大数据场景下为什么用 Parquet 而不是 CSV 或 JSON？

Parquet 是列式格式，只读取需要的列，支持谓词下推，相比 CSV 能实现 5-10 倍压缩。在 Parquet 上执行分析查询扫描的数据量要少得多，同时降低了耗时和成本。

### Apache Iceberg 在 Parquet 文件之上增加了什么？

Iceberg 为存储在 S3 上的文件增加了 ACID 事务、UPDATE/DELETE 支持、Schema 演进、时间旅行和分区演进能力 —— 无需搬移数据即可把数据湖升级为湖仓。


---

## 参考资料

- [Apache Parquet](https://parquet.apache.org/) — Apache
- [Apache ORC](https://orc.apache.org/) — Apache
- [Apache Avro](https://avro.apache.org/) — Apache
- [Apache Iceberg](https://iceberg.apache.org/) — Apache
