Envio do aplicativo Spark usando o mecanismo Spark nativo

Aplica-se a: Motor de faísca Motor de faísca acelerado com glúten

Este tópico fornece o procedimento para enviar um aplicativo Spark usando o mecanismo Spark nativo em watsonx.data no IBM Cloud.

Pré-requisitos

  • Criar um armazenamento de objetos: para armazenar o aplicativo Spark e a saída relacionada, crie um bucket de armazenamento. Para criar o Cloud Object Storage e um bucket, consulte Criando um bucket de armazenamento. Mantenha um armazenamento separado para aplicativos e dados. Registre somente os compartimentos de dados com watsonx.data.

  • Registre o Cloud Object Storage: Registre o Cloud Object Storage bucket em watsonx.data. Para registrar o Cloud Object Storage, consulte Adicionar par de catálogos de balde.

    Você pode criar diferentes Cloud Object Storage para armazenar o código do aplicativo e a saída. Registre o bucket de dados, que armazena os dados de entrada, e as tabelas watsonx.data. Não é necessário registrar o bucket de armazenamento, que mantém o código do aplicativo com watsonx.data.

  • Associe o armazenamento ao mecanismo Spark. Para obter informações sobre como se associar ao mecanismo do Spark, consulte Associação de um catálogo a um mecanismo.

Armazenamentos suportados

  • Azure Armazenamento em lago de dados (ADLS)

    Azure O Data Lake Storage (ADLS) Gen1 está obsoleto e será removido em uma versão futura. Você deve fazer a transição para o ADLS Gen2, pois o ADLS Gen1 não estará mais disponível.

  • Amazon S3

  • Google Cloud Storage (GCS)

  • Cloud Object Storage (COS)

Envio de um aplicativo Spark sem acessar o catálogo watsonx.data

Você pode enviar um aplicativo Spark executando um comando CURL. Conclua as etapas a seguir para enviar uma solicitação para Python.

Execute o seguinte comando curl para enviar o aplicativo de contagem de palavras.

Exemplo de V2 API

   curl --request POST --url https://<region>.lakehouse.cloud.ibm.com/lakehouse/api/v2/spark_engines/<spark_engine_id>/applications --header 'Authorization: Bearer <token>' --header 'Content-Type: application/json' --header 'AuthInstanceID: <crn_instance>' --data '{
       "application_details": {
           "application": "/opt/ibm/spark/examples/src/main/python/wordcount.py",
           "arguments": [
               "/opt/ibm/spark/examples/src/main/resources/people.txt"
           ]
       }
   }'

Exemplo de V3 API

   curl --request POST --url https://<region>.lakehouse.cloud.ibm.com/lakehouse/api/v3/spark_engines/<spark_engine_id>/applications --header 'Authorization: Bearer <token>' --header 'Content-Type: application/json' --header 'AuthInstanceID: <crn_instance>' --data '{
       "application_details": {
           "application": "/opt/ibm/spark/examples/src/main/python/wordcount.py",
           "arguments": [
               "/opt/ibm/spark/examples/src/main/resources/people.txt"
           ]
       }
   }'

Parâmetros:

  • <crn_instance> cRN da instância watsonx.data: O CRN da instância.
  • <region>: A região onde a instância do Spark é provisionada.
  • <spark_engine_id>: A ID do mecanismo do Spark.
  • <token> token de portador: O token de portador. Para obter mais informações sobre a geração do token, consulte Geração de um token de portador.

Envio de um aplicativo Spark acessando o catálogo watsonx.data

Para acessar dados de um catálogo associado ao mecanismo Spark e executar algumas operações básicas nesse catálogo, faça o seguinte:

Execute o comando curl a seguir:

Exemplo de V2 API

curl --request POST --url https://<region>.lakehouse.cloud.ibm.com/lakehouse/api/v2/spark_engines/<spark_engine_id>/applications --header 'Authorization: Bearer <token>' --header 'Content-Type: application/json' --header 'AuthInstanceID: <crn_instance>' --data '{
    "application_details": {
        "conf": {
            "spark.hadoop.wxd.apiKey": "Basic <encoded-api-key>"
            "spark.eventLog.logBlockUpdates.enabled":"true"
        },
        "application": "<storage>://<application-bucket-name>/iceberg.py"
    }
}'

Exemplo de V3 API

curl --request POST --url https://<region>.lakehouse.cloud.ibm.com/lakehouse/api/v3/spark_engines/<spark_engine_id>/applications --header 'Authorization: Bearer <token>' --header 'Content-Type: application/json' --header 'AuthInstanceID: <crn_instance>' --data '{
    "application_details": {
        "conf": {
            "spark.hadoop.wxd.apiKey": "Basic <encoded-api-key>"
            "spark.eventLog.logBlockUpdates.enabled":"true"
        },
        "application": "<storage>://<application-bucket-name>/iceberg.py"
    }
}'

Valores de parâmetro:

  • <encoded-api-key>: O valor deve estar no formato echo -n"ibmlhapikey_<user_id>:<user’s api key>" | base64. Aqui, <user_id> é o IBM Cloud ID do usuário cuja chave de API é usada para acessar o bucket de dados. O <IAM_APIKEY> aqui é a chave de API do usuário que está acessando o bucket do armazenamento de objetos. Para gerar uma chave de API, faça login no console watsonx.data e navegue até Profile > Profile and Settings > API Keys e gere uma nova chave de API.
  • <storage> o valor depende do tipo de armazenamento que você escolher. Deve ser s3a para Amazon S3 ou Cloud object Storage (COS), abfss para ADLS e gs para armazenamento GCS.
  • <application_bucket_name>: O nome do armazenamento de objetos que contém o código do aplicativo. Você deve passar as credenciais desse armazenamento se ele não estiver registrado em watsonx.data.

Exemplo de aplicativo Python para operações do catálogo Iceberg

A seguir, o aplicativo de amostra Python para executar operações básicas em dados armazenados em um catálogo Iceberg:

Como o esquema do Spark não é selecionado na carga útil e é selecionado no aplicativo Spark, se você receber o erro [SCHEMA_NOT_FOUND] The schema \<schema_name> ao tentar se conectar ao catálogo do Iceberg, verifique se o catálogo e o nome do esquema corretos foram fornecidos no aplicativo Spark. Certifique-se também de que o catálogo sob o qual você está procurando o esquema esteja associado ao mecanismo Spark.

from pyspark.sql import SparkSession
import os
from datetime import datetime
def init_spark():
    spark = SparkSession.builder.appName("lh-hms-cloud").enableHiveSupport().getOrCreate()
    return spark
def create_database(spark,bucket_name,catalog):
    spark.sql(f"create database if not exists {catalog}.<db_name> LOCATION 's3a://{bucket_name}/'")
def list_databases(spark,catalog):
    spark.sql(f"show databases from {catalog}").show()
def basic_iceberg_table_operations(spark,catalog):
    spark.sql(f"create table if not exists {catalog}.<db_name>.<table_name>(id INTEGER, name
    VARCHAR(10), age INTEGER, salary DECIMAL(10, 2)) using iceberg").show()
    spark.sql(f"insert into {catalog}.<db_name>.<table_name>
    values(1,'Alan',23,3400.00),(2,'Ben',30,5500.00),(3,'Chen',35,6500.00)")
    spark.sql(f"select * from {catalog}.<db_name>.<table_name>").show()
def clean_database(spark,catalog):
    spark.sql(f'drop table if exists {catalog}.<db_name>.<table_name> purge')
    spark.sql(f'drop database if exists {catalog}.<db_name> cascade')
def main():
    try:
        spark = init_spark()
        create_database(spark,"<wxd-data-bucket-name>","<wxd-data-bucket-catalog-name>")
        list_databases(spark,"<wxd-data-bucket-catalog-name>")
        basic_iceberg_table_operations(spark,"<wxd-data-bucket-catalog-name>")
    finally:
        clean_database(spark,"<wxd-data-bucket-catalog-name>")
        spark.stop()
if __name__ == '__main__':
    main()

Resolução de problemas

  • Se você receber o erro [REQUIRES_SINGLE_PART_NAMESPACE] spark_catalog requires a single-part namespace, but got \iceberg_data.results ao tentar criar uma tabela com o nome de três partes likeiceberg_data.results.resultstable, verifique se o catálogo está associado ao mecanismo de ignição. Se o Spark não estiver configurado com um catálogo Iceberg, ele não terá as configurações necessárias para reconhecer ou usar o formato de nome de três partes para o catálogo iceberg_data.

  • Como o esquema do Spark não é selecionado na carga útil e é selecionado no aplicativo Spark, se você receber o erro [SCHEMA_NOT_FOUND] The schema \<schema_name> ao tentar se conectar ao catálogo do Iceberg, verifique se o catálogo e o nome do esquema corretos foram fornecidos no aplicativo Spark. Certifique-se também de que o catálogo sob o qual você está procurando o esquema esteja associado ao mecanismo Spark.

APIs relacionadas

Para obter informações sobre a API relacionada, consulte