Comment créer et mettre à jour une définition de job Spark avec l’API Microsoft Fabric REST

L’API Fabric REST fournit un point de terminaison de service pour les opérations CRUD des éléments Fabric. Dans ce tutoriel, nous présentons un scénario de bout en bout pour créer et mettre à jour un élément de définition de job Spark. Il existe trois étapes générales :

  1. Créez un élément de définition du poste Spark avec un état initial.
  2. Chargez le fichier de définition principal et d’autres fichiers lib.
  3. Mettez à jour l’élément de définition du job Spark avec l’URL OneLake du fichier principal de définition et d’autres fichiers de libération.

Prérequis

  • Un jeton Microsoft Entra est requis pour accéder à l’API REST Fabric. La bibliothèque MSAL est recommandée pour obtenir le jeton. Pour plus d’informations, consultez l’article Prise en charge du flux d’authentification dans MSAL.
  • Un jeton de stockage est nécessaire pour accéder à l’API OneLake. Pour plus d’informations, consultez l’article MSAL pour Python.

Créez un élément de définition de job Spark avec l’état initial

L’API Fabric REST définit un point de terminaison unifié pour les opérations CRUD des éléments Fabric. Le point de terminaison a la valeur https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items.

Les détails de l’élément sont spécifiés dans le corps de la demande. Voici un exemple du corps de la demande pour créer un élément de définition de poste Spark :

{
    "displayName": "SJDHelloWorld",
    "type": "SparkJobDefinition",
    "definition": {
        "format": "SparkJobDefinitionV1",
        "parts": [
            {
                "path": "SparkJobDefinitionV1.json",
                "payload": "<REDACTED>",
                "payloadType": "InlineBase64"
            }
        ]
    }
}

Dans cet exemple, l’élément de définition du job Spark est nommé SJDHelloWorld. Le payload champ est le contenu codé en base64 de la configuration détaillée. Après le décodage, le contenu est le suivant :

{
    "executableFile":null,
    "defaultLakehouseArtifactId":"",
    "mainClass":"",
    "additionalLakehouseIds":[],
    "retryPolicy":null,
    "commandLineArguments":"",
    "additionalLibraryUris":[],
    "language":"",
    "environmentArtifactId":null
}

Voici deux fonctions d’assistance pour coder et décoder la configuration détaillée :

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

Voici l’extrait de code pour créer un élément de définition de la tâche Spark :

import requests

bearerToken = "<REDACTED>"  # Replace this token with the real AAD token

headers = {
    "Authorization": f"Bearer {bearerToken}", 
    "Content-Type": "application/json"  # Set the content type based on your request
}

payload = "<REDACTED>"

# Define the payload data for the POST request
payload_data = {
    "displayName": "SJDHelloWorld",
    "Type": "SparkJobDefinition",
    "definition": {
        "format": "SparkJobDefinitionV1",
        "parts": [
            {
                "path": "SparkJobDefinitionV1.json",
                "payload": payload,
                "payloadType": "InlineBase64"
            }
        ]
    }
}

# Make the POST request with Bearer authentication
sjdCreateUrl = f"https://api.fabric.microsoft.com//v1/workspaces/{workspaceId}/items"
response = requests.post(sjdCreateUrl, json=payload_data, headers=headers)

Charger le fichier de définition principal et d’autres fichiers lib

Un jeton de stockage est nécessaire pour charger le fichier dans OneLake. Voici une fonction d’assistance qui permet d’obtenir le jeton de stockage :

import msal

def getOnelakeStorageToken():
    app = msal.PublicClientApplication(
        "<REDACTED>",  # This field should be the client ID 
        authority="https://login.microsoftonline.com/microsoft.com")

    result = app.acquire_token_interactive(scopes=["https://storage.azure.com/.default"])

    print(f"Successfully acquired AAD token with storage audience:{result['access_token']}")

    return result['access_token']

Nous avons maintenant créé un élément de définition de job dans Spark. Pour le rendre exécutable, nous devons configurer le fichier de définition principal et les propriétés requises. Le point de terminaison pour le chargement du fichier de cet élément SJD est https://onelake.dfs.fabric.microsoft.com/{workspaceId}/{sjditemid}. Le même « workspaceId » de l’étape précédente doit être utilisé. La valeur de « sjditemid » pouvait être trouvée dans le corps de réponse de l’étape précédente. Voici l’extrait de code permettant de configurer le fichier de définition principal :

import requests

# Three steps are required: create file, append file, flush file

onelakeEndPoint = "https://onelake.dfs.fabric.microsoft.com/workspaceId/sjditemid"  # Replace the ID of workspace and item with the right one
mainExecutableFile = "main.py"  # The name of the main executable file
mainSubFolder = "Main"  # The sub folder name of the main executable file. Don't change this value


onelakeRequestMainFileCreateUrl = f"{onelakeEndPoint}/{mainSubFolder}/{mainExecutableFile}?resource=file"  # The URL for creating the main executable file via the 'file' resource type
onelakePutRequestHeaders = {
    "Authorization": f"Bearer {onelakeStorageToken}",  # The storage token can be achieved from the helper function above
}

onelakeCreateMainFileResponse = requests.put(onelakeRequestMainFileCreateUrl, headers=onelakePutRequestHeaders)
if onelakeCreateMainFileResponse.status_code == 201:
    # Request was successful
    print(f"Main File '{mainExecutableFile}' was successfully created in OneLake.")

# With the previous step, the main executable file is created in OneLake. Now we need to append the content of the main executable file

appendPosition = 0
appendAction = "append"

### Main File Append.
mainExecutableFileSizeInBytes = 83  # The size of the main executable file in bytes
onelakeRequestMainFileAppendUrl = f"{onelakeEndPoint}/{mainSubFolder}/{mainExecutableFile}?position={appendPosition}&action={appendAction}"
mainFileContents = "<REDACTED>"  # The content of the main executable file, please replace this with the real content of the main executable file
mainExecutableFileSizeInBytes = 83  # The size of the main executable file in bytes, this value should match the size of the mainFileContents

onelakePatchRequestHeaders = {
    "Authorization": f"Bearer {onelakeStorageToken}",
    "Content-Type": "text/plain"
}

onelakeAppendMainFileResponse = requests.patch(onelakeRequestMainFileAppendUrl, data = mainFileContents, headers=onelakePatchRequestHeaders)
if onelakeAppendMainFileResponse.status_code == 202:
    # Request was successful
    print(f"Successfully accepted main file '{mainExecutableFile}' append data.")

# With the previous step, the content of the main executable file is appended to the file in OneLake. Now we need to flush the file

flushAction = "flush"

### Main File flush
onelakeRequestMainFileFlushUrl = f"{onelakeEndPoint}/{mainSubFolder}/{mainExecutableFile}?position={mainExecutableFileSizeInBytes}&action={flushAction}"
print(onelakeRequestMainFileFlushUrl)
onelakeFlushMainFileResponse = requests.patch(onelakeRequestMainFileFlushUrl, headers=onelakePatchRequestHeaders)
if onelakeFlushMainFileResponse.status_code == 200:
    print(f"Successfully flushed main file '{mainExecutableFile}' contents.")
else:
    print(onelakeFlushMainFileResponse.json())

Suivez le même processus pour charger les autres fichiers lib, si nécessaire.

Mettez à jour l’élément de définition du job Spark avec l’URL OneLake du fichier principal de définition et d’autres fichiers de libération

Jusqu’à présent, nous avons créé un élément de définition de job Spark avec un certain état initial et téléchargé le fichier principal ainsi que d’autres fichiers de libération. La dernière étape consiste à mettre à jour l’élément de définition du job Spark pour définir les propriétés URL du fichier principal de définition et des autres fichiers de lib. Le point de terminaison pour mettre à jour l’élément de définition du travail Spark est https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items/{sjditemid}. Les mêmes « workspaceID » et « sjditemid » que les étapes précédentes doivent être utilisés. Voici l’extrait de code pour mettre à jour l’élément de définition du job Spark :

mainAbfssPath = f"abfss://{workspaceId}@onelake.dfs.fabric.microsoft.com/{sjditemid}/Main/{mainExecutableFile}"  # The workspaceId and sjditemid are the same as previous steps, the mainExecutableFile is the name of the main executable file
libsAbfssPath = f"abfss://{workspaceId}@onelake.dfs.fabric.microsoft.com/{sjditemid}/Libs/{libsFile}"  # The workspaceId and sjditemid are the same as previous steps, the libsFile is the name of the libs file
defaultLakehouseId = '<REDACTED>'  # Replace this with the real default lakehouse ID

updateRequestBodyJson = {
    "executableFile": mainAbfssPath,
    "defaultLakehouseArtifactId": defaultLakehouseId,
    "mainClass": "",
    "additionalLakehouseIds": [],
    "retryPolicy": None,
    "commandLineArguments": "",
    "additionalLibraryUris": [libsAbfssPath],
    "language": "Python",
    "environmentArtifactId": None}

# Encode the bytes as a Base64-encoded string
base64EncodedUpdateSJDPayload = json_to_base64(updateRequestBodyJson)

# Print the Base64-encoded string
print("Base64-encoded JSON payload for SJD Update:")
print(base64EncodedUpdateSJDPayload)

# Define the API URL
updateSjdUrl = f"https://api.fabric.microsoft.com//v1/workspaces/{workspaceId}/items/{sjditemid}/updateDefinition"

updatePayload = base64EncodedUpdateSJDPayload
payloadType = "InlineBase64"
path = "SparkJobDefinitionV1.json"
format = "SparkJobDefinitionV1"
Type = "SparkJobDefinition"

# Define the headers with Bearer authentication
bearerToken = "<REDACTED>"  # Replace this token with the real AAD token

headers = {
    "Authorization": f"Bearer {bearerToken}", 
    "Content-Type": "application/json"  # Set the content type based on your request
}

# Define the payload data for the POST request
payload_data = {
    "displayName": "sjdCreateTest11",
    "Type": Type,
    "definition": {
        "format": format,
        "parts": [
            {
                "path": path,
                "payload": updatePayload,
                "payloadType": payloadType
            }
        ]
    }
}


# Make the POST request with Bearer authentication
response = requests.post(updateSjdUrl, json=payload_data, headers=headers)
if response.status_code == 200:
    print("Successfully updated SJD.")
else:
    print(response.json())
    print(response.status_code)

Pour récapituler tout le processus, l’API Fabric REST et l’API OneLake sont nécessaires pour créer et mettre à jour un élément de définition de job Spark. L’API Fabric REST est utilisée pour créer et mettre à jour l’élément de définition du job Spark. L’API OneLake est utilisée pour charger le fichier de définition principal et d’autres fichiers lib. Le fichier de définition principal et d’autres fichiers lib sont d’abord chargés dans OneLake. Ensuite, les propriétés URL du fichier principal de définition et des autres fichiers de lib sont définies dans l’élément de définition du job Spark.