AWS 大数据深度剖析(第 6 部分):端到端数据管道 —— 从数据源到特征存储

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

zhuermu··10 分钟
big-dataawsdata-pipelineend-to-endevent-drivendata-warehouse-layers

第 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 ODS14+ 张 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_userAurora usersuser_iddt
ods_postAurora postspost_iddt, hr
ods_followAurora follows(follower_id, followee_id)dt
ods_doc_userDocumentDB user_profileuser_iddt
ods_doc_msgDocumentDB messagesmsg_iddt, hr
ods_event埋点(MSK→Firehose)event_iddt, 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 而非 DMSAurora 原生路径,亚秒级延迟,零运维
DocumentDB 使用 DMSDocumentDB 没有 Zero-ETL
埋点:API GW + MSK + FirehoseMSK 多订阅同时支持离线和实时
ETL:以 EMR Serverless 为主 + Athena CTAS 处理轻量作业最优成本组合
编排:MWAA任务依赖复杂;Airflow 表达能力强
数据格式:Parquet + Iceberg列式压缩 + ACID + Time Travel + PIT
Schema 管理:Glue Schema Registry50+ 种事件类型时必需

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

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

本章小结

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

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

接下来的四章将进入 ML 侧(推荐系统):

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

参考资料

  1. Apache Airflow documentation — Apache
  2. AWS Glue Developer Guide — AWS Documentation
  3. Apache Iceberg — Apache