# AWS 大数据深度剖析（第三部分）：数据接入 —— DMS、Zero-ETL、Firehose 与 MSK

> 四类数据源，四条接入管道 —— 学习使用 DMS 做 CDC、Aurora Zero-ETL、MSK 上的 Kafka，以及 Firehose 微批处理，将数据落地到你的 S3 数据湖。

- 作者: zhuermu
- 发布: 2025-05-12
- 网页版: https://zhuermu.com/blog/bigdata-deep-dive-part-3-data-ingestion/

---
> 客户场景中的核心问题：**MySQL、DocumentDB、Elasticsearch 以及客户端事件埋点** —— 每一类数据源如何落地到 S3？
>
> 本章覆盖四条接入管道以及所涉及的 7 项 AWS 服务（DMS、Zero-ETL、OpenSearch Ingestion、Firehose、MSK、KDS 和 Lambda）。

## 全景概览

![数据接入概览](/images/blog/bigdata-deep-dive/03-data-ingestion.svg)

四类数据源对应四条接入管道：

| # | 数据源 | 推荐管道 | 理由 |
|---|---|---|---|
| 1 | Aurora / RDS MySQL | **Zero-ETL to SageMaker Lakehouse** | 亚秒级延迟、完全托管、AWS 推荐方案 |
| 2 | 自建 MySQL / DocumentDB | DMS 落地到 S3，再由 Glue 合并进 Iceberg | DMS 支持 60 多种异构源 |
| 3 | Elasticsearch / OpenSearch | OpenSearch Ingestion 落地到 S3 | DMS 不支持将 ES 作为源 |
| 4 | 客户端事件埋点 | API Gateway to MSK to Firehose to S3 | 多订阅者扇出，可直接对接实时特征管道 |

我们逐一深入。

## 管道 1：使用 DMS 对业务数据库做 CDC

### 什么是 DMS？

**AWS Database Migration Service** 是 AWS 的托管数据库迁移与同步服务。

它最初为“将本地 Oracle 迁移到 AWS RDS”而设计，如今已演进为通用的**异构数据库同步管道**。

它支持：
- 60 多种数据库作为源（MySQL、PostgreSQL、Oracle、SQL Server、MongoDB、DocumentDB、Redis、DynamoDB、Kafka、S3 等）
- 30 多种端点作为目标（包括 S3、Kinesis 和 OpenSearch）

### 两种工作模式

**Full Load**（全量初始加载）：
- 通过 `SELECT * FROM table` 读取整张表并写入目标
- 大表会导致源库承受数小时的高负载（**这是 DMS 唯一会给源库带来显著压力的阶段**）

**CDC**（增量同步）：
- 把自己伪装成 MySQL **副本**，订阅 binlog
- 将每个 binlog 事件转换后转发给目标
- 对源库压力极小（相当于一个真实的副本）
- 延迟约 1 分钟（可调优至秒级）

实际部署中通常会组合两个阶段：“Full Load + CDC” —— 先跑一次全量初始加载，再切换到 CDC 做持续同步。

### DMS 的“三件套”

配置 DMS 需要创建三个对象：

| 概念 | 用途 |
|---|---|
| **Replication Instance（复制实例）** | 执行同步的“工人”。你需要选择实例规格（起步为 dms.t3.medium），7×24 小时运行，成本约每月 $50 |
| **Endpoint（端点）** | 源与目标的连接信息（主机、凭证、表选择规则） |
| **Replication Task（复制任务）** | 打包“从哪个源端点到哪个目标端点、同步哪些表、使用哪种模式（Full Load / CDC / Full + CDC）” |

### DMS 输出到 S3

写入 S3 时，DMS 输出 **Parquet 文件**（也支持 CSV，但不推荐）。目录结构如下：

```
s3://my-bucket/dms-raw/
└── poc-mysql-source/
    └── orders/
        ├── LOAD00000001.parquet     ← Full Load initialization files
        ├── LOAD00000002.parquet
        └── 20260510-100001234.parquet  ← CDC incremental files (includes Op column: I/U/D)
```

**重要提示**：DMS 输出的是**原始 Parquet**，**并非 Iceberg 表** ——
- 重复行：同一行的多次更新会产生多条记录
- 没有 ACID 保证
- 下游分析直接读取会得到混乱的结果

因此，在 DMS 输出之后，你需要一个 **Glue Job 定期做 MERGE**，合并进规范的 Iceberg ODS 表：

```sql
-- Run inside a Glue Spark Job
MERGE INTO ods_user t
USING (
  SELECT * FROM dms_raw_user 
  WHERE dt = '2026-05-10'
  -- For rows with multiple changes, keep only the latest
  QUALIFY ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY commit_ts DESC) = 1
) s
ON t.user_id = s.user_id
WHEN MATCHED AND s.Op = 'D' THEN DELETE
WHEN MATCHED AND s.Op IN ('U', 'I') THEN UPDATE SET *
WHEN NOT MATCHED AND s.Op IN ('U', 'I') THEN INSERT *;
```

### DMS 的局限

- 你需要自行管理复制实例（CPU/内存监控、扩缩容）
- Schema 变更需要手动重新配置
- 大事务（百万行的 UPDATE）可能引发延迟尖刺

这些痛点催生了下一个方案：**Zero-ETL**。

## 管道 1 升级版：Aurora / RDS Zero-ETL

### 什么是 Zero-ETL？

“Zero-ETL”是 AWS 在 re:Invent 2022 上提出的产品愿景 —— 让 ETL 中的“E（Extract）”和“L（Load）”**无需用户编写代码或管理集群**。

并不是说真的没有 ETL —— 而是 **AWS 帮你全托管了 ETL**。

### Zero-ETL 家族（截至 2026 年 5 月）

| 源 | 目标 | GA 状态 | 适用场景 |
|---|---|---|---|
| Aurora MySQL / PostgreSQL | Redshift | 2023-06 GA | BI 分析 |
| **Aurora MySQL** | **SageMaker Lakehouse**（S3 Tables） | **2025-06 GA** | **数据湖场景 —— 我们的推荐选择** |
| RDS MySQL | Redshift / Lakehouse | 自 2025 年起陆续 GA（主要商用区域已可用；**GovCloud/中国区/部分较新区域仍处于预览或不支持状态** —— 签约前请确认目标区域） | 同上 |
| DynamoDB | Redshift / OpenSearch | 2023+ GA | 全文检索 |
| **DynamoDB to SageMaker Lakehouse** | Lakehouse | 2025 GA（新） | KV 数据入湖 |
| **SaaS（Salesforce/SAP/ServiceNow/Zendesk）to Lakehouse** | Lakehouse | 2024-2025 GA | 跨系统数据集成 |
| 自管 MySQL | S3 / S3 Tables（基于 DMS） | 2025+ GA | 自建 MySQL |

> 本文默认采用 **Aurora MySQL to SageMaker Lakehouse**。如果客户使用 RDS for MySQL，请对照上表确认目标区域的 GA 状态；否则回退到“基于 DMS 的 Zero-ETL”或传统 DMS。

### Zero-ETL to SageMaker Lakehouse 的工作原理

```
Aurora Primary
    │ binlog (still binlog — no black magic)
    ▼
AWS Internal Managed Zero-ETL Service (serverless)
    │
    ▼
S3 Tables (auto-creates Iceberg tables) + auto-registered in Glue Catalog
```

关键差异：

| | DMS（含基于 DMS 的 Zero-ETL） | Zero-ETL to Lakehouse |
|---|---|---|
| 初始化 | `SELECT *` 全表读取 —— 对大表压力大 | **使用 RDS 快照** —— 完全不触碰源库查询引擎 |
| CDC | 订阅 binlog | 订阅 binlog |
| 延迟 | 约 1 分钟 | **亚秒级** |
| 输出 | 原始 Parquet（需进一步处理成 Iceberg） | **直接输出 Iceberg 表** |
| 元数据 | 手动运行 Glue Crawler | **自动注册到 Glue Catalog** |
| 运维 | 需管理复制实例 | **完全 serverless** |
| 计费 | 复制实例 + 数据量 | 变更量 + S3 |

### 选型标准

```
   Aurora / RDS MySQL          → Zero-ETL to SageMaker Lakehouse (recommended)
   Self-hosted MySQL (EC2/IDC) → DMS-based Zero-ETL to S3 Tables
   DocumentDB / other sources  → Traditional DMS → S3 → Glue merge into Iceberg
   Complex ETL / custom logic  → Traditional DMS (preserves full control)
```

**一个常见误解**：很多客户以为“Zero-ETL 比 DMS 好 100 倍”。而实际上：
- **在 CDC 阶段，两者底层都使用 binlog** —— 对源库的压力相当
- 真正的差异在于**初始化阶段**（快照 vs 全表 SELECT）以及**运维模式**
- 当客户问“DMS 会不会定期做 SELECT？”时 —— 不会，在增量阶段 DMS 读取的是 binlog 事件

## 管道 2：DocumentDB 到 S3

### 前提：启用 Change Streams

DocumentDB 是 AWS 的“兼容 MongoDB API”服务（它并不是真正的 MongoDB）。

DocumentDB 不产生 binlog，而是使用 **Change Streams** —— 源自 MongoDB 协议的“操作日志订阅 API”。

DMS 要对 DocumentDB 做 CDC，需要满足以下条件：
- DocumentDB 版本 4.0 或更高
- 集群参数组中启用 `change_stream_log_retention_duration`
- DocumentDB Change Streams **会增加源库的 I/O 和成本**（写入开销增加 10-20%）

### 管道流程

```
DocumentDB (4.0+, Change Streams ON)
    │
    ▼
DMS Replication Task (source endpoint = DocDB, target = S3)
    │
    ▼
S3: dms-raw/docdb/<collection>/...parquet
    │
    ▼
Glue Job MERGE INTO Iceberg
    │
    ▼
ods_doc_user / ods_doc_post
```

关键说明：
- DocumentDB 是文档数据库 —— DMS 在写入 Parquet 前会把 BSON 转换为 JSON/嵌套结构
- 嵌套结构可以在 Athena 中用 `dot notation`（点号表示法）查询：`SELECT user.profile.age FROM ...`

## 管道 3：通过 OpenSearch Ingestion 把 Elasticsearch 导入 S3

### 为什么不能用 DMS

**关键事实**：DMS 不支持将 ES / OpenSearch 作为**源**（只能作为**目标**）。

这是业界常见的坑：客户在方案里写“用 DMS 把 ES 同步到 S3”，直到 POC 阶段才发现根本行不通。

### 什么是 OpenSearch Ingestion（OSI）？

OSI 是 AWS 于 2023 年推出的托管数据接入服务。底层其实是**托管的 Data Prepper**（一款开源工具）。

特性：
- Serverless（按 OCU = OpenSearch Compute Unit 计费）
- 内置 OpenSearch / Elasticsearch 源（使用 scroll API / PIT 做持续读取）
- 内置多种 sink：S3、OpenSearch、Lambda、Kafka

```
OpenSearch Service (managed)
    │ scroll API / Point-in-Time
    ▼
OpenSearch Ingestion Pipeline (YAML configuration)
    │
    ▼
S3 (Parquet / JSON)
```

### 自建 ES 怎么办？

OSI **只支持 AWS 托管的 OpenSearch Service** 作为源。

对于部署在 EC2 或本地的自建 ES，你有以下选择：
1. **Logstash** + S3 输出插件（最常见）
2. 自定义 Lambda / EMR Job，使用 scroll API / PIT 拉取数据
3. 数据量小的话，定期做 `_search` 全量导出

### 但先问一个问题

**ES 里存的到底是什么数据？**

ES 的常见用途：
- 业务数据的搜索副本（原始数据在 MySQL，ES 只用于搜索）—— **直接从 MySQL 接入更好**，别从 ES 拉
- 日志索引（例如 ELK 栈的日志）—— 未必需要进数仓，CloudWatch 或 S3 归档可能就够了
- 只存在于 ES 中的业务数据（少见）—— 必须从 ES 入湖

结论是：**先确认 ES 里的数据是否与 MySQL 重复**，再决定要不要搭这条管道。

## 管道 4：客户端事件埋点 —— 四大组件协同工作

### 基础管道（客户的原始设计）

```
Client SDK ──▶ API Gateway ──▶ Firehose ──▶ S3 (Parquet)
```

这条管道**能用**，但对于**推荐场景来说不够** —— 下面会解释原因。

### API Gateway

托管的 API 网关，分两种类型：

| 类型 | 价格（每百万请求） | 延迟 | 推荐用途 |
|---|---|---|---|
| REST API | $3.5 | 约 30ms | 复杂功能（缓存、用量计划） |
| **HTTP API** | **$1.0** | 约 20ms | 事件埋点 / 简单代理（推荐） |

对于事件埋点，REST API 过于重量级 —— **用 HTTP API**：便宜 70%，延迟更低。

### Lambda（数据富化）

API Gateway 可以直接路由到 Firehose，但**强烈建议在中间加一层 Lambda**：

```
SDK reports: { user_id, event_type, ts, ... }
   ↓
Lambda enrichment:
  + server_ts (server-side timestamp — prevents client clock tampering)
  + ip + geo (resolve IP to geographic location)
  + app_version (extract from User-Agent)
  + authentication (AppKey + HMAC)
  - filter invalid / replayed events
   ↓
Push to downstream
```

为什么**不让客户端直接写入**？因为客户端时间戳不可靠、客户端拿不到 IP 地址，而鉴权必须在服务端完成。

### Amazon Data Firehose（前身为 Kinesis Data Firehose，于 2024 年 2 月更名）

Firehose 是一条“**托管的微批传送带**”：
- 你逐条 PutRecord 投递数据
- 它在内存中缓冲（默认：达到 5 MB 或 60 秒，以先到者为准）
- 自动转换为 Parquet、压缩，并按时间分区写入 S3

特性：
- 完全 serverless（无 broker，无消费端代码）
- 自动重试、背压处理和扩缩容
- **每条投递流只能有一个主目标** —— 这是一个关键限制（虽然可以启用到 S3 的源备份，但那只是故障/审计副本，无法作为独立消费者用于实时回放）

```
Firehose output file naming (dynamic partitioning):
s3://bucket/raw/events/event_type=click/dt=2026-05-10/hr=13/
   firehose-events-1-2026-05-10-13-23-01-xxx.parquet
```

### 为什么仅靠 Firehose 不够：多订阅者问题

回到推荐场景的真实需求：

```
The same event data needs to be consumed by N downstream systems:
  1. Offline data warehouse (every event lands in S3, T+1 model training)
  2. Real-time features (Flink computes last 5 clicks → DynamoDB, millisecond latency)
  3. Real-time fraud detection (Lambda detects anomalous logins → block)
  4. Real-time dashboard (Flink computes GMV → push to frontend)
```

Firehose 的设计本质上是**生产者到单一 sink** —— 每条流只能有一个主目标。要服务 4 个独立消费者，你要么创建 4 条 Firehose 流（数据量 4 倍、成本 4 倍），要么从 S3 回读（会彻底破坏实时性）。这正是为什么你需要在前面加一层 Kafka/KDS 做“扇出”。

### 改进后的架构：在前面加上 MSK/KDS

正确的架构：

```
Client SDK
   ↓
API Gateway (HTTP API)
   ↓
Lambda (auth / enrichment)
   ↓
Amazon MSK (Kafka)            ← Multi-subscriber message bus
   ├──▶ Consumer Group 1: Firehose → S3 Iceberg (offline)
   ├──▶ Consumer Group 2: Managed Flink → DynamoDB (real-time features)
   ├──▶ Consumer Group 3: Lambda → fraud detection
   └──▶ Consumer Group 4: Flink → real-time dashboard
```

## Amazon MSK vs Kinesis Data Streams

![Kafka MSK Topics 与 Partitions](/images/blog/bigdata-deep-dive/03-msk-kafka.svg)

### 什么是 MSK？

**Amazon MSK = Managed Streaming for Apache Kafka** —— 一个由 AWS 托管的 Apache Kafka 集群。

**核心模型**（理解这三个概念，你就理解了 Kafka）：

| 概念 | 类比 |
|---|---|
| **Topic（主题）** | 一个频道（一类数据），例如 `events` / `cdc.user` |
| **Partition（分区）** | 把一个 topic 拆成 N 个分片以提升并行度；同一分区内的消息是有序的 |
| **Consumer Group（消费者组）** | 一组消费者，彼此瓜分分区；不同的组之间完全独立 |

“**多订阅者**”的本质：创建多个消费者组，同一份数据由每个组按各自的节奏独立消费。

### MSK 部署模式

| 模式 | 特点 | 计费 |
|---|---|---|
| **MSK Provisioned** | 你选择 broker 实例类型、可用区和数量；灵活度最高 | 按实例 + 存储 |
| **MSK Serverless** | 自动扩缩容；无需运维，但功能略少 | 按分区 + 吞吐量 |
| **MSK Connect** | 托管的 Kafka Connect，用于运行各类连接器 | 按 worker 实例 |

对于每日数亿事件的埋点场景，采用 Provisioned、跨 3 个可用区部署 3 台 m7g.large broker，成本约每月 $400-500。

### Kinesis Data Streams（KDS）

KDS 是 AWS 自研的“类 Kafka”流服务（比 MSK 更早推出）。

**MSK vs KDS 对比**：

| | MSK | KDS |
|---|---|---|
| 协议 | Apache Kafka | AWS 自研 |
| 生态 | 完整的 Kafka 生态（Flink、Spark、Connect、KSQL） | AWS 原生（Lambda、Firehose、Flink） |
| 学习曲线 | 中等（需要 Kafka 知识） | 低（API 简单） |
| 运维 | Provisioned 需管理 broker；Serverless 无需运维 | 完全 serverless |
| 成本（每日 1 亿事件） | 略低 | 略高 |
| 顺序保证 | 同一分区内有序 | 同一 shard 内有序 |

**如何选择**：
- 团队熟悉 Kafka，或下游使用 Flink → 选 MSK
- 完全 AWS 原生且想要 serverless → 选 KDS
- 当日量级相当（数十亿以内）时，成本相近 —— 按团队熟悉度来选

### Schema Registry（事件埋点的必备项）

有 50 种事件类型、每种字段各不相同 —— 没有 schema 管理，混乱是必然的。

**AWS Glue Schema Registry**：
- 集中注册每个 topic 的 Avro / JSON Schema
- 生产者写入前校验
- 消费者读取时解析
- 支持 Schema Evolution（模式演进，含兼容性检查）

### 完整的改进架构

```
┌─────────────┐
│ Client SDK  │
└──────┬──────┘
       ▼
┌─────────────┐    ┌───────────┐
│ API Gateway │───▶│  Lambda   │ Auth + Enrichment (server_ts, geo, ip)
│ (HTTP API)  │    └─────┬─────┘
└─────────────┘          │
                         ▼
                  ┌─────────────┐
                  │  MSK Topic  │ events (12 partitions, across 3 AZs)
                  │  (events)   │
                  └──┬──┬──┬───┘
       ┌─────────────┘  │  └────────────────┐
       ▼                ▼                   ▼
  ┌─────────┐    ┌────────────┐      ┌──────────────┐
  │Firehose │    │Managed     │      │Lambda (Fraud │
  │  → S3   │    │Flink       │      │ Detection)   │
  │(offline)│    │→ DynamoDB  │      │              │
  └─────────┘    └────────────┘      └──────────────┘
```

## 决策表：接入新数据源时该考虑什么

| 步骤 | 问题 | 决策 |
|---|---|---|
| 1 | 是 Aurora / RDS MySQL 吗？ | 是 → Zero-ETL to Lakehouse；否 → 下一步 |
| 2 | 是 DMS 支持的源吗？ | 是 → DMS；否 → 下一步 |
| 3 | 是 ES / OpenSearch 吗？ | 是 → OSI（托管）/ Logstash（自建） |
| 4 | 是流式事件（埋点 / 日志）吗？ | 是 → API GW + (MSK) + Firehose |
| 5 | 以上都不是？ | Lambda / Glue Job 自定义 ETL |

## 本章小结

| 管道 | 关键服务 | 一句话总结 |
|---|---|---|
| Aurora MySQL 到 S3 | Zero-ETL to Lakehouse | 亚秒级 + serverless + 直接输出 Iceberg |
| 自建 MySQL 到 S3 | 基于 DMS 的 Zero-ETL | 面向自管数据库的 Zero-ETL 路径 |
| 异构源到 S3 | 传统 DMS | 支持 60 多种数据库；需自行管理复制实例 |
| ES 到 S3 | OpenSearch Ingestion / Logstash | DMS 不支持将 ES 作为源 |
| 事件埋点到 S3 + 实时 | API GW + MSK + Firehose + Flink | 多订阅者扇出，兼顾离线 + 实时双路径 |

下一篇：数据进入 S3 之后，如何让 Athena 和 Spark 把它识别为“表”？

---

## 常见问题

### 什么是 CDC？为什么它比轮询更适合数据接入？

CDC（Change Data Capture，变更数据捕获）订阅数据库的 binlog，实时捕获每一次 INSERT、UPDATE 和 DELETE，对源库影响极小；而轮询会给源库带来额外查询压力，并且会漏掉 DELETE 操作。

### 什么时候该用 Aurora Zero-ETL，什么时候该用 AWS DMS？

对 Aurora MySQL/PostgreSQL 源使用 Aurora Zero-ETL —— 它具备亚秒级延迟且完全托管。对于 Zero-ETL 不支持的异构源（如 DocumentDB、自建 MySQL 或 Oracle），则使用 DMS。


---

## 参考资料

- [Apache Kafka documentation](https://kafka.apache.org/documentation/) — Apache
- [Amazon Kinesis Data Streams](https://docs.aws.amazon.com/streams/latest/dev/introduction.html) — AWS Documentation
- [Debezium](https://debezium.io/) — Debezium
