Spark 宣言型パイプライン
Spark 宣言型パイプライン (SDP) は、AWS Glue 6.0 でバッチおよびストリーミングデータパイプラインを構築するための宣言型フレームワークです。SDP では、SQL または Python を使用してデータの外観を定義すると、フレームワークが自動的に実行プランを決定し、データセット間の依存関係を解決し、独立したブランチを並行実行します。
SDP は、読み取り、書き込み、カタログ登録、実行順序付けのための命令型のボイラープレートコードを排除することで、パイプライン開発を簡素化します。フレームワークがパイプラインインフラストラクチャに対処するため、ユーザーはビジネストランスフォーメーションに集中することができます。
SDP は、AWS Glue バージョン 6.0 以降で有効です。
SDP の概念
パイプラインは、YAML マニフェストファイル (spark-pipeline.yml) と、1 つ以上の SQL または Python 変換ファイルで構成されます。SDP は自動的に以下を行います。
-
テーブル参照から DAG を推測して依存関係を解決する
-
手動オーケストレーションなしで実行順序を決定する
-
独立したブランチを並行して実行してスループットを最大化する
-
チェックポイントを介したストリーミングテーブルの増分状態を管理する
-
マテリアライズ時にカタログに出力テーブルを登録する
データセットタイプ
SDP では、以下の 3 つのデータセットタイプを使用できます。
- ストリーミングテーブル
-
最後の実行以降の新しいデータのみを処理します。チェックポイントを使用してジョブ実行全体の状態を維持します。ストリーミングテーブルは、取り込み、イベントストリーム、IoT データ、変更データキャプチャ、追加のみのソースに使用します。
- マテリアライズドビュー
-
実行のたびにデータセットを完全に再計算します。出力は、常にソースデータの現在の状態を反映します。集計、結合、概要分析、レポートにマテリアライズドビューを使用します。
- 一時ビュー
-
セッションスコープであり、永続化またはカタログ化されません。中間変換とステージングロジックには一時ビューを使用します。
重要
現在のバージョンでは、マテリアライズドビューが常に完全に再計算されます。増分更新はサポートされていません。増分ワークロードには、ストリーミングテーブルを使用します。
前提条件
SDP を使用するには、以下が必要です。
-
AWS Glue バージョン 6.0
-
パイプラインストレージ用の Amazon S3 の場所 (チェックポイント、メタデータ)
-
データカタログ統合の場合 (オプション):
--enable-glue-datacatalogをtrueに設定します。または、Spark 設定でカタログ設定を直接設定することもできます。 -
永続テーブルストレージの場合:
spark.sql.warehouse.dirを Amazon S3 パスに設定するか、パイプライン YAML のdatabase:フィールドを設定し、AWS Glue データベースに Amazon S3 パスまで設定されたLocationUriがあることを確認します -
ストリーミングテーブルによるクロスラン増分処理では、ストリーミングテーブルのデータとチェックポイントの状態が Amazon S3 に保持されている必要があります。例えば、Hive テーブルまたは AWS Glue マネージドテーブル (Iceberg 以外) では、データベース
LocationUriを Amazon S3 パスに設定する必要がありますが、Apache Iceberg テーブルはテーブルメタデータ自体を管理します。
重要
Hive または AWS Glue マネージド (Iceberg 以外) テーブルにパイプライン YAML の database: フィールドを使用する場合は、対応する AWS Glue データベースの LocationUri を Amazon S3 パスに設定する必要があります。LocationUri は、マネージドストリーミングテーブル (およびその _spark_metadata ログ) を Amazon S3 に配置するものです。これにより、それらのテーブルの永続的なクロスラン増分処理が可能になります。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" }'
パイプラインの作成
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 |
いいえ | 出力テーブルの AWS Glue データベース。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=...) を組み合わせて使用します。
ステップ 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(オプション。データカタログにテーブルを登録)
以下の例では、AWS 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" }'
パイプラインの実行
StartJobRun を使用して SDP パイプラインを実行します。実行時に渡されたジョブ引数を使用して、実行動作を制御できます。
実行モード
パイプラインの実行を制御するには、以下の引数を StartJobRun に渡します。
--conf spark.glue.sdp.jobMode-
実行モードの制御:
RUN(デフォルト): パイプラインを正常に実行します。VALIDATE: データを書き込まずに YAML 構文、依存関係の解決、SQL/Python のコンパイルをチェックするドライランを実行します。
--conf spark.glue.sdp.runMode-
更新するデータセットを制御します。
- デフォルト (実行モードフラグなし)
-
すべてのデータセットを実行します。マテリアライズドビューは完全に再計算します。ストリーミングテーブルは、最後のチェックポイント以降の新しいデータのみを処理します。
--refresh <dataset>-
指定されたデータセットのみを更新します。ストリーミングテーブルは新しいデータの増分処理を行います。マテリアライズドビューは完全に再計算します。
--full-refresh <dataset>-
指定されたデータセットのみをリセットして再計算します。ストリーミングテーブルの場合、これはチェックポイントをリセットし、すべてのデータを最初から再処理します。
--full-refresh-all-
すべてのデータセットをリセットして再計算します。
SDP での Iceberg テーブルの使用
Apache Iceberg は、Hive または AWS Glue マネージドテーブルで使用されるファイルベースの _spark_metadata ログに依存しないため、耐久性のあるクロスラン増分処理を必要とするストリーミングテーブルに対して推奨されるテーブル形式です。SDP で Iceberg を設定する方法は 2 つあります。データカタログへの出力テーブルの登録の有無に応じて適切な方を使用してください。
オプション 1: データカタログを使用した Iceberg (推奨)
このオプションは、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: を AWS Glue データベースに設定します。AWS Glue ジョブを作成するときは、--enable-spark-declarative-pipeline を true に設定します。Iceberg テーブルには --enable-glue-datacatalog を設定しないでください。
注記
Iceberg カタログは、テーブルデータとメタデータ用として独自の warehouse の場所を使用します。このオプションを使用する場合、Iceberg テーブル自体に spark.sql.warehouse.dir またはデータベース LocationUri を設定する必要はありません。
オプション 2: ファイルベースの (Hadoop) カタログを使用した Iceberg
データカタログへの登録が不要である場合は、このオプションを使用します。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 は、AWS SDK を介してデータカタログに直接テーブルを登録するため、Iceberg には --enable-glue-datacatalog を設定しません。データカタログに Iceberg テーブルを登録しようとして、デフォルトの SparkSessionCatalog (type: hive) と --enable-glue-datacatalog を組み合わせて Iceberg を設定しないでください。AWS Glue 6.0 ではこの組み合わせは失敗します。
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")
考慮事項と制限事項
SDP を使用する場合は、以下を考慮してください。
-
マテリアライズドビューは常に完全に再計算します。増分更新はサポートされていません。増分ワークロードには、ストリーミングテーブルを使用します。
-
ストリーミングテーブル Python API。
@dp.append_flow(target=...)でdp.create_streaming_table()を使用します。 -
ストリーミングテーブルのクロスラン増分処理。ストリーミングテーブルは、データとチェックポイントの状態が Amazon S3 に保持されている場合にのみ、クロスラン増分処理をサポートします。Hive テーブルまたは AWS Glue マネージドテーブル (Iceberg 以外) では、データベース
LocationUriを Amazon S3 パスに設定する必要がありますが、Apache Iceberg テーブルはテーブルメタデータ自体を管理します。 -
Hive または AWS Glue マネージド (Iceberg 以外) テーブルに必要なデータベース LocationUri。
catalog-impl=GlueCatalogとともにデータカタログを使用する Iceberg テーブルには不要です。詳細については、「前提条件」を参照してください。 -
データ品質の期待値。インラインデータ品質アノテーションは、現在の SDP フレームワークではサポートされていません。
-
ダウンストリームのクエリ関数での
withColumnの回避。ダウンストリームデータセット (マテリアライズドビューなど) がspark.table(...)を使用してアップストリームパイプラインデータセットから読み取り、.withColumn(...)を適用する場合、SDP が 2 回目以降の実行でデータセット間の依存関係を検出できない可能性があります。これにより、ダウンストリームは前回の実行から古いデータを読み取ることになります (1 回分の遅延)。この問題を回避するには、.select(...)を使用する代わりに、派生した列を.withColumn(...)内で表現します。また、クエリ関数内でプラン解決 (.schemaまたは.collectなど) を強制するオペレーションは避けてください。 -
移行ツールなし。他のパイプラインフレームワークからの自動移行はサポートされていません。テーブルを段階的に移行します。SDP は既存のカタログテーブルから読み取ることができます。
-
スケジューリング。SDP ジョブは、他の AWS Glue ジョブ (AWS Glue トリガー、Amazon EventBridge、Apache Airflow) と同じスケジューリングメカニズムを使用します。
命令スクリプトから SDP への移行
既存の命令的な Spark スクリプトを SDP に段階的に移行できます。
-
1 回の
spark.sql(...).write.saveAsTable(...)呼び出しをCREATE MATERIALIZED VIEWSQL ステートメントに変換することで、1 つのテーブルから開始します。 -
テーブルを増分的に追加します。SDP は混合依存関係を処理します。SDP テーブルは、パイプラインに含まれていない既存のカタログテーブルから読み取ることができます。
-
移行中は、両方のパターンを並行して実行します。SDP ジョブと命令型ジョブは共存できます。
SDP は、既存のデータカタログテーブル、外部テーブル、クロスデータベース参照など、SparkSession を介してアクセス可能な任意のテーブルを参照できます。