AWS 大数据深度剖析(第三部分):数据接入 —— DMS、Zero-ETL、Firehose 与 MSK
四类数据源,四条接入管道 —— 学习使用 DMS 做 CDC、Aurora Zero-ETL、MSK 上的 Kafka,以及 Firehose 微批处理,将数据落地到你的 S3 数据湖。
客户场景中的核心问题:MySQL、DocumentDB、Elasticsearch 以及客户端事件埋点 —— 每一类数据源如何落地到 S3?
本章覆盖四条接入管道以及所涉及的 7 项 AWS 服务(DMS、Zero-ETL、OpenSearch Ingestion、Firehose、MSK、KDS 和 Lambda)。
全景概览
四类数据源对应四条接入管道:
| # | 数据源 | 推荐管道 | 理由 |
|---|---|---|---|
| 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 表:
-- 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,你有以下选择:
- Logstash + S3 输出插件(最常见)
- 自定义 Lambda / EMR Job,使用 scroll API / PIT 拉取数据
- 数据量小的话,定期做
_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
什么是 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 把它识别为“表”?
参考资料
- Apache Kafka documentation — Apache
- Amazon Kinesis Data Streams — AWS Documentation
- Debezium — Debezium