AWS 大数据深度剖析(第二部分):S3、Parquet 与 Apache Iceberg 详解
掌握现代数据湖的存储基石 —— S3 对象存储、Parquet 列式格式,以及 Iceberg 如何为 S3 上的文件加上 ACID 事务能力。
本章深入讲解数据湖的两个基础主题:
- 数据存在哪里(S3)
- 数据采用什么格式(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 离不开列式
一个具体的例子
设想一张 events 表,1 亿行、50 列。你执行:
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 |
+------------------------------------------+
关键设计决策:
- Row Group(行组):行被切分成 ~128 MB 的组;每组内部按列存储(在扫描吞吐与随机访问之间取得平衡)
- 按列压缩/编码:每一列都可以针对其数据特征选用最优算法(字典编码、RLE、位打包)
- 列统计信息:每个 Column Chunk 记录 min/max/null 计数,使查询引擎能够跳过整个 Column Chunk
谓词下推(Predicate Pushdown)
这是列式格式的杀手级特性。考虑:
SELECT * FROM events WHERE user_id = 99999;
当执行引擎读取一个 Parquet 文件时:
- 它读取 footer,发现该文件的
user_id范围是[100000, 200000] - 整个文件被跳过 —— 一个字节的行数据都不用读
同理:
- 文件级 min/max 可跳过整个文件
- Row Group 级 min/max 可跳过整个行组
- Page 级 min/max 可跳过页
经过层层跳过,引擎实际可能只需物理读取 1% 的数据。
Parquet 最佳实践(必做)
- 目标文件大小:128 MB 到 512 MB
- 太小:文件数量过多,元数据开销大,查询慢
- 太大:并行度差
- 按
dt(日期)分区- 路径:
.../events/dt=2026-05-10/part-001.parquet WHERE dt='2026-05-10'直接命中目录,跳过所有其他日期
- 路径:
- 按业务维度做二级分区(例如
event_type、app_id),但要保持基数低(少于几千)—— 否则会引发小文件爆炸 - 定期执行合并(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 的分层元数据架构
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)
每一次写操作:
- 写入新的 Parquet 数据文件
- 写入一个新的 manifest 来登记这些文件
- 写入一个新的 snapshot 来引用这批 manifest
- 写入一个新的
metadata.json,把当前快照指向新版本 - 更新 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,压缩效果更佳 |
| Parquet | AWS 数据湖的标准文件格式 |
| 谓词下推 | Parquet 列统计信息让引擎跳过整个文件 —— 对性能至关重要 |
| Iceberg | 位于 S3 文件之上的「表格式」层,增加了 ACID、UPDATE/DELETE 和时间旅行 |
| 时点正确性 | 推荐系统必须使用每日快照分区或缓慢变化维 Type 2,以避免特征泄漏 |
下一篇:数据如何从源系统流入 S3。
参考资料
- Apache Parquet — Apache
- Apache ORC — Apache
- Apache Avro — Apache
- Apache Iceberg — Apache