# AWS 大数据深度解析（第五部分）：EMR、Glue ETL、Flink 与管道编排

> 对比 EMR Serverless、Glue ETL、Managed Flink，选出合适的计算引擎；再用 MWAA（Airflow）与 Step Functions 编排数据管道。

- 作者: zhuermu
- 发布: 2025-05-14
- 网页版: https://zhuermu.com/blog/bigdata-deep-dive-part-5-compute-orchestration/

---
> 当数据落入 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 集群、运行作业，然后销毁集群。

```python
# 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](/images/blog/bigdata-deep-dive/05-orchestration.svg)

谁来管理这个 DAG？有两位竞争者。

### MWAA（Managed Workflows for Apache Airflow）

开源 Apache Airflow（最初由 Airbnb 创建）的托管版本。

**Airflow 的工作方式**：用 Python 描述一个 DAG。

```python
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）：**

```json
{
  "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 部分的各项服务串联起来，为客户场景绘制出**端到端的数据管道**。

---

## 常见问题

### 什么时候该用 EMR Serverless，什么时候该用 Glue ETL？

标准 Spark ETL 作业、需要自动扩缩且不想管理集群时，用 Glue ETL。当你需要对 Spark/Hive/Presto 完全掌控、使用自定义库，或工作负载能受益于 EMR 的各项优化时，用 EMR Serverless。

### 管道编排该用 MWAA（Airflow）还是 Step Functions？

对于包含 10 个以上任务、存在跨团队依赖、需要回填（backfill）的复杂 DAG，用 MWAA。对于更简单的工作流（5-20 个步骤），若希望完全 Serverless、按状态转换计费，则用 Step Functions。


---

## 参考资料

- [Apache Spark documentation](https://spark.apache.org/docs/latest/) — Apache
- [Apache Flink documentation](https://nightlies.apache.org/flink/flink-docs-stable/) — Apache
- [Apache Airflow documentation](https://airflow.apache.org/docs/) — Apache
- [Amazon EMR Management Guide](https://docs.aws.amazon.com/emr/latest/ManagementGuide/emr-what-is-emr.html) — AWS Documentation
