# AWS 大数据深度剖析（第一部分）：数据湖、数据仓库与湖仓一体革命

> 理解大数据的核心概念——数据湖 vs. 数据仓库 vs. 湖仓一体，OLTP vs. OLAP，以及为什么现代分析架构都在向 S3 收敛。

- 作者: zhuermu
- 发布: 2025-05-10
- 网页版: https://zhuermu.com/blog/bigdata-deep-dive-part-1-fundamentals/

---
> 本章没有代码——只有思维模型。一旦你把这些概念内化于心，之后遇到的每一个 AWS 服务都会自动归位到大数据宇宙中它应有的位置。

---

## 为什么会出现"大数据"这个东西

回想一下你写过的最早的后端代码：一张 MySQL `users` 表加一张 `orders` 表，业务就跑得美滋滋。

然后某一天，产品经理对你说：

> "我要**最近 30 天的日活跃用户数，按城市和设备型号拆分，并排除通过营销活动获取的用户**——明天就要。"

你写下 SQL：

```sql
SELECT dt, city, model, COUNT(DISTINCT user_id)
FROM   user_activity_log              -- this table already has 5 billion rows
WHERE  dt BETWEEN '2026-04-10' AND '2026-05-10'
  AND  user_id NOT IN (SELECT user_id FROM marketing_users)
GROUP BY dt, city, model;
```

你提交了查询。MySQL 磨了 4 个小时，把主库 CPU 顶到 100%，业务团队则收到了雪片般的"下单失败"客诉。

**这正是"大数据"被创造出来所要解决的问题：**

1. **数据量大到单个数据库扛不住**（从数亿行到数十 PB）
2. **分析查询和业务事务必须分开运行**——否则它们会争抢资源、相互拖垮
3. **数据格式异构**：MySQL 行、Elasticsearch 全文索引、DocumentDB 文档、JSON 事件日志……你需要一个地方把它们全部汇聚起来
4. **机器学习需要访问数据**：ML 工程师为了拿到一份训练样本，需要扫描数千万行数据——SQL 太慢，他们需要 Spark 能直接读取的 Parquet 文件

大数据生态里的每一项技术——数据湖、Parquet、Iceberg、Spark、Flink、Athena、Glue、SageMaker——都是在回应这四个挑战中的一个或多个。

---

## OLTP vs. OLAP：两种数据库哲学

这是你需要内化的第一个概念分野。

![OLTP vs OLAP](/images/blog/bigdata-deep-dive/01-oltp-vs-olap.svg)

### OLTP（联机事务处理）

**面向事务。** 每次操作只触及少数几行，但要求**极致的速度、强一致性和完整的 ACID 保证。**

- 你打开一个外卖 App 下单：一条 INSERT 写入订单，一条 UPDATE 扣减库存，一条 UPDATE 扣除余额——三条 SQL 在 100ms 内完成
- 用户表、关注表、点赞表——都是 OLTP 工作负载
- 代表：**MySQL、PostgreSQL、Aurora、MongoDB、DocumentDB**

特征：
- **行式存储**：一行的所有列连续存放（读取整行很快）
- **规范化 schema**：避免冗余；多表 JOIN 是常态
- **索引**：B+Tree 支持点查
- 数据量：单库通常在 GB 到低量级 TB

### OLAP（联机分析处理）

**面向分析。** 扫描数十亿行做聚合；不要求毫秒级延迟——重要的是**吞吐量**。

- "过去 30 天每个城市每天的 GMV 是多少？"——就是这类查询
- 代表：**Athena、Redshift、Snowflake、BigQuery、Spark SQL**

特征：
- **列式存储**：每一列独立存放（计算 SUM(amount) 只读 amount 这一列，不碰其他 99 列）
- **反规范化 schema**：100+ 列的宽表很正常；尽量避免 JOIN
- 数据量：TB 到 EB

### 为什么不能用一个数据库同时干这两件事

不是说做不到——而是这两种工作负载在性能画像上根本互不兼容：

| | OLTP | OLAP |
|---|---|---|
| 每次操作涉及行数 | 1-10 | 数千万到数十亿 |
| 期望延迟 | 毫秒级 | 秒级到分钟级 |
| 写入频率 | 高（每次用户操作） | 低（批量导入） |
| 一致性 | 强一致性 | 最终一致性即可 |
| 最优物理存储 | 行式 | 列式 |

如果强行把分析查询压到 MySQL 上，分析会慢，**而且**事务会被拖垮。所以现代架构总是把二者分开：

```
OLTP (business DB)  ──sync──▶  OLAP (data warehouse / data lake)
   MySQL                         S3 + Iceberg + Athena
```

"**如何把数据从 OLTP 搬到 OLAP？**"——这正是 **CDC / DMS / Zero-ETL** 所做的事（见下文 1.5 节和第 03 章）。

---

## 数据仓库、数据湖、湖仓一体：三代架构

这是大数据架构演进的主线。每一代都在解决上一代的痛点。

![Lakehouse vs Warehouse](/images/blog/bigdata-deep-dive/01-lakehouse-vs-warehouse.svg)

### 第一代：数据仓库（1990 年代）

**代表**：Teradata、IBM DB2 Warehouse，以及后来的 Redshift / Snowflake。

**做法**：
- 专用硬件（早期的 MPP——大规模并行处理）
- 存储与计算紧耦合
- 写入前必须先定义 schema（写时模式，schema-on-write）
- 主要服务于 BI 仪表盘和报表

**优势**：查询快、完整 SQL 支持、ACID 事务。
**劣势**：
- 只能存储结构化数据（JSON、视频和日志无法入库）
- ML 访问只能走 JDBC——把数据拉出来很慢
- 存储和计算一起扩容——加存储就得加计算节点（昂贵）
- 厂商锁定

### 第二代：数据湖（2010 年代）

**代表**：Hadoop HDFS、S3 + Hive。

**核心革命**：
- **存算分离**：S3 / HDFS 只负责存储；Spark / Hive 负责计算
- **读时模式（schema-on-read）**：想写什么就写什么（JSON / CSV / Parquet），读的时候再解析
- **低成本存储**：S3 每 GB 只要几分钱

**优势**：
- 存储任意格式
- ML 友好：Spark / Pandas 直接读 Parquet
- 可承载 EB 级数据
- 多个引擎可读取同一份数据

**劣势**：即**"数据沼泽"（Data Swamp）**问题——
- 没有 ACID；无法 UPDATE 或 DELETE 单行
- schema 混乱——没人知道一张表到底有多少列
- 一不小心就产生数百万个小文件，让查询慢到无法忍受
- 治理、审计和访问控制薄弱

### 第三代：湖仓一体（2020 年代至今）

**代表**：Databricks Delta Lake、Apache Iceberg、Apache Hudi。

**核心思想**：在数据湖文件之上加一层**表格式（table format）**，赋予 S3 目录以数据库的能力。

| 数据湖痛点 | 湖仓一体如何解决 |
|---|---|
| 无法 UPDATE/DELETE | Iceberg 维护元数据，追踪"哪些文件仍然有效"；一次 UPDATE 实际上是写入新文件并把旧文件标记为失效 |
| 没有 ACID | Iceberg 用乐观锁 + 元数据快照来实现 ACID |
| 无法查看历史 | 每次写入都会创建一个快照；Time Travel 让你回退到任意历史版本 |
| schema 混乱 | Iceberg 强制约束 schema；Schema Evolution（模式演进）是受控且显式的 |
| 小文件 | Compaction（合并）任务定期把小文件合并 |

**最终结果**：湖的成本 + 仓的体验 + ML 友好性——三者兼得。

> **本参考架构采用湖仓一体模式**：S3（存储）+ Iceberg（表格式）+ Glue Catalog（元数据）+ Athena/EMR/SageMaker（多引擎）。

---

## 批处理 vs. 流处理

第二个需要区分的思维模型：数据是**攒成批一起处理**，还是**随到随处理、逐条消费**？

![Batch vs Stream](/images/blog/bigdata-deep-dive/01-batch-vs-stream.svg)

### 批处理

**特征**：
- 数据先在某处攒起来（S3 / 数据库）
- 按调度触发（每天凌晨 / 每小时整点）
- 一次性处理一大批

**示例**：
- 凌晨 2 点跑昨天的 GMV 报表
- 每天重新训练一次推荐模型
- 数据仓库分层加工：ODS 到 DWD 到 DWS 到 ADS

**典型工具**：EMR Spark、AWS Glue、Athena CTAS、Redshift。

### 流处理

**特征**：
- 数据一到达就立即处理
- 7×24 小时运行
- 状态管理、开窗、乱序事件和水位线（watermark）是日常要处理的问题

**示例**：
- 实时反欺诈（立即拦截可疑登录）
- 实时大屏（双十一 GMV 滚动计数器）
- 推荐系统的实时特征（最近 5 次点击）
- 实时告警

**典型工具**：Flink、Kafka Streams、Spark Streaming、Lambda + Kinesis Data Streams。

### 如何选择

| 业务需求 | 选择 |
|---|---|
| T+1 报表、模型训练 | 批处理 |
| 可接受分钟级延迟 | 批处理（按小时 / 微批） |
| 秒级延迟，且有重算需求 | 流处理 |
| 始终需要"当前最新值" | 流处理 |

**关键原则：能用批就别用流。** 流处理在运维、故障恢复和一致性保障上都要难上一个数量级。先从批开始，验证它能跑通，再考虑上流——这是一条朴素但重要的工程经验法则。

> 在我们的参考架构中，客户的延迟要求是 **T+1**，因此离线管道以批处理为主。只有未来的实时特征管道才需要流处理（Flink）。

---

## CDC：把 OLTP 数据搬进数据湖

我们已经讲完了三代架构以及批处理 vs. 流处理，现在来处理最实际的问题：**MySQL 里的数据到底怎么进 S3？**

直觉给出的答案："写个定时任务，每 5 分钟跑一次 `SELECT * WHERE updated_at > 'last_time'` 把变更导出来。"

这个直觉**是错的**。下图解释了原因：

![CDC Flow](/images/blog/bigdata-deep-dive/01-cdc-flow.svg)

### 错误做法：周期性 SELECT 轮询

```sql
-- Run every 5 minutes
SELECT * FROM orders WHERE updated_at > '2026-05-10 14:00:00';
```

问题：
1. **给主库带来压力**：全表扫描把主库 CPU 顶到 100%，拖垮生产流量
2. **无法捕获 DELETE**：一行一旦被删除，它的 updated_at 也随之消失
3. **依赖应用维护的字段**：每次更新真的都会改 updated_at 吗？应用真的维护对了吗？
4. **延迟受限于轮询间隔**：要做到秒级延迟，你就得每秒查询一次——本质上是在对自己发动 DDoS

### 正确做法：CDC（变更数据捕获）

CDC 的核心思想：**不要去查表——去订阅数据库的复制日志。**

MySQL 内置了一个机制，叫 `binlog`（二进制日志）。它正是 MySQL 用来做主从复制的东西——源库（主库）上每一次 INSERT、UPDATE 和 DELETE 都会写入 binlog，从库读取并回放它。

**CDC 工具实际上做的事**：它把自己伪装成一个 MySQL **从库**，订阅 binlog，逐行解析每一个事件，然后转发到下游。

```
App ─SQL─▶ MySQL source ─binlog─▶ DMS (disguised as replica) ─▶ S3 / Kafka / any downstream
```

优势：
- 对源库压力极小（反正它本来就要复制给真正的从库）
- 完整捕获 INSERT、UPDATE 和 DELETE
- 秒级延迟
- schema 变更也会被捕获

PostgreSQL 用逻辑复制槽（logical replication slot）；MongoDB / DocumentDB 用变更流（change streams）——原理都一样。

> 在 AWS 上，CDC 主要由 **DMS（Database Migration Service）** 和 **Aurora Zero-ETL** 实现。第 03 章会详细讲解它们。

---

## 数据分层：ODS、DWD、DWS、ADS

数据一旦落进湖里，你不能就把原始数据扔在那儿等着被查询。你必须做**分层加工**，原因有三：

1. **查询性能**：原始数据有冗余字段和嵌套结构，直接查很慢
2. **可复用性**：像 DAU 这样的指标应该只算一次、供 100 个仪表盘消费，而不是每个都重算一遍
3. **数据治理**：清洗、维度补全、去重和指标口径对齐，应该在一个统一的层里完成

### 经典四层模型

| 层 | 全称 | 用途 | 实际示例 |
|---|---|---|---|
| **ODS** | Operational Data Store | 原始数据备份层。与源系统一一对应，几乎不做转换 | `ods_users`：MySQL users 表的镜像 |
| **DWD** | Data Warehouse Detail | 明细数据层。清洗 + 标准化 + 维度 JOIN | `dwd_user_action`：事件日志与用户画像 + IP 转地理位置映射关联后的结果 |
| **DWS** | Data Warehouse Summary | 轻度汇总层。按主题/维度组合预聚合 | `dws_user_daily`：每个用户每天的曝光、点赞和关注 |
| **ADS** | Application Data Store | 应用/集市层。直接供下游系统消费的最终产出 | `ads_user_features`：供推荐模型使用的 100 维用户特征宽表 |

### 数据流

```
Source Systems             Data Lake
─────────────             ───────────────────────────────────────────────
MySQL  ─CDC──▶      ods_*  ──transform──▶  dwd_*  ──aggregate──▶  dws_*  ──serve──▶  ads_*
ES     ─OSI──▶                                                                        │
DocDB  ─CDC──▶                                                                        ▼
Events ─Firehose─▶                                                   BI dashboards / ML / online services
```

每一层都由 Iceberg 表构成，且每一层都是通过 SQL（Athena CTAS / Spark SQL）对上一层做转换而产出。编排由 MWAA（托管 Airflow）或 Step Functions 负责。

---

## 离线 vs. 在线：两个完全不同的世界

最后一个、也是最常被混淆的概念。**数据仓库**和**在线服务存储**是两个截然不同的层，各自的职责完全不同。

| 维度 | 离线层（数据仓库） | 在线层（推理服务） |
|---|---|---|
| 存储 | S3 + Iceberg | DynamoDB / Redis / OpenSearch |
| 消费方 | ML 训练 / BI | 面向用户的请求服务 |
| 查询模式 | SQL 批量扫描 | 键值点查 |
| 延迟 | 秒级到分钟级 | **毫秒级** |
| QPS | 几十到几百 | 几万到几十万 |
| 成本模型 | 按存储 + 扫描字节数付费 | 按 QPS + 容量付费 |

### 一个具体的例子

一个用户打开他的社交媒体信息流，App 必须在 50ms 内返回个性化推荐。在这次请求内部：

1. 查询用户特征（年龄、城市、近期兴趣标签）——**不能查数据仓库**；必须从 DynamoDB 这样的 KV 存储做点查
2. 拉取召回候选集（这个用户可能感兴趣的 1000 个物品）——来自 DynamoDB / OpenSearch
3. 排序模型对全部 1000 个候选打分——SageMaker Endpoint
4. 返回 Top 10

**为什么不能直接查数据仓库？**
- 光是 Athena 的启动 + 解析 + 排队开销就有数百毫秒到数秒——200ms 的预算根本容不下
- Athena 按扫描字节数收费；每次推荐扫描几 MB、乘以 10 万 QPS，一天就能把你的预算烧光
- Athena 是 OLAP 分析引擎，不是高 QPS 的 OLTP 服务——它的架构从根本上就不适配高并发点查

所以现代推荐架构总是长这样：

```
Offline warehouse (S3 + Iceberg)       ──daily batch sync──▶      Online storage (DynamoDB + Redis)
ads_user_features                                                  user_features (KV)
ads_recall_pool                                                    recall_candidates (KV)
                                                                          │
                                                                          ▼
                                                            Recommendation service (ms-level response)
```

> **第 06 章**将画出我们参考架构的完整离线管道；**第 08 章**讲解如何在不同的在线存储选项之间做选择。

---

## 本章小结

| 概念 | 一句话总结 |
|---|---|
| OLTP vs. OLAP | 业务数据库 vs. 数据仓库——两种不同的工作负载，物理存储结构从根本上不同 |
| 数据仓库 | 上一代架构：存算耦合，只能存结构化数据，对 ML 不友好 |
| 数据湖 | 存算分离，什么都能存，但没有 ACID——容易沦为沼泽 |
| 湖仓一体 | 数据湖 + 表格式（Iceberg）——两全其美 |
| 批处理 vs. 流处理 | 攒批一起处理 vs. 随到随逐条处理；能用批就别用流 |
| CDC | 订阅数据库 binlog 做实时同步——永远不要轮询 |
| 数据分层 | ODS 到 DWD 到 DWS 到 ADS，每一层都是一张 Iceberg 表 |
| 离线 vs. 在线 | 数据仓库不是在线存储；50ms 的推荐响应查不了 Athena |

---

## 常见问题

### 数据湖和数据仓库有什么区别？

数据湖以低成本存储任意格式的原始数据（结构化、半结构化、非结构化），而数据仓库存储经过清洗、结构化、并针对 SQL 查询优化的数据。数据湖提供灵活性，数据仓库提供性能。

### 什么是湖仓一体（Lakehouse）架构？

湖仓一体将数据湖低成本、读时模式（schema-on-read）的灵活性，与数据仓库的 ACID 事务和 SQL 性能结合在一起，通常是在 S3 这样的对象存储之上使用 Apache Iceberg 等开放表格式来实现。


---

## 参考资料

- [Apache Hadoop](https://hadoop.apache.org/) — Apache
- [Apache Spark documentation](https://spark.apache.org/docs/latest/) — Apache
- [AWS Well-Architected Framework](https://docs.aws.amazon.com/wellarchitected/latest/framework/welcome.html) — AWS Documentation
