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 formatoecho -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 sers3apara Amazon S3 ou Cloud object Storage (COS),abfsspara ADLS egspara 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.resultsao 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