# AWS 大数据深度剖析（第 6 部分）：端到端数据管道 —— 从数据源到特征存储

> 把所有环节串联起来：追踪一个点击事件如何从客户端 SDK 经过 API Gateway、MSK、Firehose、S3、数仓分层（ODS→DWD→DWS→ADS），最终写入 DynamoDB 用于实时服务。

- 作者: zhuermu
- 发布: 2025-05-15
- 网页版: https://zhuermu.com/blog/bigdata-deep-dive-part-6-end-to-end-pipeline/

---
第 03 至 05 章分别详细介绍了各个服务。本章将**把它们全部串联起来**，并把整套架构映射到一个真实的社交 App 场景。

本章不引入任何新服务。目标是：读完本章后，你能用一张图讲清楚整个数据侧的架构。

## 客户场景回顾

来自客户方案文档的关键事实：

- **业务**：面向用户推荐的社交 App（关注 / 信息流 / 你可能认识的人）
- **数据源**：MySQL（业务数据库）、ES（搜索）、DocumentDB（文档）、客户端埋点
- **DAU**：100 万+
- **DocumentDB 版本**：5.0（支持 Change Streams）
- **延迟要求**：T+1（离线优先；实时留待后续阶段考虑）
- **现有 Kafka**：无（需从零搭建）

## 端到端架构图

![端到端](/images/blog/bigdata-deep-dive/06-end-to-end.svg)

整条管道分为 5 层：

| 层级 | 职责 |
|---|---|
| 第 1 层：数据源 | Aurora MySQL / DocumentDB / OpenSearch / 客户端埋点 / 第三方数据 |
| 第 2 层：接入通道 | Aurora Zero-ETL / DMS / OpenSearch Ingestion / API GW + MSK + Firehose / EventBridge |
| 第 3 层：S3 + Iceberg ODS | 14+ 张 Iceberg 表，按数据源分组 |
| 第 4 层：分层处理 | DWD / DWS / ADS，由 MWAA 编排 |
| 第 5 层：下游消费方 | BI / ML 训练 / 在线层同步 / 实时管道 |

## 管道拆解：5 条相互独立的数据流

把整条管道拆成 5 条相对独立的子管道，更便于理解：

### 管道 A：业务数据库 CDC（Aurora 到数据湖）

```
Aurora MySQL (orders, users, posts, follows)
     | binlog
     v
Aurora Zero-ETL to SageMaker Lakehouse  (sub-second latency)
     | (AWS-managed, zero ops)
     v
S3 Tables (Iceberg) + auto-registered in Glue Catalog
     |
     v
ods_user / ods_post / ods_follow / ods_order
```

**要点：**
- 选择 Zero-ETL 而非传统 DMS（Zero-ETL 是 Aurora 的推荐路径）
- 数据落在 **S3 Tables**（AWS 原生的 Iceberg 存储）—— 你在控制台看不到底层文件，但 Athena/SageMaker 可以直接读取
- Glue Catalog 自动注册表；schema 随数据源演进

### 管道 B：DocumentDB CDC

```
DocumentDB 5.0 (Change Streams enabled)
     | change streams
     v
DMS Replication Task (source = DocDB, target = S3)
     |
     v
S3 dms-raw/docdb/<collection>/...parquet
     |
     v
Glue Job hourly MERGE → ods_doc_user / ods_doc_msg (Iceberg)
```

**要点：**
- DocumentDB 必须为 4.0+ 版本并启用 Change Streams（会增加源库 I/O）
- DMS 输出的是原始 Parquet，**并非 Iceberg** —— 需要一个 Glue Job 合并进 Iceberg 表
- `s3://.../dms-raw/` 目录是暂存区；生产表位于 `s3://.../warehouse/ods/`

### 管道 C：OpenSearch 到 S3

```
OpenSearch Service (managed)
     | scroll API / PIT
     v
OpenSearch Ingestion Pipeline (yaml)
     |
     v
S3 ods_es_search/ (Parquet)
```

**要点：**
- **先确认是否有必要** —— 如果 ES 只是 MySQL 数据的搜索副本，那么直接从 MySQL 接入更直接
- DMS 不支持 ES 作为数据源
- 对于自建 ES，改用 Logstash + S3 output

### 管道 D：埋点（核心管道，分两个阶段）

> 客户当前的需求是 **T+1** 且没有 Kafka 基础设施。第 1 阶段不应急于上 MSK —— 等到第 3 阶段真正需要实时特征时再说。下面同时展示**简化版（第 1 阶段）**和**目标架构（第 3 阶段及以后）**。

#### 第 1 阶段：简化版（满足 T+1 已足够）

```
Client SDK → API Gateway (HTTP API) → Lambda (auth + enrichment) → Firehose → S3 ods_event
```

特点：
- 完全 Serverless，运维负担极小
- 月成本约 $4K（API Gateway 是最大的开销项）
- 缺点：Firehose 只能投递到单一目的地（S3）；日后转实时需要改造架构

#### 第 3 阶段及以后：目标架构（实时特征上线时）

```
Client SDK
     |
     v
API Gateway (HTTP API, auth)
     |
     v
Lambda (enrichment: server_ts, geo, ip, app_ver)
     |
     v
Amazon MSK (Kafka, topic=events, 12 partitions, across 3 AZs)
     |
     |-- Consumer Group "offline"  --> Firehose --> S3 ods_event (Parquet)
     |-- Consumer Group "realtime" --> Managed Flink --> DynamoDB user_realtime_features
     +-- Consumer Group "risk"     --> Lambda fraud detection
```

**迁移成本**：把 Lambda 的投递目标从 Firehose 改为 MSK；下游挂接以 Kafka 为源的 Firehose（自 2024 年起支持 —— MSK 作为 Firehose 的数据源）。应用层无需任何改动。

**要点：**
- API Gateway 使用 **HTTP API**（比 REST API 便宜约 70%）
- Lambda 不可或缺（鉴权 + 补全 + server_ts）
- **MSK 的价值**：多订阅者的消息总线 —— 同一份数据被 N 个下游消费方独立消费，且支持数据回放
- Schema 通过 Glue Schema Registry 管理（一旦埋点超过 50+ 种事件类型，这就是必需的）

下面的时间线图展示了双路径消费：

![时间线](/images/blog/bigdata-deep-dive/06-data-timeline.svg)

同一个点击事件：**500ms 内落入 DynamoDB** 供推荐服务使用，**60 秒内落入 S3** 供模型训练使用。

### 管道 E：第三方数据

```
EventBridge Schedule (hourly / daily)
     |
     v
Lambda (calls third-party HTTP APIs)
     |
     v
S3 ods_3rd_channel/dt=.../*.json
```

**要点：**
- 使用 EventBridge 而非 cron 或 Lambda Scheduled —— 更规范
- 失败必须重试；建议采用小批量 + 幂等写入

## 分层处理（每日例行任务）

![编排](/images/blog/bigdata-deep-dive/05-orchestration.svg)

由 MWAA 的 Airflow DAG 编排，每日凌晨运行：

```
02:00 Wait for Zero-ETL / DMS daily data completeness (sentinel task)
02:30 |- Glue Job: dwd_user_action_clean       (cleanse + IP→geo + join user dim)
      +- Glue Job: dwd_post_enrich              (post + tags + engagement counts)
03:00 |- Athena CTAS: dws_user_daily            (user daily metric aggregation)
      |- Athena CTAS: dws_post_daily            (post daily metrics)
      +- Athena CTAS: dws_pair_interaction      (user-pair interaction)
03:30 |- EMR Serverless: ads_user_features     (user 100+ dimension feature wide table)
      |- EMR Serverless: ads_post_features     (content features)
      |- EMR Serverless: ads_sample_follow     (follow prediction samples, PIT-correct)
      +- EMR Serverless: ads_recall_u2u_cf     (collaborative filtering recall pool)
04:30 |- Glue Job: sync_user_features_to_ddb   (write to DynamoDB)
      +- Glue Job: sync_recall_pool_to_ddb     (write to DynamoDB)
05:00 |- SageMaker Training Job: train_recall_two_tower
      +- SageMaker Training Job: train_rank_lightgbm
06:00 Lambda: deploy SageMaker Endpoint (blue/green or canary)
06:30 DQ Check + report (success / failure → Slack / email)
```

每一步：
- 失败最多重试 2 次
- 仍然失败 → 告警值班工程师
- 整个 DAG 的 SLA：7 小时

## 数据资产清单

以下是该客户场景的核心 Iceberg 表清单：

### ODS 层（数据源镜像）

| 表 | 数据源 | 主键 | 分区 |
|---|---|---|---|
| ods_user | Aurora users | user_id | dt |
| ods_post | Aurora posts | post_id | dt, hr |
| ods_follow | Aurora follows | (follower_id, followee_id) | dt |
| ods_doc_user | DocumentDB user_profile | user_id | dt |
| ods_doc_msg | DocumentDB messages | msg_id | dt, hr |
| ods_event | 埋点（MSK→Firehose） | event_id | dt, hr, event_type |

### DWD 层（明细 —— 清洗与补全）

| 表 | 说明 |
|---|---|
| dwd_user_action | 事件 + 用户维度 + IP→geo |
| dwd_post | 帖子明细 + 标签 + 互动计数 |
| dwd_user_relation | 关注关系的缓慢变化维度（SCD）表 |

### DWS 层（汇总 —— 轻量聚合）

| 表 | 说明 |
|---|---|
| dws_user_daily | 用户日级：曝光、点击、关注、停留时长 |
| dws_post_daily | 帖子日级：曝光、点击、点赞、分享 |
| dws_pair_interaction | 用户对互动累计值（用于 U2U 协同过滤） |

### ADS 层（应用 / 数据集市）

| 表 | 说明 | 用途 |
|---|---|---|
| ads_user_features | 用户 100+ 维特征宽表 | 推荐特征 + 在线同步 |
| ads_post_features | 内容特征 | 推荐特征 |
| ads_sample_follow | 关注预测样本表（label + 特征 PIT 快照） | 模型训练 |
| ads_sample_ctr | 点击率预测样本表 | 模型训练 |
| ads_recall_u2u_cf | 协同过滤召回池（用户 → top-K 候选） | 在线召回同步 |

## 关键设计决策回顾

回过头审视整套架构，每一个决策都有其依据：

| 决策 | 依据 |
|---|---|
| 使用 Lakehouse（S3 + Iceberg）而非 Redshift | 对 ML 友好 + 成本更优 + 可扩展 |
| Aurora 使用 Zero-ETL 而非 DMS | Aurora 原生路径，亚秒级延迟，零运维 |
| DocumentDB 使用 DMS | DocumentDB 没有 Zero-ETL |
| 埋点：API GW + MSK + Firehose | MSK 多订阅同时支持离线和实时 |
| ETL：以 EMR Serverless 为主 + Athena CTAS 处理轻量作业 | 最优成本组合 |
| 编排：MWAA | 任务依赖复杂；Airflow 表达能力强 |
| 数据格式：Parquet + Iceberg | 列式压缩 + ACID + Time Travel + PIT |
| Schema 管理：Glue Schema Registry | 50+ 种事件类型时必需 |

## 原始架构 vs. 改进架构对比

| # | 客户的原始做法 | 改进后的做法 | 改进点 |
|---|---|---|---|
| 1 | MySQL → S3（未指明方式） | **Aurora Zero-ETL → Lakehouse** | 推荐路径，亚秒级，零运维 |
| 2 | ES → S3 使用 DMS | **OpenSearch Ingestion** | DMS 不支持 ES 作为数据源 |
| 3 | DocumentDB → S3 | DMS + DocumentDB Change Streams | 必须先启用 Change Streams；注意成本影响 |
| 4 | S3 + Athena 数据仓库 | **加入 Iceberg + Parquet + 分区** | 解决 UPDATE / DELETE / PIT 需求 |
| 5 | ODS → DWD 分层 | **明确的 4 层 + 引擎组合 + Iceberg PIT** | PIT 正确性对推荐场景是强制要求 |
| 6 | API GW → Firehose → S3 | **API GW → Lambda → MSK → Firehose** | MSK 支持多订阅；为实时管道预留能力 |

## 本章小结

至此，**数据侧**架构已完整覆盖。你现在应该能够：

- 画出完整的数据管道图
- 讲清楚每条管道中每个服务的作用
- 为每个决策阐述"为什么选 X 而非 Y"
- 把一切映射到客户的实际数据资产清单

接下来的四章将进入 **ML 侧**（推荐系统）：

- 第 07 章：推荐系统基础（召回 / 排序 / 特征 / 模型）
- 第 08 章：在线特征与召回存储（DynamoDB / Redis / OpenSearch kNN / Neptune）
- 第 09 章：SageMaker 与 ML 平台（Feature Store / Training / Endpoint）
- 第 10 章：完整的端到端架构 + 成本估算

---

## 常见问题

### ODS、DWD、DWS 和 ADS 数仓分层分别是什么？

ODS（操作数据存储）保存原始接入的数据。DWD（明细层）对其进行清洗和补全。DWS（汇总层）聚合成指标。ADS（应用层）构建宽表特征和 ML 样本，供下游直接消费。

### 单个点击事件如何流经整条数据管道？

客户端 SDK → API Gateway → Lambda → MSK topic → 同时被 Firehose（→ S3/Iceberg 离线）和 Flink（→ DynamoDB 实时特征）消费。同一个事件同时服务于训练和推理。


---

## 参考资料

- [Apache Airflow documentation](https://airflow.apache.org/docs/) — Apache
- [AWS Glue Developer Guide](https://docs.aws.amazon.com/glue/latest/dg/what-is-glue.html) — AWS Documentation
- [Apache Iceberg](https://iceberg.apache.org/) — Apache
