Cómo crear y actualizar una definición de trabajo en Spark con la API REST de Microsoft Fabric

La API REST de Fabric proporciona un punto final de servicio para operaciones CRUD de elementos Fabric. En este tutorial, repasamos un escenario de extremo a extremo sobre cómo crear y actualizar un elemento de definición de trabajo en Spark. Tres pasos de alto nivel son necesarios:

  1. Crea un elemento de definición de trabajo en Spark con algún estado inicial.
  2. Cargue el archivo de definición principal y otros archivos lib.
  3. Actualiza el elemento de definición de trabajo de Spark con la URL de OneLake del archivo principal de definición y otros archivos de liberación.

Requisitos previos

  • Se requiere un token de Microsoft Entra para acceder a la API REST de Fabric. Se recomienda la biblioteca MSAL para obtener el token. Para más información, consulte Compatibilidad con el flujo de autenticación de MSAL.
  • Se requiere un token de almacenamiento para acceder a la API de OneLake. Para más información, consulte MSAL para Python.

Crea un elemento de definición de trabajo en Spark con el estado inicial

La API Fabric REST define un punto final unificado para las operaciones CRUD de los elementos Fabric. El extremo es https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items.

Los detalles del elemento se especifican dentro del cuerpo de la solicitud. Aquí tienes un ejemplo del cuerpo de la solicitud para crear un elemento de definición de trabajo en Spark:

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

En este ejemplo, el elemento de definición de trabajo de Spark se denomina SJDHelloWorld. El payload campo es el contenido codificado en base64 de la configuración detallada. Después de la descodificación, el contenido es:

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

Estas son dos funciones de asistencia para codificar y descodificar la configuración detallada:

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

Aquí tienes el fragmento de código para crear un elemento de definición de trabajo en 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)

Cargar el archivo de definición principal y otros archivos lib

Se requiere un token de almacenamiento para cargar el archivo a OneLake. Esta es una función auxiliar para obtener el token de almacenamiento:

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']

Ahora tenemos creado un elemento de definición de trabajo en Spark. Para que sea ejecutable, es necesario configurar el archivo de definición principal y las propiedades necesarias. El punto de conexión para cargar el archivo para este elemento SJD es https://onelake.dfs.fabric.microsoft.com/{workspaceId}/{sjditemid}. Se debe usar el mismo "workspaceId" del paso anterior. El valor de "sjditemid" podía encontrarse en el cuerpo de respuesta del paso anterior. Este es el fragmento de código para configurar el archivo de definición 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())

Siga el mismo proceso para cargar los otros archivos lib si es necesario.

Actualizar el elemento de definición de trabajo de Spark con la URL de OneLake del archivo principal de definición y otros archivos de libación

Hasta ahora, hemos creado un elemento de definición de trabajo en Spark con un estado inicial y hemos subido el archivo principal de definición y otros archivos de liberación. El último paso es actualizar el elemento de definición de trabajo de Spark para establecer las propiedades de URL del archivo principal de definición y otros archivos de liberación. El punto final para actualizar el elemento de definición de tareas de Spark es https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items/{sjditemid}. Se deben usar los mismos "workspaceId" y "sjditemid" de los pasos anteriores. Aquí tienes el fragmento de código para actualizar el elemento de definición del trabajo de 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)

Para resumir todo el proceso, se necesitan tanto la API REST de Fabric como la API OneLake para crear y actualizar un elemento de definición de trabajo en Spark. La API Fabric REST se utiliza para crear y actualizar el elemento de definición de trabajo de Spark. La API de OneLake se usa para cargar el archivo de definición principal y otros archivos lib. El archivo de definición principal y otros archivos lib se cargan primero en OneLake. Luego, las propiedades URL del archivo principal de definición y otros archivos de libación se establecen en el elemento de definición del trabajo de Spark.