AWS 大数据深度解析(第五部分):EMR、Glue ETL、Flink 与管道编排
对比 EMR Serverless、Glue ETL、Managed Flink,选出合适的计算引擎;再用 MWAA(Airflow)与 Step Functions 编排数据管道。
当数据落入 S3 并注册到 Glue Catalog 后,真正的活儿才开始:分层 ETL 处理、流式处理与调度。
本章涉及的服务:EMR Serverless / AWS Glue ETL / Managed Flink / Lambda / MWAA / Step Functions。
计算引擎全景
在数据湖中要”算东西”有很多选择。这里做一个分类:
| 类别 | 服务 | 最适合场景 |
|---|---|---|
| SQL 引擎(轻量) | Athena CTAS / INSERT | 用简单 SQL 就能表达的转换 |
| Serverless 上的 Spark | AWS Glue ETL / EMR Serverless | 中到重量级的批处理 |
| EC2 上的 Spark | EMR on EC2 | 高度定制化 / 极致 Spot 省钱 |
| 流式处理 | Amazon Managed Service for Apache Flink | 实时特征 / 实时欺诈检测 |
| 轻量函数 | AWS Lambda | 短任务 / 小数据 / 数据补充增强 |
如何选择:
- 一条 SQL 就能搞定 → Athena CTAS(最便宜)
- 需要 Python UDF / 复杂转换 / 大数据集 → 优先 EMR Serverless(通常比 Glue 便宜、性能更好)
- 团队已熟悉 Glue Studio 可视化设计器 / 中等数据量 → Glue ETL Job
- 需要 Spot 实例 / 自定义 Hadoop 组件 → EMR on EC2
- 实时 → Managed Flink
AWS Glue ETL
Glue 的两重身份
请注意,Glue 是一个伞形产品,包含多个子服务:
- Glue Data Catalog(第四部分已讲)—— 元数据
- Glue ETL —— Spark 作业引擎
- Glue Crawler —— 自动表发现
- Glue Studio —— 可视化拖拽式 ETL
- Glue DataBrew —— 数据清洗 UI
- Glue Schema Registry —— Schema 管理
本节聚焦于 Glue ETL。
Glue ETL = Serverless Spark
其核心是:托管的 Apache Spark + Python(PySpark)。你编写 Spark Job 代码,AWS 负责拉起一个 Spark 集群、运行作业,然后销毁集群。
# A typical Glue Job (PySpark)
from pyspark.sql import SparkSession
from awsglue.context import GlueContext
glueContext = GlueContext(SparkSession.builder.getOrCreate())
# Read from an Iceberg table
df = glueContext.create_data_frame.from_catalog(
database="poc_social_layla",
table_name="ods_event"
).filter("dt = '2026-05-10'")
# Transform
df_clean = df.dropDuplicates(['event_id'])
# Write to DWD Iceberg table
df_clean.writeTo("poc_social_layla.dwd_user_action").append()
计费:DPU x 时间
DPU = Data Processing Unit = 4 vCPU + 16 GB 内存。
计费:
- Glue ETL:$0.44 每 DPU-小时(最低 1 分钟)
- Glue Flex(低优先级,便宜 35%):$0.29 每 DPU-小时
- Streaming Job:$0.44 每 DPU-小时
真实成本示例:每天 30 GB 数据,10 DPU x 10 分钟 = 1.67 DPU-小时 x $0.44 ~ $0.73/天 ~ $22/月。
Glue 5.0(2024-2025 GA)的改进
Glue 5.0 将 Spark 升级到 3.5,原生集成最新版本的 Iceberg / Delta / Hudi,并引入了:
- 自动 Iceberg 压缩 / 快照过期清理(由 S3 Tables / Glue Catalog 托管)—— 无需自己编写 OPTIMIZE 作业
- 启动时间从约 1 分钟缩短到约 30 秒
- 在 Spark 内部强制执行 Lake Formation 的行级 / 列级权限
Glue 尚存的痛点
- 相比开源 Spark 启动仍然较慢(短作业不划算)
- 同等规模下,EMR Serverless 通常仍便宜 15-25%
经验法则依旧是:重作业 → EMR Serverless / 中等作业 → Glue 5.0 / 轻量 SQL → Athena CTAS。
官方文档:
- Glue ETL:https://docs.aws.amazon.com/glue/
- 写入 Iceberg:https://docs.aws.amazon.com/glue/latest/dg/aws-glue-programming-etl-format-iceberg.html
EMR / EMR Serverless
EMR 版本演进
| 版本 | 说明 |
|---|---|
| EMR on EC2 | 经典模式:拉起一个运行 Hadoop / Spark / Hive 的 EC2 集群。最灵活、最便宜(Spot),但需要运维 |
| EMR on EKS | 在 Kubernetes 上运行 Spark |
| EMR Serverless(2022+) | 完全 Serverless,AWS 全托管,按 vCPU + GB-小时计费 |
我们的推荐:EMR Serverless。
EMR Serverless 执行模型
Your Spark Job (PySpark / Scala JAR)
|
v
EMR Serverless Application
(Apache Spark / Hive — your choice)
|
v
AWS automatically spins up workers (elastic on demand)
|
v
Job completes, results land in S3, resources released
关键特性:
- 亚秒级启动(配合预初始化容量模式)
- 按精确资源用量计费(vCPU-小时 + 内存 GB-小时 + 存储)
- 无最低消费
EMR Serverless vs Glue ETL
| EMR Serverless | Glue ETL | |
|---|---|---|
| 价格 | 略低 | 略高 |
| 启动时间 | 秒级(预初始化)/ 数十秒 | 1-2 分钟 |
| Spark 版本 | 更接近上游开源版 | AWS 分支版,略微落后 |
| 易用性 | 中 | 高(Glue Studio 可视化) |
| 推荐 | 重量级批处理首选 | 团队已在用 / 简单作业 |
官方文档: EMR Serverless:https://docs.aws.amazon.com/emr/latest/EMR-Serverless-UserGuide/
Amazon Managed Service for Apache Flink
什么是 Flink,为什么不用 Spark Streaming
Apache Flink 是一个开源流处理引擎。相比 Spark Streaming,它在真正的流式处理上表现出色:
| Spark Streaming | Flink | |
|---|---|---|
| 模型 | 微批(秒级) | 真流式(事件级,毫秒级) |
| 状态管理 | 弱 | 强(RocksDB 状态后端、Savepoints) |
| Exactly-once | 复杂 | 原生支持 |
| 时间语义 | 一般 | 优秀(事件时间 + watermark) |
对于复杂的实时计算(例如”某用户在最近 5 分钟内连续 3 次登录失败”),Flink 的体验远胜 Spark Streaming。
AWS 托管 Flink 服务
旧名:Kinesis Data Analytics for Apache Flink
新名(2023 年 8 月更名):Amazon Managed Service for Apache Flink
关键特性:
- 托管的 Flink 集群
- 按 KPU 计费(Kinesis Processing Unit = 1 vCPU + 4 GB)
- 与 MSK / KDS / Firehose / DynamoDB / S3 集成
在客户场景中的角色:实时特征
MSK Topic: events
|
v
Managed Flink:
- Group by user_id
- Maintain "last 5 clicks" state (RocksDB)
- On each new event → update state → write to DynamoDB
|
v
DynamoDB user_realtime_features:
user_id=12345 → { last_5_clicks: [item_a, item_b, ...] }
|
v
Recommendation service point-queries at inference time (milliseconds)
实用建议:
- 根据吞吐量确定 KPU 数量(每个 KPU 约 10,000 事件/秒)
- 状态量大时启用 RocksDB 后端
- 将 checkpoint 间隔设为 1-5 分钟以支持故障恢复
官方文档: https://docs.aws.amazon.com/managed-flink/
AWS Lambda
什么是 Lambda
Serverless 函数计算。你提供代码,由事件触发执行:
- 单次调用最长时长:15 分钟
- 内存:128 MB 到 10 GB
- 按毫秒计费
在数据管道中的角色
| 位置 | 用途 |
|---|---|
| API Gateway 之后 | 鉴权 + 事件补充增强 |
| MSK / KDS 消费者 | 简单的实时处理 |
| S3 PUT 触发 | 文件一落地立即处理 |
| EventBridge / Cron | 轻量的周期性任务 |
| Glue / EMR 触发器 | 拉起下游作业 |
Lambda 不擅长的场景
- 长任务(> 15 分钟)→ 用 ECS / Step Functions
- 大内存(> 10 GB)→ 用 EMR
- 持久化状态 → 用 DynamoDB / RDS
管道编排:MWAA vs Step Functions
数据仓库的 ETL 从来不是一个作业就能搞定的。它是由数十个作业按依赖关系串联而成的一个 DAG(有向无环图):
谁来管理这个 DAG?有两位竞争者。
MWAA(Managed Workflows for Apache Airflow)
开源 Apache Airflow(最初由 Airbnb 创建)的托管版本。
Airflow 的工作方式:用 Python 描述一个 DAG。
from airflow import DAG
from airflow.providers.amazon.aws.operators.glue import GlueJobOperator
from airflow.providers.amazon.aws.operators.athena import AthenaOperator
from datetime import datetime
with DAG('daily_warehouse', start_date=datetime(2026,5,1), schedule='0 2 * * *') as dag:
dwd_clean = GlueJobOperator(
task_id='dwd_clean',
job_name='dwd_user_action_clean',
)
dws_aggr = AthenaOperator(
task_id='dws_aggr',
query="INSERT INTO dws_user_daily SELECT ... FROM dwd_user_action ...",
workgroup='poc-social-layla',
)
ads_features = GlueJobOperator(
task_id='ads_features',
job_name='ads_user_features',
)
dwd_clean >> dws_aggr >> ads_features
优势:
- 强大的 DAG 表达能力(条件分支、动态生成、SubDAG)
- 200+ 个 Operator(含 Glue / EMR / Athena / SageMaker)
- 直观的 Web UI(查看 DAG 状态、重试、回填)
- 内置重试、SLA 与告警
痛点:
- MWAA 起步门槛高:一个 mw1.small 基础容量约 $300/月;再加 1-2 个弹性 worker 会落在 $300-400/月
- Airflow 有学习曲线(DAG 调度概念、execution_date 时区陷阱)
- 升级 Airflow 版本很痛苦
AWS Step Functions
完全 AWS 原生,按状态转换计费(每百万次状态转换 $25)。
工作流用 JSON 描述(ASL = Amazon States Language):
{
"StartAt": "DWD",
"States": {
"DWD": {
"Type": "Task",
"Resource": "arn:aws:states:::glue:startJobRun.sync",
"Parameters": {"JobName": "dwd_user_action_clean"},
"Next": "DWS"
},
"DWS": {
"Type": "Task",
"Resource": "arn:aws:states:::athena:startQueryExecution.sync",
"Parameters": {"QueryString": "INSERT INTO ..."},
"Next": "ADS"
},
"ADS": {
"Type": "Task",
"Resource": "arn:aws:states:::glue:startJobRun.sync",
"Parameters": {"JobName": "ads_user_features"},
"End": true
}
}
}
优势:
- 完全 Serverless,按用量付费
- 与 200+ AWS 服务集成(可直接调用,无需 Lambda 包装)
- 可视化 DAG(执行过程中实时查看每一步的状态)
- 支持错误重试、并行分支、Map State
痛点:
- DAG 表达能力不如 Airflow 灵活(动态 DAG、复杂条件分支较弱)
- ASL JSON 规模一大就难以维护(建议用 CDK / Terraform 生成)
如何选择
| 场景 | 选择 |
|---|---|
| 数仓批处理 ETL,10+ 任务 | MWAA(成熟的 Airflow 生态) |
| 数仓批处理 ETL,简单的 5-20 个任务 | Step Functions(便宜、Serverless) |
| 跨团队复杂调度,需要 Web UI | MWAA |
| 单一业务线,不想运维 | Step Functions |
| 混合:MWAA 作为主调度器 + Step Functions 处理子工作流 | 两者兼用 |
官方文档:
- MWAA:https://docs.aws.amazon.com/mwaa/
- Step Functions:https://docs.aws.amazon.com/step-functions/
客户场景:编排示例
Daily 02:00 (UTC+8): MWAA DAG kicks off
+-- 02:00 ods_user_full_load_check (prerequisite: DMS / Zero-ETL completed for the day)
+-- 02:30 dwd_user_action_clean (Glue Job)
+-- 02:30 dwd_post_enrich (Glue Job)
+-- 03:00 dws_user_daily (Athena CTAS)
+-- 03:30 ads_user_features (EMR Serverless)
+-- 03:30 ads_sample_follow (EMR Serverless)
+-- 04:00 sync_to_dynamodb (Glue Job writes to DynamoDB)
+-- 04:30 train_recall_model (SageMaker Training Job)
+-- 04:30 train_rank_model (SageMaker Training Job)
+-- 05:30 deploy_endpoint (Lambda calls SageMaker API)
+-- 06:00 dq_check_report (Slack alert / email)
Each step:
- Auto-retry 2 times on failure
- Still failing → page on-call engineer
- Overall DAG SLA: 7 hours
本章小结
| 服务 | 一句话总结 |
|---|---|
| Athena CTAS | 简单 SQL 转换的最便宜之选 |
| Glue ETL | 带 Studio 可视化编辑器的托管 Spark |
| EMR Serverless | 重量级批处理 —— 更便宜也更快 |
| Managed Flink | 实时流处理 |
| Lambda | 短任务 / 数据增强 / 触发器 |
| MWAA | 借助 Airflow 生态进行复杂 DAG 调度 |
| Step Functions | 简单 DAG,Serverless,AWS 原生 |
至此,计算与编排层就介绍完了。在下一章,我们会把第 3-5 部分的各项服务串联起来,为客户场景绘制出端到端的数据管道。
参考资料
- Apache Spark documentation — Apache
- Apache Flink documentation — Apache
- Apache Airflow documentation — Apache
- Amazon EMR Management Guide — AWS Documentation