Hoe maak en werk je een Spark-taakdefinitie bij met V2-formaat via de Microsoft Fabric REST API

Spark job definitie (SJD) is een type Fabric-item waarmee gebruikers Apache Spark-taken in Fabric kunnen definiëren en uitvoeren. De Spark job definition API v2 stelt gebruikers in staat om Spark-job definitie-items te maken en bij te werken met een nieuw formaat genaamd SparkJobDefinitionV2. Het belangrijkste voordeel van het gebruik van de v2-indeling is dat gebruikers het belangrijkste uitvoerbare bestand en andere bibliotheekbestanden kunnen beheren met één API-aanroep, in plaats van opslag-API te gebruiken om bestanden afzonderlijk te uploaden, is er geen opslagtoken meer nodig voor het beheren van bestanden.

Vereiste voorwaarden

  • Er is een Microsoft Entra-token vereist om toegang te krijgen tot de Fabric REST API. De MSAL-bibliotheek (Microsoft Authentication Library) wordt aanbevolen om het token op te halen. Zie ondersteuning voor verificatiestromen in MSAL voor meer informatie.

De Fabric REST API definieert een geïntegreerd eindpunt voor CRUD-operaties van Fabric-items. Het eindpunt is https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items.

Overzicht van het formaat van Spark-job definitie v2

In de payload van het beheren van een Spark-taakdefinitie-item wordt het definition veld gebruikt om de gedetailleerde opzet van het Spark-taakdefinitie-item te specificeren. Het definition veld bevat twee subvelden: format en parts. Het format veld specificeert het formaat van het Spark-taakdefinitie-item, dat voor het v2-formaat zou moeten zijn SparkJobDefinitionV2 .

Het parts veld is een array die de gedetailleerde opzet van het Spark-taakdefinitie-item bevat. Elk item in de parts matrix vertegenwoordigt een deel van de gedetailleerde installatie. Elk deel bevat drie subvelden: path, payloaden payloadType. Het path veld geeft het pad van het deel op, het payload veld geeft de inhoud van het deel dat base64 is gecodeerd en het payloadType veld specificeert het type van de nettolading, dat moet zijn InlineBase64.

Belangrijk

Dit v2-formaat ondersteunt alleen Spark-taakdefinities met bestandsformaten van .py of .scala. De .jar bestandsindeling wordt niet ondersteund.

Maak een Spark-taakdefinitie-item aan met het hoofddefinitiebestand en andere libbestanden

In het volgende voorbeeld maken we een Spark-taakdefinitie-item dat:

  1. De naam is SJDHelloWorld.
  2. Het hoofddefinitiebestand is main.py, wat betekent dat je een CSV-bestand uit het standaard lakehouse leest en als een Delta-tabel opslaat naar hetzelfde lakehouse.
  3. Een ander lib-bestand is libs.py, dat een hulpprogrammafunctie heeft om de naam van het CSV-bestand en de Delta-tabel te retourneren.
  4. Het standaard lakehouse is ingesteld op een specifieke lakehouse item-ID.

Hieronder volgt de gedetailleerde payload voor het maken van het Spark-taakdefinitie-item.

{
  "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"
      }
    ]
  }
}

Als u de gedetailleerde installatie wilt decoderen of coderen, kunt u de volgende helperfuncties in Python gebruiken. Er zijn ook andere onlinehulpprogramma's, zoals https://www.base64decode.org/ die dezelfde taak kunnen uitvoeren.

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

Een HTTP-code 202-antwoord geeft aan dat het Spark-taakdefinitie-item succesvol is aangemaakt.

Krijg Spark-jobdefinitie met definitieonderdelen onder v2-formaat

Met het nieuwe v2-formaat worden bij het ontvangen van een Spark-taakdefinitie-item met definitiedelen de bestandsinhoud van het hoofddefinitiebestand en andere libbestanden allemaal opgenomen in de responspayload, base64 gecodeerd onder het parts veld. Hier is een voorbeeld van het krijgen van een Spark-taakdefinitie-item met definitiedelen:

  1. Maak eerst een POST-aanvraag naar het eindpunt https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items/{sjditemid}/getDefinitionParts?format=SparkJobDefinitionV2. Zorg ervoor dat de waarde van de queryparameter 'format' SparkJobDefinitionV2 is.
  2. Controleer vervolgens in de antwoordheaders de HTTP-statuscode. Een HTTP-code 202 geeft aan dat de aanvraag is geaccepteerd. Kopieer de x-ms-operation-id waarde uit de antwoordheaders.
  3. Breng ten slotte een GET-aanvraag naar het eindpunt https://api.fabric.microsoft.com/v1/operations/{operationId} met de gekopieerde x-ms-operation-id waarde om het bewerkingsresultaat op te halen. In de response payload bevat het definition veld de gedetailleerde opzet van het Spark-taakdefinitie-item, inclusief het hoofddefinitiebestand en andere libbestanden onder het parts veld.

Werk het Spark-taakdefinitie-item bij met het hoofddefinitiebestand en andere libbestanden onder v2-formaat

Om een bestaand Spark-taakdefinitie-item bij te werken met het hoofddefinitiebestand en andere libbestanden onder het v2-formaat, kun je een vergelijkbare payloadstructuur gebruiken als de create-operatie. Hier is een voorbeeld van het bijwerken van het Spark-taakdefinitie-item dat in de vorige sectie is aangemaakt:

{
  "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"
      }
    ]
  }
}

Met de bovenstaande payload worden de volgende wijzigingen aangebracht in de bestanden:

  1. Het main.py-bestand wordt bijgewerkt met nieuwe inhoud.
  2. De lib1.py wordt verwijderd uit dit Spark-taakdefinitie-item en ook uit de OneLake-opslag.
  3. Een nieuw lib2.py-bestand wordt toegevoegd aan dit Spark-taakdefinitie-item en geüpload naar de OneLake-opslag.

Om het Spark-taakdefinitie-item bij te werken, doe je een POST-verzoek aan het eindpunt https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items/{sjditemid} met bovenstaande payload. Een HTTP-code 202-antwoord geeft aan dat het Spark-taakdefinitie-item succesvol is bijgewerkt.