

本文為英文版的機器翻譯版本，如內容有任何歧義或不一致之處，概以英文版為準。

# Spark 宣告管道
<a name="spark-declarative-pipelines"></a>

Spark 宣告管道 (SDP) 是用於在 AWS Glue 6.0 中建置批次和串流資料管道的宣告性架構。使用 SDP，您可以定義資料使用 SQL 或 Python 的外觀，而架構會自動決定執行計畫、解決資料集之間的相依性，並平行執行獨立分支。

SDP 透過消除讀取、寫入、目錄註冊和執行順序的強制性樣板程式碼，簡化管道開發。當架構處理管道基礎設施時，您會專注於業務轉型。

SDP 可在 6.0 版和更新 AWS Glue 版本中使用。

## 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，您需要下列項目：
+ AWS Glue 6.0 版
+ 管道儲存的 Amazon S3 位置 （檢查點、中繼資料）
+ 針對 Data Catalog 整合 （選用）：`--enable-glue-datacatalog`設為 `true`。或者，您可以透過 Spark 組態直接設定目錄設定。
+ 對於持久性資料表儲存：`spark.sql.warehouse.dir`設定為 Amazon S3 路徑，或設定管道 YAML 中的 `database:` 欄位， AWS Glue 並確保資料庫已將 `LocationUri`設定為 Amazon S3 路徑
+ 對於使用串流資料表進行跨執行增量處理：使用 Iceberg 資料表 (Hive 受管串流資料表不支援跨執行增量處理）

**重要**  
如果您在管道 YAML 中使用 `database:` 欄位，對應的 AWS Glue 資料庫必須將其`LocationUri`設定為 Amazon S3 路徑。透過 主控台建立的資料庫通常具有空白的 `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 | 否 | 輸出資料表的目標資料庫。資料庫必須存在且具有LocationUri集合。 | 
| 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=...)`。`@dp.streaming_table` 裝飾項目無法使用。

### 步驟 3：上傳至 Amazon S3
<a name="spark-declarative-pipelines-step3-upload"></a>

將您的管道檔案上傳至 Amazon S3，格式為：
+ 包含 `spark-pipeline.yml`和 `transformations/`目錄`.zip`的檔案
+ 包含相同結構的 Amazon S3 字首 （目錄）

### 步驟 4：建立和執行 AWS Glue 任務
<a name="spark-declarative-pipelines-step4-create-job"></a>

使用下列參數建立 AWS Glue 任務：
+ `--enable-spark-declarative-pipeline`：`true`（必要 — 啟用 SDP 模式）
+ `ScriptLocation`：管道定義 zip 或 Amazon S3 字首 (SDP 管道需要）
+ `--enable-glue-datacatalog`：`true`（選用 — 在 Data Catalog 中註冊資料表）

下列範例使用 CLI 建立 SDP AWS 任務：

```
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>

您可以使用 執行 SDP 管道`StartJobRun`。您可以使用執行時間傳遞的任務引數來控制執行行為。

### 執行模式
<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`  
執行所有資料集。具體化視觀表完全重新計算。串流資料表只會處理自上次檢查點以來的新資料。  
`--refresh <dataset_name>`  
僅執行指定的資料集。對於串流資料表，這會遞增處理新資料。對於具體化視觀表，這只會執行該視觀表的完整重新計算。  
`--full-refresh`  
重設並重新計算所有資料集。對於串流資料表，這會重設檢查點，並從頭開始重新處理所有資料。  
`--full-refresh-all`  
捨棄所有資料表，並從頭開始重新處理整個管道。

## 搭配 SDP 使用 Iceberg 資料表
<a name="spark-declarative-pipelines-iceberg"></a>

對於需要跨執行增量處理的串流資料表，請使用 Apache Iceberg。Hive 受管串流資料表會在本機存放中繼資料，而不會在任務執行期間保留。

若要設定 Iceberg，請將下列項目新增至您的`spark-pipeline.yml`組態區段：

```
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"
```

設定 Iceberg 後，您可以獲得下列優點：
+ 串流資料表會在任務執行期間維持 Amazon S3 中的檢查點狀態
+ 每次執行都會建立新的資料檔案和 Iceberg 快照
+ 後續執行會從上次遞交的偏移恢復
+ 透過 Iceberg 的快照機制保留完整資料表歷史記錄

## 考量和限制
<a name="spark-declarative-pipelines-considerations"></a>

當您使用 SDP 時，請考慮下列事項：
+ **具體化視觀表一律完全重新計算** - 不支援增量重新整理。針對增量工作負載使用串流資料表。
+ **串流資料表 Python API** — `dp.create_streaming_table()`搭配 使用`@dp.append_flow(target=...)`。目前版本中無法使用`@dp.streaming_table`裝飾項目。
+ **跨執行增量處理需要 Iceberg** — 具有 Hive 或 AWS Glue 受管目錄的串流資料表不支援跨任務執行的增量處理。針對持久性增量狀態使用 Iceberg 資料表。
+ **資料庫 LocationUri 必要** — 如果您在管道 YAML `database:`中指定 ， AWS Glue 資料庫必須將其`LocationUri`設定為 Amazon S3 路徑。如果沒有它，管道會失敗。
+ **資料品質預期** — 目前的 SDP 架構不支援內嵌資料品質註釋。
+ **避免`withColumn`使用下游查詢函數** — 當下游資料集 （例如具體化視觀表） 使用 從上游管道資料集讀取`spark.table(...)`並套用 時`.withColumn(...)`，SDP 可能無法在第二個和後續執行時偵測資料集之間的相依性。這會導致下游從上一次執行讀取過時的資料 （一次性執行延遲）。為了避免此問題，請在內部表達衍生的資料欄`.select(...)`，而不是使用 `.withColumn(...)`。同時避免在查詢函數內強制計畫解析的任何操作 （例如 `.schema`或 `.collect`)。
+ **無遷移工具** — 不支援從其他管道架構自動遷移。遞增遷移資料表 — SDP 可以從現有的目錄資料表讀取。
+ **排程** — SDP 任務使用與其他 AWS Glue 任務相同的排程機制 (AWS 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 存取的任何資料表，包括現有的資料目錄資料表、外部資料表和跨資料庫參考。