本文為英文版的機器翻譯版本,如內容有任何歧義或不一致之處,概以英文版為準。
Spark 宣告管道
Spark 宣告管道 (SDP) 是用於在 AWS Glue 6.0 中建置批次和串流資料管道的宣告性架構。使用 SDP,您可以定義資料使用 SQL 或 Python 的外觀,而架構會自動決定執行計畫、解決資料集之間的相依性,並平行執行獨立分支。
SDP 透過消除讀取、寫入、目錄註冊和執行順序的強制性樣板程式碼,簡化管道開發。當架構處理管道基礎設施時,您會專注於業務轉型。
SDP 可在 6.0 版和更新 AWS Glue 版本中使用。
SDP 概念
管道包含 YAML 資訊清單檔案 (spark-pipeline.yml) 和一或多個 SQL 或 Python 轉換檔案。SDP 會自動:
-
從資料表參考推斷 DAG 以解決相依性
-
在沒有手動協同運作的情況下決定執行順序
-
平行執行獨立分支以達到最大輸送量
-
透過檢查點管理串流資料表的增量狀態
-
在具體化時註冊目錄中的輸出資料表
資料集類型
SDP 提供三種資料集類型:
- 串流資料表
-
僅處理自上次執行以來的新資料。使用檢查點維護任務執行的狀態。將串流資料表用於擷取、事件串流、IoT 資料、變更資料擷取和僅附加來源。
- 具體化視觀表
-
在每次執行時完全重新計算資料集。輸出一律會反映來源資料的目前狀態。使用具體化檢視進行彙總、聯結、摘要分析和報告。
- 暫時檢視
-
工作階段範圍,且未保留或編目。使用暫時檢視進行中繼轉換和預備邏輯。
重要
具體化視觀表一律會在目前版本中執行完整重新計算。它們不支援增量重新整理。針對增量工作負載使用串流資料表。
先決條件
若要使用 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" }'
建立管道
若要建立 SDP 管道,請完成下列步驟。
步驟 1:建立管道 YAML
建立名為 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:寫入轉換
在 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
將您的管道檔案上傳至 Amazon S3,格式為:
-
包含
spark-pipeline.yml和transformations/目錄.zip的檔案 -
包含相同結構的 Amazon S3 字首 (目錄)
步驟 4:建立和執行 AWS Glue 任務
使用下列參數建立 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" }'
執行中管道
您可以使用 執行 SDP 管道StartJobRun。您可以使用執行時間傳遞的任務引數來控制執行行為。
執行模式
將下列引數傳遞至 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 資料表
對於需要跨執行增量處理的串流資料表,請使用 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 的快照機制保留完整資料表歷史記錄
考量和限制
當您使用 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
您可以逐步將現有的必要 Spark 指令碼遷移至 SDP:
-
從一個資料表開始 — 將單一
spark.sql(...).write.saveAsTable(...)呼叫轉換為CREATE MATERIALIZED VIEWSQL 陳述式。 -
遞增新增資料表 — SDP 處理混合相依性。SDP 資料表可以從不屬於管道的現有目錄資料表讀取。
-
在轉換期間平行執行這兩個模式 — SDP 任務和強制性任務可以共存。
SDP 可以參考可透過 SparkSession 存取的任何資料表,包括現有的資料目錄資料表、外部資料表和跨資料庫參考。