Jak stworzyć i zaktualizować definicję zadania w Spark w formacie V2 za pomocą Microsoft Fabric REST API

Definicja zadania Spark (SJD) to rodzaj elementu Fabric, który pozwala użytkownikom definiować i uruchamiać zadania Apache Spark w Fabric. API definicji zadań Spark v2 pozwala użytkownikom tworzyć i aktualizować elementy definicji zadań Spark w nowym formacie o nazwie SparkJobDefinitionV2. Główną zaletą korzystania z formatu w wersji 2 jest to, że umożliwia użytkownikom zarządzanie głównym plikiem wykonywalnym i innymi plikami biblioteki przy użyciu jednego wywołania interfejsu API. Zamiast używania interfejsu API magazynu do oddzielnego przesyłania plików, do zarządzania plikami nie jest potrzebny żaden token magazynu.

Wymagania wstępne

  • Token Microsoft Entra jest wymagany do uzyskania dostępu do interfejsu API REST Fabric. Biblioteka MSAL (Biblioteka uwierzytelniania firmy Microsoft) jest zalecana do uzyskania tokenu. Aby uzyskać więcej informacji, zobacz Obsługa przepływu uwierzytelniania w usłudze MSAL.

API Fabric REST definiuje zunifikowany punkt końcowy dla operacji CRUD dla elementów Fabric. Punkt końcowy jest https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items.

Przegląd formatu definicji zadania w Spark v2

W ładunku zarządzania elementem definition definicji zadania Spark pole służy do określenia szczegółowej konfiguracji elementu definicji zadania Spark. Pole definition zawiera dwa pola podrzędne: format i parts. Pole format określa format elementu definicji zadania w Spark, który powinien dotyczyć SparkJobDefinitionV2 formatu v2.

Pole parts to tablica zawierająca szczegółową konfigurację elementu definicji zadania w Sparku. Każdy element w tablicy parts reprezentuje część szczegółowej konfiguracji. Każda część zawiera trzy pola podrzędne: path, payloadi payloadType. Pole path określa ścieżkę części, payload pole określa zawartość części, która jest zakodowana w formacie base64, a payloadType pole określa typ ładunku, który powinien mieć wartość InlineBase64.

Ważne

Ten format v2 obsługuje tylko definicje zadań Spark z formatami plików .py lub .scala. Format pliku .jar nie jest obsługiwany.

Stwórz element definicji zadania w Spark z głównym plikiem definicji oraz innymi plikami lib

W poniższym przykładzie utworzymy element definicji zadania w Spark, który:

  1. Nazwa to SJDHelloWorld.
  2. Główny plik definicyjny to main.py, czyli odczyt pliku CSV z domyślnego lakehouse i zapis jako tabela Delta z powrotem do tego samego lakehouse.
  3. Inny plik lib to libs.py, który ma funkcję narzędzia, która zwraca nazwę pliku CSV i tabelę delta.
  4. Domyślny domek nad jeziorem jest ustawiony na konkretny identyfikator przedmiotu przy domku nad jeziorem.

Poniżej znajduje się szczegółowy ładunek do tworzenia elementu definicji zadania 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"
      }
    ]
  }
}

Aby zdekodować lub zakodować szczegółową konfigurację, możesz użyć następujących funkcji pomocnika w języku Python. Istnieją również inne narzędzia online, takie jak https://www.base64decode.org/ te, które mogą wykonywać to samo zadanie.

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

Odpowiedź z kodem HTTP 202 wskazuje, że element definicji zadania w Spark został pomyślnie utworzony.

Pobierz definicję pracy w Spark z częściami definicyjnymi w formacie v2

W nowym formacie v2, gdy otrzymujemy element definicji zadania Spark z elementami definicji, zawartość pliku głównego pliku definicji oraz innych plików lib jest zawarta w payloadie odpowiedzi, base64 zakodowanym pod polem parts . Oto przykład uzyskania elementu definicji pracy w Spark z elementami definicyjnymi:

  1. Najpierw utwórz żądanie POST do punktu końcowego https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items/{sjditemid}/getDefinitionParts?format=SparkJobDefinitionV2. Upewnij się, że wartość parametru zapytania formatu to SparkJobDefinitionV2.
  2. Następnie w nagłówkach odpowiedzi sprawdź kod stanu HTTP. Kod HTTP 202 wskazuje, że żądanie zostało pomyślnie zaakceptowane. Skopiuj wartość x-ms-operation-id z nagłówków odpowiedzi.
  3. Na koniec utwórz żądanie GET do punktu końcowego https://api.fabric.microsoft.com/v1/operations/{operationId} z skopiowaną x-ms-operation-id wartością, aby uzyskać wynik operacji. W ładunku definition odpowiedzi pole zawiera szczegółową konfigurację elementu definicji zadania Spark, w tym główny plik definicji oraz inne pliki lib pod tym parts polem.

Zaktualizuj element definicji zadania w Spark o główny plik definicji oraz inne pliki lib w formacie v2

Aby zaktualizować istniejący element definicji zadania w Spark o główny plik definicji oraz inne pliki lib w formacie v2, możesz użyć podobnej struktury payload jak operacja create. Oto przykład aktualizacji elementu definicji zadania w Spark utworzony w poprzedniej sekcji:

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

W przypadku powyższego ładunku następujące zmiany są wprowadzane do plików:

  1. Plik main.py zostanie zaktualizowany o nową zawartość.
  2. lib1.py jest usuwany z tej definicji zadania Spark oraz usuwany z pamięci OneLake.
  3. Do tego elementu definicji zadania w Spark dodaje się nowy plik lib2.py i przesyła go do pamięci OneDrive.

Aby zaktualizować element definicji zadania Spark, wykonaj żądanie POST do punktu końcowego https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items/{sjditemid} z powyższym ładunkiem. Odpowiedź HTTP kodu 202 wskazuje, że element definicji zadania w Spark został pomyślnie zaktualizowany.