AWS 大数据深度剖析(第一部分):数据湖、数据仓库与湖仓一体革命

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

zhuermu··12 分钟
big-dataawsdata-lakedata-warehouselakehouseoltpolap

本章没有代码——只有思维模型。一旦你把这些概念内化于心,之后遇到的每一个 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%,业务团队则收到了雪片般的”下单失败”客诉。

这正是”大数据”被创造出来所要解决的问题:

  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

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

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

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

OLTPOLAP
每次操作涉及行数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

第一代:数据仓库(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/DELETEIceberg 维护元数据,追踪”哪些文件仍然有效”;一次 UPDATE 实际上是写入新文件并把旧文件标记为失效
没有 ACIDIceberg 用乐观锁 + 元数据快照来实现 ACID
无法查看历史每次写入都会创建一个快照;Time Travel 让你回退到任意历史版本
schema 混乱Iceberg 强制约束 schema;Schema Evolution(模式演进)是受控且显式的
小文件Compaction(合并)任务定期把小文件合并

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

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


批处理 vs. 流处理

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

Batch vs Stream

批处理

特征

  • 数据先在某处攒起来(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

错误做法:周期性 SELECT 轮询

-- 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. 数据治理:清洗、维度补全、去重和指标口径对齐,应该在一个统一的层里完成

经典四层模型

全称用途实际示例
ODSOperational Data Store原始数据备份层。与源系统一一对应,几乎不做转换ods_users:MySQL users 表的镜像
DWDData Warehouse Detail明细数据层。清洗 + 标准化 + 维度 JOINdwd_user_action:事件日志与用户画像 + IP 转地理位置映射关联后的结果
DWSData Warehouse Summary轻度汇总层。按主题/维度组合预聚合dws_user_daily:每个用户每天的曝光、点赞和关注
ADSApplication 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 + IcebergDynamoDB / 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

参考资料

  1. Apache Hadoop — Apache
  2. Apache Spark documentation — Apache
  3. AWS Well-Architected Framework — AWS Documentation