AWS 大数据深度剖析(第一部分):数据湖、数据仓库与湖仓一体革命
理解大数据的核心概念——数据湖 vs. 数据仓库 vs. 湖仓一体,OLTP vs. OLAP,以及为什么现代分析架构都在向 S3 收敛。
本章没有代码——只有思维模型。一旦你把这些概念内化于心,之后遇到的每一个 AWS 服务都会自动归位到大数据宇宙中它应有的位置。
为什么会出现”大数据”这个东西
回想一下你写过的最早的后端代码:一张 MySQL users 表加一张 orders 表,业务就跑得美滋滋。
然后某一天,产品经理对你说:
“我要最近 30 天的日活跃用户数,按城市和设备型号拆分,并排除通过营销活动获取的用户——明天就要。”
你写下 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%,业务团队则收到了雪片般的”下单失败”客诉。
这正是”大数据”被创造出来所要解决的问题:
- 数据量大到单个数据库扛不住(从数亿行到数十 PB)
- 分析查询和业务事务必须分开运行——否则它们会争抢资源、相互拖垮
- 数据格式异构:MySQL 行、Elasticsearch 全文索引、DocumentDB 文档、JSON 事件日志……你需要一个地方把它们全部汇聚起来
- 机器学习需要访问数据:ML 工程师为了拿到一份训练样本,需要扫描数千万行数据——SQL 太慢,他们需要 Spark 能直接读取的 Parquet 文件
大数据生态里的每一项技术——数据湖、Parquet、Iceberg、Spark、Flink、Athena、Glue、SageMaker——都是在回应这四个挑战中的一个或多个。
OLTP vs. OLAP:两种数据库哲学
这是你需要内化的第一个概念分野。
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 章)。
数据仓库、数据湖、湖仓一体:三代架构
这是大数据架构演进的主线。每一代都在解决上一代的痛点。
第一代:数据仓库(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. 流处理
第二个需要区分的思维模型:数据是攒成批一起处理,还是随到随处理、逐条消费?
批处理
特征:
- 数据先在某处攒起来(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' 把变更导出来。”
这个直觉是错的。下图解释了原因:
错误做法:周期性 SELECT 轮询
-- Run every 5 minutes
SELECT * FROM orders WHERE updated_at > '2026-05-10 14:00:00';
问题:
- 给主库带来压力:全表扫描把主库 CPU 顶到 100%,拖垮生产流量
- 无法捕获 DELETE:一行一旦被删除,它的 updated_at 也随之消失
- 依赖应用维护的字段:每次更新真的都会改 updated_at 吗?应用真的维护对了吗?
- 延迟受限于轮询间隔:要做到秒级延迟,你就得每秒查询一次——本质上是在对自己发动 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
数据一旦落进湖里,你不能就把原始数据扔在那儿等着被查询。你必须做分层加工,原因有三:
- 查询性能:原始数据有冗余字段和嵌套结构,直接查很慢
- 可复用性:像 DAU 这样的指标应该只算一次、供 100 个仪表盘消费,而不是每个都重算一遍
- 数据治理:清洗、维度补全、去重和指标口径对齐,应该在一个统一的层里完成
经典四层模型
| 层 | 全称 | 用途 | 实际示例 |
|---|---|---|---|
| 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 内返回个性化推荐。在这次请求内部:
- 查询用户特征(年龄、城市、近期兴趣标签)——不能查数据仓库;必须从 DynamoDB 这样的 KV 存储做点查
- 拉取召回候选集(这个用户可能感兴趣的 1000 个物品)——来自 DynamoDB / OpenSearch
- 排序模型对全部 1000 个候选打分——SageMaker Endpoint
- 返回 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 |
参考资料
- Apache Hadoop — Apache
- Apache Spark documentation — Apache
- AWS Well-Architected Framework — AWS Documentation