AWS 大数据深度剖析(第三部分):数据接入 —— DMS、Zero-ETL、Firehose 与 MSK

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

zhuermu··16 分钟
big-dataawsdmszero-etlcdckafkamskfirehosedata-ingestion

客户场景中的核心问题:MySQL、DocumentDB、Elasticsearch 以及客户端事件埋点 —— 每一类数据源如何落地到 S3?

本章覆盖四条接入管道以及所涉及的 7 项 AWS 服务(DMS、Zero-ETL、OpenSearch Ingestion、Firehose、MSK、KDS 和 Lambda)。

全景概览

数据接入概览

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

#数据源推荐管道理由
1Aurora / RDS MySQLZero-ETL to SageMaker Lakehouse亚秒级延迟、完全托管、AWS 推荐方案
2自建 MySQL / DocumentDBDMS 落地到 S3,再由 Glue 合并进 IcebergDMS 支持 60 多种异构源
3Elasticsearch / OpenSearchOpenSearch Ingestion 落地到 S3DMS 不支持将 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 / PostgreSQLRedshift2023-06 GABI 分析
Aurora MySQLSageMaker Lakehouse(S3 Tables)2025-06 GA数据湖场景 —— 我们的推荐选择
RDS MySQLRedshift / Lakehouse自 2025 年起陆续 GA(主要商用区域已可用;GovCloud/中国区/部分较新区域仍处于预览或不支持状态 —— 签约前请确认目标区域)同上
DynamoDBRedshift / OpenSearch2023+ GA全文检索
DynamoDB to SageMaker LakehouseLakehouse2025 GA(新)KV 数据入湖
SaaS(Salesforce/SAP/ServiceNow/Zendesk)to LakehouseLakehouse2024-2025 GA跨系统数据集成
自管 MySQLS3 / 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

什么是 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 对比

MSKKDS
协议Apache KafkaAWS 自研
生态完整的 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 到 S3Zero-ETL to Lakehouse亚秒级 + serverless + 直接输出 Iceberg
自建 MySQL 到 S3基于 DMS 的 Zero-ETL面向自管数据库的 Zero-ETL 路径
异构源到 S3传统 DMS支持 60 多种数据库;需自行管理复制实例
ES 到 S3OpenSearch Ingestion / LogstashDMS 不支持将 ES 作为源
事件埋点到 S3 + 实时API GW + MSK + Firehose + Flink多订阅者扇出,兼顾离线 + 实时双路径

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

参考资料

  1. Apache Kafka documentation — Apache
  2. Amazon Kinesis Data Streams — AWS Documentation
  3. Debezium — Debezium