

# Spark 声明式管道
<a name="spark-declarative-pipelines"></a>

Spark 声明式管道（SDP）是一个声明式框架，用于在 Amazon Glue 6.0 中构建批处理和流式传输数据管道。使用 SDP，您可以通过 SQL 或 Python 定义数据的目标形态，而框架将自动确定执行计划、解析数据集之间的依赖关系，并行运行相互独立的分支。

SDP 省去了用于读取、写入、目录注册和执行排序的命令式样板代码，从而简化了管道开发过程。您可以专注于业务转换，而管道基础设施则由框架负责处理。

SDP 在 Amazon Glue 6.0 及更高版本中可用。

## SDP 概念
<a name="spark-declarative-pipelines-concepts"></a>

管道由一个 YAML 清单文件（`spark-pipeline.yml`）和一个或多个 SQL 或 Python 转换文件组成。SDP 自动执行以下操作：
+ 通过从表引用推断 DAG 来解析依赖关系
+ 确定执行顺序，而无需进行手动编排
+ 并行运行独立分支，以实现最大吞吐量
+ 通过检查点管理流式传输表的增量状态
+ 实体化时将输出表注册到目录中

### 数据集类型
<a name="spark-declarative-pipelines-dataset-types"></a>

SDP 中提供三种数据集类型：

流式传输表  
仅处理自上次运行以来的新数据。使用检查点跨作业执行维护状态。将流式传输表用于数据摄取、事件流、IoT 数据、变更数据捕获和仅追加数据来源。

实体化视图  
每次运行时都会完全重新计算数据集。输出始终反映源数据的当前状态。使用实体化视图进行聚合、联接、摘要分析和报告操作。

临时视图  
会话范围限定，且不进行持久化存储或编目登记。将临时视图用于中间转换和暂存逻辑。

**重要**  
在当前版本中，实体化视图始终执行完全重新计算。它们不支持增量刷新。使用流式传输表处理增量工作负载。

## 先决条件
<a name="spark-declarative-pipelines-prerequisites"></a>

要使用 SDP，您需要以下内容：
+ Amazon Glue 版本 6.0
+ 用于存储管道的 Amazon S3 位置（检查点、元数据）
+ 对于数据目录集成（可选）：将 `--enable-glue-datacatalog` 设置为 `true`。或者，您也可以通过 Spark 配置直接配置目录设置。
+ 对于永久表存储：将 `spark.sql.warehouse.dir` 设置为 Amazon S3 路径，或在管道 YAML 中设置 `database:` 字段，并确保 Amazon Glue 数据库已将 `LocationUri` 配置为 Amazon S3 路径
+ 要使用流式传输表进行跨运行增量处理，流式传输表的数据和检查点状态必须在 Amazon S3 上持久保存。例如，Hive 或 Amazon Glue 托管（非 Iceberg）表要求将数据库 `LocationUri` 设置为 Amazon S3 路径，而 Apache Iceberg 表则自行管理其表元数据。

**重要**  
如果您在管道 YAML 中为 Hive 或 Amazon Glue 托管（非 Iceberg）表使用 `database:` 字段，则相应的 Amazon Glue 数据库必须将其 `LocationUri` 设置为 Amazon S3 路径。`LocationUri` 将托管流式传输表（及其 `_spark_metadata` 日志）放置在 Amazon S3 上，从而为这些表提供持久的跨运行增量处理能力。使用 Data Catalog 且采用 `catalog-impl=GlueCatalog`（选项 1）的 Iceberg 表不需要数据库 `LocationUri`。使用明确的 Amazon S3 位置创建或更新数据库：  

```
aws glue create-database --database-input '{
  "Name":"my_pipeline_db",
  "LocationUri":"s3://my-bucket/warehouse/my_pipeline_db"
}'
```

## 创建管道
<a name="spark-declarative-pipelines-creating"></a>

要创建 SDP 管道，请完成以下步骤。

### 步骤 1：创建管道 YAML
<a name="spark-declarative-pipelines-step1-yaml"></a>

创建一个名为 `spark-pipeline.yml` 的文件：

```
name: my_analytics_pipeline
catalog: spark_catalog
database: analytics_db
storage: s3://my-bucket/pipeline-storage/
libraries:
  - glob:
      include: transformations/**
configuration:
  spark.sql.shuffle.partitions: "4"
```

下表描述了管道 YAML 字段。


| 字段 | 必填 | 描述 | 
| --- | --- | --- | 
| name | 是 | 管道的名称。 | 
| catalog | 否 | 要使用的目录。默认值为 spark\_catalog。 | 
| database | 否 | 用于存放输出表的 Amazon Glue 数据库。有关 LocationUri 要求，请参阅[先决条件](#spark-declarative-pipelines-prerequisites)。 | 
| storage | 是 | 用于存储管道检查点和元数据的 Amazon S3 路径。 | 
| libraries | 是 | 要包含的转换文件的 Glob 模式。 | 
| configuration | 否 | Spark 配置属性。 | 

### 步骤 2：编写转换
<a name="spark-declarative-pipelines-step2-transformations"></a>

在 `transformations/` 目录中创建转换文件。您可以在同一个管道中使用 SQL、Python 或此两者。

**SQL 示例**（`transformations/silver.sql`）：

```
CREATE MATERIALIZED VIEW silver_sales AS
SELECT *, UPPER(region) as clean_region
FROM bronze_sales
WHERE amount > 0;

CREATE MATERIALIZED VIEW gold_summary AS
SELECT clean_region, COUNT(*) as order_count, SUM(amount) as total_revenue
FROM silver_sales
GROUP BY clean_region;
```

**Python 示例**（`transformations/bronze.py`）：

```
from pyspark import pipelines as dp
from pyspark.sql import DataFrame, SparkSession

spark = SparkSession.active()

@dp.materialized_view(comment="Raw sales data from S3")
def bronze_sales() -> DataFrame:
    return spark.read.format("csv").option("header", "true") \
        .option("inferSchema", "true") \
        .load("s3://source-bucket/raw-data/sales/")
```

**Python 流式传输表示例**（`transformations/events.py`）：

```
from pyspark import pipelines as dp
from pyspark.sql import DataFrame, SparkSession

spark = SparkSession.active()

dp.create_streaming_table(
    "streaming_events",
    comment="Incremental event ingestion",
    schema="event_id STRING, event_type STRING, timestamp LONG, payload STRING"
)

@dp.append_flow(target="streaming_events")
def ingest_events() -> DataFrame:
    return (
        spark.readStream.format("json")
        .schema("event_id STRING, event_type STRING, timestamp LONG, payload STRING")
        .load("s3://source-bucket/events/")
    )
```

**注意**  
对于 Python 中的流式传输表，请将 `dp.create_streaming_table()` 与 `@dp.append_flow(target=...)` 结合使用。

### 步骤 3：上传到 Amazon S3
<a name="spark-declarative-pipelines-step3-upload"></a>

将您的管道文件上传到 Amazon S3 作为以下任一内容：
+ 一个包含 `spark-pipeline.yml` 和 `transformations/` 目录的 `.zip` 文件
+ 包含相同结构的 Amazon S3 前缀（目录）

### 步骤 4：创建并运行 Amazon Glue 作业
<a name="spark-declarative-pipelines-step4-create-job"></a>

使用以下参数创建 Amazon Glue 作业：
+ `--enable-spark-declarative-pipeline`：`true`（必填；启用 SDP 模式）
+ `ScriptLocation`：管道定义 zip 或 Amazon S3 前缀（SDP 管道的必填项）
+ `--enable-glue-datacatalog`：`true`（可选；在 Data Catalog 中注册数据表）

以下示例使用 Amazon CLI 创建 SDP 作业：

```
aws glue create-job \
  --name my-sdp-pipeline \
  --role arn:aws:iam::123456789012:role/MyGlueRole \
  --glue-version 6.0 \
  --worker-type G.1X --number-of-workers 2 \
  --command '{"Name":"glueetl","ScriptLocation":"s3://my-bucket/pipelines/my_pipeline.zip"}' \
  --default-arguments '{
      "--enable-spark-declarative-pipeline": "true",
      "--enable-glue-datacatalog": "true"
  }'
```

## 运行管道
<a name="spark-declarative-pipelines-running"></a>

您可以使用 `StartJobRun` 运行 SDP 管道。您可以使用在运行时传递的作业参数来控制执行行为。

### 运行模式
<a name="spark-declarative-pipelines-run-modes"></a>

将以下参数传递给 `StartJobRun` 以控制管道执行：

`--conf spark.glue.sdp.jobMode`  
控制执行模式：  
+ `RUN`（默认）：正常执行管道。
+ `VALIDATE`：执行试运行，在不写入任何数据的情况下检查 YAML 语法、依赖关系解析和 SQL/Python 编译。

`--conf spark.glue.sdp.runMode`  
控制刷新哪些数据集：    
默认（无运行模式标志）  
运行所有数据集。实体化视图完全重新计算；流式传输表仅处理自上次检查点以来的新数据。  
`--refresh <dataset>`  
仅刷新指定的数据集。流式传输表以增量方式处理新数据；实体化视图则进行完整重新计算。  
`--full-refresh <dataset>`  
仅重置和重新计算指定的数据集。对于流式传输表，这会重置检查点并重新处理所有数据。  
`--full-refresh-all`  
重置并重新计算所有数据集。

## 将 Iceberg 表与 SDP 结合使用
<a name="spark-declarative-pipelines-iceberg"></a>

对于需要持久、跨运行增量处理的流式传输表，Apache Iceberg 是推荐的表格式，因为它不依赖于 Hive 或 Amazon Glue 托管表所使用的基于文件的 `_spark_metadata` 日志。您可以通过两种方式使用 SDP 配置 Iceberg，具体取决于是否要在数据目录中注册输出表。

### 选项 1：使用数据目录的 Iceberg（推荐）
<a name="spark-declarative-pipelines-iceberg-glue-catalog"></a>

如果您希望将 Iceberg 表注册到数据目录中（使用 `table_type=ICEBERG`），以便通过 Amazon Redshift 和 Amazon EMR 等其他引擎查询这些表，请使用此选项。表数据和元数据存储在 Amazon S3 中，并保留跨运行增量状态。将以下内容添加到您 `spark-pipeline.yml` 的 `configuration` 部分：

```
configuration:
  spark.sql.catalog.glue_catalog: "org.apache.iceberg.spark.SparkCatalog"
  spark.sql.catalog.glue_catalog.catalog-impl: "org.apache.iceberg.aws.glue.GlueCatalog"
  spark.sql.catalog.glue_catalog.io-impl: "org.apache.iceberg.aws.s3.S3FileIO"
  spark.sql.catalog.glue_catalog.warehouse: "s3://my-bucket/iceberg-warehouse"
```

在您的管道 YAML 中，设置 `catalog: glue_catalog` 并将 `database:` 设置为 Amazon Glue 数据库。创建 Amazon Glue 作业时，将 `--enable-spark-declarative-pipeline` 设置为 `true`。请勿为 Iceberg 表设置 `--enable-glue-datacatalog`。

**注意**  
Iceberg 目录使用自己的 `warehouse` 位置存储表数据和元数据。使用此选项时，无需为 Iceberg 表本身设置 `spark.sql.warehouse.dir` 或数据库 `LocationUri`。

### 选项 2：使用基于文件的（Hadoop）目录的 Iceberg
<a name="spark-declarative-pipelines-iceberg-hadoop-catalog"></a>

当您不需要数据目录注册时，请使用此选项。Iceberg 元数据在 Amazon S3 中是基于文件的，并且这些表**未**在数据目录中注册。将以下内容添加到您 `spark-pipeline.yml` 的 `configuration` 部分：

```
configuration:
  spark.sql.catalog.spark_catalog: "org.apache.iceberg.spark.SparkSessionCatalog"
  spark.sql.catalog.spark_catalog.type: "hadoop"
  spark.sql.catalog.spark_catalog.warehouse: "s3://my-bucket/iceberg-warehouse"
```

**注意**  
请注意有关此选项的以下信息：  
使用此选项创建的表不会注册到数据目录中，因此无法通过查询引擎进行查询。
此选项使用自己的 `warehouse` 位置来存储表。

如果您需要在数据目录中注册表并可从其他引擎查询，请使用选项 1。仅当不需要注册数据目录时，才使用选项 2。

**注意**  
`--enable-glue-datacatalog` 作业参数将 Spark Hive 元数据仓连接到 Hive（非 Iceberg）表的数据目录。对于 Iceberg 表，`catalog-impl=GlueCatalog` 通过 Amazon SDK 直接在数据目录中注册这些表，因此不要为 Iceberg 设置 `--enable-glue-datacatalog`。在 Amazon Glue 6.0 上，请勿同时使用 Iceberg 的默认 `SparkSessionCatalog`（`type: hive`）配置和 `--enable-glue-datacatalog`，试图将 Iceberg 表注册到数据目录；这种组合会失败。

配置 Iceberg 后，您可以获得以下优势：
+ 流式传输表在作业运行之间将检查点状态保存在 Amazon S3 中
+ 每次运行都会创建新的数据文件和 Iceberg 快照
+ 后续运行从上次提交的偏移量处恢复
+ 完整的表历史记录通过 Iceberg 的快照机制保留

您也可以将 Iceberg 表作为流式数据源进行增量读取。在奖章架构中，下游流式传输表在每次运行时只能处理新提交到上游 Iceberg 表的行。以下示例以增量方式将数据从 Iceberg `bronze` 表读取到 `silver` 流式传输表中：

```
from pyspark import pipelines as dp
from pyspark.sql import SparkSession

spark = SparkSession.active()

dp.create_streaming_table(
    "silver",
    comment="Incremental silver layer built from the bronze Iceberg table"
)

@dp.append_flow(target="silver")
def from_bronze():
    # Incremental read from the Iceberg bronze table; each run processes only new rows.
    return spark.readStream.table("bronze")
```

## 注意事项和限制
<a name="spark-declarative-pipelines-considerations"></a>

在使用 SDP 时，请考虑以下几点：
+ **实体化视图始终完全重新计算**。不支持增量刷新。使用流式传输表处理增量工作负载。
+ **流式传输表 Python API**。将 `dp.create_streaming_table()` 与 `@dp.append_flow(target=...)` 结合使用。
+ **对流式传输表进行跨运行增量处理**。仅当流式传输表的数据和检查点状态持久保存在 Amazon S3 上时，才支持跨运行增量处理。Hive 或 Amazon Glue 托管（非 Iceberg）表要求将数据库 `LocationUri` 设置为 Amazon S3 路径，而 Apache Iceberg 表则自行管理其表元数据。
+ **Hive 或 Amazon Glue 托管（非 Iceberg）表需要设置数据库 LocationUri**。使用数据目录且采用 `catalog-impl=GlueCatalog` 的 Iceberg 表不需要设置该项。有关更多信息，请参阅 [先决条件](#spark-declarative-pipelines-prerequisites)。
+ **数据质量预期**。当前 SDP 框架不支持内联数据质量标注。
+ **避免在下游查询函数中使用 `withColumn`**。当下游数据集（例如实体化视图）使用 `spark.table(...)` 从上游管道数据集读取数据并应用 `.withColumn(...)` 时，SDP 可能无法在第二次及后续运行中检测到数据集之间的依赖关系。这会导致下游读取上次运行中的陈旧数据（单次运行延迟）。为避免此问题，请在 `.select(...)` 内表达派生列，而不是使用 `.withColumn(...)`。还要避免在查询函数中强制执行计划解析的任何操作（例如 `.schema` 或 `.collect`）。
+ **无迁移工具**。不支持从其他管道框架自动迁移。以增量方式迁移表：SDP 可从现有目录表中读取数据。
+ **调度**。SDP 作业与其他 Amazon Glue 作业使用相同的调度机制（Amazon Glue 触发器、Amazon EventBridge、Apache Airflow）。

## 从命令式脚本迁移到 SDP
<a name="spark-declarative-pipelines-migrating"></a>

您可以将现有的命令式 Spark 脚本以增量方式迁移到 SDP：

1. 从一张表开始：将单个 `spark.sql(...).write.saveAsTable(...)` 调用转换为 `CREATE MATERIALIZED VIEW` SQL 语句。

1. 以增量方式添加表。SDP 处理混合依赖关系。SDP 表可以从不属于管道的现有目录表中读取数据。

1. 在过渡期间并行运行两种模式。SDP 作业与命令式作业可以共存。

SDP 可以引用任何可通过 SparkSession 访问的表，包括现有的 Data Catalog 表、外部表以及跨数据库引用。