AWS 大数据深度解析(第五部分):EMR、Glue ETL、Flink 与管道编排

对比 EMR Serverless、Glue ETL、Managed Flink,选出合适的计算引擎;再用 MWAA(Airflow)与 Step Functions 编排数据管道。

zhuermu··14 分钟
big-dataawsemrglue-etlflinkmwaaairflowstep-functionsdata-pipeline

当数据落入 S3 并注册到 Glue Catalog 后,真正的活儿才开始:分层 ETL 处理、流式处理与调度。

本章涉及的服务:EMR Serverless / AWS Glue ETL / Managed Flink / Lambda / MWAA / Step Functions


计算引擎全景

在数据湖中要”算东西”有很多选择。这里做一个分类:

类别服务最适合场景
SQL 引擎(轻量)Athena CTAS / INSERT用简单 SQL 就能表达的转换
Serverless 上的 SparkAWS Glue ETL / EMR Serverless中到重量级的批处理
EC2 上的 SparkEMR 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

官方文档:


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 ServerlessGlue ETL
价格略低略高
启动时间秒级(预初始化)/ 数十秒1-2 分钟
Spark 版本更接近上游开源版AWS 分支版,略微落后
易用性高(Glue Studio 可视化)
推荐重量级批处理首选团队已在用 / 简单作业

官方文档: EMR Serverless:https://docs.aws.amazon.com/emr/latest/EMR-Serverless-UserGuide/


什么是 Flink,为什么不用 Spark Streaming

Apache Flink 是一个开源流处理引擎。相比 Spark Streaming,它在真正的流式处理上表现出色:

Spark StreamingFlink
模型微批(秒级)真流式(事件级,毫秒级)
状态管理强(RocksDB 状态后端、Savepoints)
Exactly-once复杂原生支持
时间语义一般优秀(事件时间 + watermark)

对于复杂的实时计算(例如”某用户在最近 5 分钟内连续 3 次登录失败”),Flink 的体验远胜 Spark Streaming

旧名: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

谁来管理这个 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 UIMWAA
单一业务线,不想运维Step Functions
混合:MWAA 作为主调度器 + 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 部分的各项服务串联起来,为客户场景绘制出端到端的数据管道

参考资料

  1. Apache Spark documentation — Apache
  2. Apache Flink documentation — Apache
  3. Apache Airflow documentation — Apache
  4. Amazon EMR Management Guide — AWS Documentation