Como criar e atualizar uma definição de trabalho Spark com formato V2 via API REST do Microsoft Fabric

A definição de trabalho Spark (SJD) é um tipo de item Fabric que permite aos usuários definir e executar jobs Apache Spark no Fabric. A API de definição de jobs do Spark v2 permite que os usuários criem e atualizem itens de definição de jobs do Spark com um novo formato chamado SparkJobDefinitionV2. O principal benefício de usar o formato v2 é que ele permite que os usuários gerenciem o arquivo executável principal e outros arquivos de biblioteca com uma única chamada de API, em vez de usar a API de armazenamento para carregar arquivos separadamente, nenhum token de armazenamento é necessário para gerenciar arquivos.

Pré-requisitos

  • Um token do Microsoft Entra é necessário para acessar a API REST do Fabric. A biblioteca MSAL (Biblioteca de Autenticação da Microsoft) é recomendada para obter o token. Para obter mais informações, consulte o suporte ao fluxo de autenticação na MSAL.

A API REST do Fabric define um endpoint unificado para operações CRUD dos itens Fabric. O ponto de extremidade é https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items.

Visão geral do formato da definição de trabalho do Spark v2

No payload de gerenciamento de um item de definição de trabalho Spark, o definition campo é usado para especificar a configuração detalhada do item de definição de trabalho Spark. O definition campo contém dois subcampos: format e parts. O format campo especifica o formato do item de definição de trabalho do Spark, que deve ser SparkJobDefinitionV2 para o formato v2.

O parts campo é um array que contém a configuração detalhada do item de definição do trabalho do Spark. Cada item na parts matriz representa uma parte da configuração detalhada. Cada parte contém três subcampos: path, payloade payloadType. O path campo especifica o caminho da parte, o payload campo especifica o conteúdo da parte codificada em base64 e o payloadType campo especifica o tipo da carga, que deve ser InlineBase64.

Importante

Esse formato v2 suporta apenas definições de trabalhos Spark com formatos de arquivo .py ou .scala. Não há suporte para o formato de arquivo .jar.

Crie um item de definição de trabalho no Spark com o arquivo principal de definição e outros arquivos de lib

No exemplo a seguir, vamos criar um item de definição de trabalho no Spark que:

  1. O nome é SJDHelloWorld.
  2. O arquivo principal de definição é main.py, que consiste em ler um arquivo CSV da sua casa padrão e salvar como uma tabela Delta de volta na mesma casa do lago.
  3. Outro arquivo lib é libs.py, que tem uma função de utilitário para retornar o nome do arquivo CSV e da tabela Delta.
  4. A casa padrão do lago está definida para um ID específico de item da casa do lago.

A seguir está a carga útil detalhada para criar o item de definição do trabalho 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"
      }
    ]
  }
}

Para decodificar ou codificar a configuração detalhada, você pode usar as seguintes funções auxiliares no Python. Há também outras ferramentas online, como https://www.base64decode.org/, que podem realizar o mesmo trabalho.

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

Uma resposta com código HTTP 202 indica que o item de definição de trabalho do Spark foi criado com sucesso.

Obtenha definição de trabalho do Spark com partes de definição no formato v2

Com o novo formato v2, ao receber um item de definição de trabalho Spark com partes de definição, o conteúdo do arquivo principal de definição e outros arquivos lib são todos incluídos na carga útil de resposta, codificada em base64 sob o parts campo. Aqui está um exemplo de conseguir um item de definição de trabalho do Spark com partes de definição:

  1. Primeiro, faça uma solicitação POST para o ponto de extremidade https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items/{sjditemid}/getDefinitionParts?format=SparkJobDefinitionV2. Verifique se o valor do parâmetro de consulta de formato é SparkJobDefinitionV2.
  2. Em seguida, nos cabeçalhos de resposta, verifique o código de status HTTP. Um código HTTP 202 indica que a solicitação foi aceita com êxito. Copie o valor x-ms-operation-id dos cabeçalhos de resposta.
  3. Por fim, faça uma solicitação GET para o endpoint https://api.fabric.microsoft.com/v1/operations/{operationId} com o valor copiado x-ms-operation-id para obter o resultado da operação. No payload de resposta, o definition campo contém a configuração detalhada do item de definição do trabalho Spark, incluindo o arquivo principal de definição e outros arquivos lib sob o parts campo.

Atualize o item de definição do trabalho do Spark com o arquivo principal de definição e outros arquivos de lib no formato v2

Para atualizar um item de definição de trabalho Spark existente com o arquivo principal de definição e outros arquivos lib sob o formato v2, você pode usar uma estrutura de payload semelhante à da operação de criação. Aqui está um exemplo de atualização do item de definição de trabalho do Spark criado na seção anterior:

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

Com a payload acima, as seguintes alterações são feitas nos arquivos:

  1. O arquivo main.py é atualizado com novo conteúdo.
  2. A lib1.py é excluída deste item de definição de trabalho do Spark e também removida do armazenamento do OneLake.
  3. Um novo arquivo lib2.py é adicionado a esse item de definição de trabalho do Spark e enviado para o armazenamento OneLake.

Para atualizar o item de definição de trabalho do Spark, faça uma requisição POST ao endpoint https://api.fabric.microsoft.com/v1/workspaces/{workspaceId}/items/{sjditemid} com a carga útil acima. Uma resposta do código HTTP 202 indica que o item de definição do trabalho do Spark foi atualizado com sucesso.