AWS 大数据深度剖析(第 6 部分):端到端数据管道 —— 从数据源到特征存储
把所有环节串联起来:追踪一个点击事件如何从客户端 SDK 经过 API Gateway、MSK、Firehose、S3、数仓分层(ODS→DWD→DWS→ADS),最终写入 DynamoDB 用于实时服务。
第 03 至 05 章分别详细介绍了各个服务。本章将把它们全部串联起来,并把整套架构映射到一个真实的社交 App 场景。
本章不引入任何新服务。目标是:读完本章后,你能用一张图讲清楚整个数据侧的架构。
客户场景回顾
来自客户方案文档的关键事实:
- 业务:面向用户推荐的社交 App(关注 / 信息流 / 你可能认识的人)
- 数据源:MySQL(业务数据库)、ES(搜索)、DocumentDB(文档)、客户端埋点
- DAU:100 万+
- DocumentDB 版本:5.0(支持 Change Streams)
- 延迟要求:T+1(离线优先;实时留待后续阶段考虑)
- 现有 Kafka:无(需从零搭建)
端到端架构图
整条管道分为 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+ 种事件类型,这就是必需的)
下面的时间线图展示了双路径消费:
同一个点击事件: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 —— 更规范
- 失败必须重试;建议采用小批量 + 幂等写入
分层处理(每日例行任务)
由 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 章:完整的端到端架构 + 成本估算
参考资料
- Apache Airflow documentation — Apache
- AWS Glue Developer Guide — AWS Documentation
- Apache Iceberg — Apache