Sparkジョブ定義(SJD)は、ユーザーがApache SparkジョブをFabric上で定義し実行できるようにするFabricの一種です。 Sparkジョブ定義API v2では、ユーザーが新しいフォーマット「 SparkJobDefinitionV2」でSparkジョブ定義項目を作成・更新できます。 v2 形式を使用する主な利点は、ユーザーが 1 回の API 呼び出しでメインの実行可能ファイルやその他のライブラリ ファイルを管理できる点です。ストレージ API を使用してファイルを個別にアップロードする代わりに、ファイルを管理するためにこれ以上ストレージ トークンは必要ありません。
[前提条件]
- Fabric REST API にアクセスするには、Microsoft Entra トークンが必要です。 トークンを取得するには、MSAL (Microsoft Authentication Library) ライブラリを使用することをお勧めします。 詳細については、「 MSAL での認証フローのサポート」を参照してください。
Fabric REST APIは、FabricアイテムのCRUD操作のための統一エンドポイントを定義しています。 エンドポイントが https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items です。
Sparkジョブ定義 v2フォーマット概要
Sparkジョブ定義項目の管理におけるペイロードでは、 definition フィールドがSparkジョブ定義項目の詳細設定を指定するために使われます。
definition フィールドには、formatとpartsの 2 つのサブフィールドが含まれています。
formatフィールドはSparkジョブ定義項目のフォーマットを指定しており、v2フォーマットではSparkJobDefinitionV2されるべきです。
partsフィールドは、Sparkジョブ定義項目の詳細設定を含む配列です。
parts配列内の各項目は、詳細なセットアップの一部を表します。 各パーツには、 path、 payload、 payloadTypeの 3 つのサブフィールドが含まれています。
path フィールドはパーツのパスを指定し、payload フィールドは base64 でエンコードされたパーツの内容を指定し、payloadType フィールドはペイロードの種類を指定します。これはInlineBase64する必要があります。
Important
このv2フォーマットは、.pyまたは.scalaのファイル形式のみでSparkジョブ定義をサポートしています。 .jarファイル形式はサポートされていません。
メインの定義ファイルやその他のリブファイルでSparkジョブ定義項目を作成します
以下の例では、次のSparkジョブ定義項目を作成します。
- 名前は
SJDHelloWorld。 - メイン定義ファイルは
main.pyで、これはデフォルトのレイクハウスからCSVファイルを読み込み、同じレイクハウスにDeltaテーブルとして保存することです。 - その他の lib ファイルは
libs.pyで、CSV ファイルと Delta テーブルの名前を返すユーティリティ関数があります。 - デフォルトのレイクハウスは特定のレイクハウスアイテムIDに設定されています。
以下はSparkジョブ定義項目の作成に関する詳細なペイロードです。
{
"displayName": "SJDHelloWorld",
"type": "SparkJobDefinition",
"definition": {
"format": "SparkJobDefinitionV2",
"parts": [
{
"path": "SparkJobDefinitionV1.json",
"payload": "<REDACTED>",
"payloadType": "InlineBase64"
},
{
"path": "Main/main.py",
"payload": "<REDACTED>",
"payloadType": "InlineBase64"
},
{
"path": "Libs/lib1.py",
"payload": "<REDACTED>",
"payloadType": "InlineBase64"
}
]
}
}
詳細なセットアップをデコードまたはエンコードするには、Python で次のヘルパー関数を使用できます。 同じジョブを実行できる https://www.base64decode.org/ などの他のオンライン ツールもあります。
import base64
def json_to_base64(json_data):
# Serialize the JSON data to a string
json_string = json.dumps(json_data)
# Encode the JSON string as bytes
json_bytes = json_string.encode('utf-8')
# Encode the bytes as Base64
base64_encoded = base64.b64encode(json_bytes).decode('utf-8')
return base64_encoded
def base64_to_json(base64_data):
# Decode the Base64-encoded string to bytes
base64_bytes = base64_data.encode('utf-8')
# Decode the bytes to a JSON string
json_string = base64.b64decode(base64_bytes).decode('utf-8')
# Deserialize the JSON string to a Python dictionary
json_data = json.loads(json_string)
return json_data
HTTPコード202の応答は、Sparkジョブ定義項目が正常に作成されたことを示します。
Sparkジョブの定義と定義パーツをv2形式で取得できます
新しいv2形式では、定義パーツを含むSparkジョブ定義項目を取得する際、メイン定義ファイルやその他のリブファイルのファイル内容がすべて、 parts フィールドの下にエンコードされたレスポンスペイロード(base64)に含まれます。 以下は定義パーツ付きのSparkジョブ定義項目を取得する例です:
- まず、エンドポイント
https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items/{sjditemid}/getDefinitionParts?format=SparkJobDefinitionV2に POST 要求を行います。 形式クエリ パラメーターの値がSparkJobDefinitionV2されていることを確認します。 - 次に、応答ヘッダーで、HTTP 状態コードを確認します。 HTTP コード 202 は、要求が正常に受け入れられたかどうかを示します。 応答ヘッダーから
x-ms-operation-id値をコピーします。 - 最後に、コピーした
https://api.fabric.microsoft.com/v1/operations/{operationId}値を使用してエンドポイントx-ms-operation-idに GET 要求を行って、操作の結果を取得します。 レスポンスペイロードのdefinitionフィールドには、Sparkジョブ定義項目の詳細設定が含まれており、メインの定義ファイルやpartsフィールドの下にある他のlibファイルが含まれます。
Sparkジョブの定義項目をメインの定義ファイルや他のリブファイルでv2形式で更新してください
既存のSparkジョブ定義項目にメインの定義ファイルやv2形式のlibファイルで更新するには、create操作と同様のペイロード構造を使えます。 前節で作成したSparkジョブ定義項目の更新例は以下の通りです:
{
"displayName": "SJDHelloWorld",
"type": "SparkJobDefinition",
"definition": {
"format": "SparkJobDefinitionV2",
"parts": [
{
"path": "SparkJobDefinitionV1.json",
"payload": "<REDACTED>",
"payloadType": "InlineBase64"
},
{
"path": "Main/main.py",
"payload": "<REDACTED>",
"payloadType": "InlineBase64"
},
{
"path": "Libs/lib2.py",
"payload": "<REDACTED>",
"payloadType": "InlineBase64"
}
]
}
}
上記のペイロードでは、ファイルに次の変更が加えられます。
- main.py ファイルが新しいコンテンツで更新されます。
- この lib1.py はこのSparkジョブ定義項目から削除され、OneLakeストレージからも削除されます。
- 新しい lib2.py ファイルがこのSparkジョブ定義項目に追加され、OneLakeストレージにアップロードされます。
Sparkジョブ定義項目を更新するには、上記のペイロードを含むエンドポイント https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items/{sjditemid} にPOSTリクエストを送ります。 HTTPコード202の応答は、Sparkジョブ定義項目が正常に更新されたことを示します。