Envío de la aplicación Spark utilizando el motor Spark nativo
Se aplica a: motor Spark Motor Spark acelerado por gluten
Este tema proporciona el procedimiento para enviar una aplicación Spark mediante el motor Spark nativo en watsonx.data en IBM Cloud.
Requisitos previos
-
Crear un almacenamiento de objetos : Para almacenar la aplicación Spark y la salida relacionada, cree un bucket de almacenamiento. Para crear Cloud Object Storage y un bucket, consulta Crear un bucket de almacenamiento. Mantenga un almacenamiento separado para la aplicación y los datos. Registrar sólo cubos de datos con watsonx.data.
-
Registrar el Cloud Object Storage: Registrar Cloud Object Storage bucket en watsonx.data. Para registrar Cloud Object Storage bucket, consulte Addding bucket catalog pair.
Puedes crear diferentes Cloud Object Storage buckets para almacenar el código de la aplicación y la salida. Registra el bucket de datos, que almacena los datos de entrada, y las tablas watsonx.data. No es necesario registrar el bucket de almacenamiento, que mantiene el código de la aplicación con watsonx.data.
-
Asociar el almacenamiento con el motor Spark. Para obtener información sobre cómo asociarse al motor Spark, consulte Asociar un catálogo a un motor.
Almacenes compatibles
-
Azure Almacenamiento en lagos de datos (ADLS)
Azure Data Lake Storage (ADLS) Gen1 está obsoleto y se eliminará en una próxima versión. Debe pasar a ADLS Gen2, ya que ADLS Gen1 dejará de estar disponible.
-
Amazon S3
-
Google Cloud Storage (GCS)
-
Cloud Object Storage (COS)
Envío de una aplicación Spark sin acceder al catálogo watsonx.data
Puede enviar una aplicación Spark ejecutando un comando CURL. Complete los siguientes pasos para enviar una solicitud Python.
Ejecute el siguiente comando curl para enviar la aplicación de recuento de palabras.
Ejemplo de API V2
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"
]
}
}'
Ejemplo de API V3
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>: El CRN de la instancia watsonx.data.<region>: La región en la que se aprovisiona la instancia de Spark.<spark_engine_id>: El ID de motor del motor Spark.<token>: La ficha de portador. Para obtener más información sobre la generación del token, consulte Generación de un token de portador.
Envío de una aplicación Spark accediendo al catálogo watsonx.data
Para acceder a los datos de un catálogo que está asociado con el motor Spark y realizar algunas operaciones básicas en ese catálogo, haga lo siguiente:
Ejecute el siguiente mandato curl:
Ejemplo de API V2
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"
}
}'
Ejemplo de API V3
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ámetros:
<encoded-api-key>: El valor debe tener el formatoecho -n"ibmlhapikey_<user_id>:<user’s api key>" | base64. Aquí, <user_id> es el IBM Cloud del usuario cuya clave api se utiliza para acceder al bucket de datos. El<IAM_APIKEY>aquí es la clave de API del usuario que accede al cubo del almacén de objetos. Para generar una clave de API, inicia sesión en la consola watsonx.data y navega hasta Perfil > Perfil y configuración > Claves de API y genera una nueva clave de API.<storage>el valor depende del tipo de almacenamiento que elijas. Debe sers3apara el almacenamiento en la nube ( Amazon S3, COS),abfsspara ADLS ygspara el almacenamiento en GCS.<application_bucket_name>: El nombre del almacenamiento de objetos que contiene el código de su aplicación. Debes pasar las credenciales de este almacenamiento si no está registrado en watsonx.data.
Ejemplo de aplicación Python para el catálogo Iceberg Operaciones
A continuación se muestra el ejemplo de aplicación Python para realizar operaciones básicas sobre los datos almacenados en un catálogo Iceberg:
Como el esquema de Spark no se selecciona en la carga útil, sino en la aplicación de Spark, si recibe el error [SCHEMA_NOT_FOUND] The schema \<schema_name> al intentar conectarse al catálogo de Iceberg, asegúrese de que en
la aplicación de Spark se proporcionan el catálogo y el nombre del esquema correctos. Asegúrese también de que el catálogo en el que busca el esquema está asociado al motor 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()
Resolución de problemas
-
Si recibe el error,
[REQUIRES_SINGLE_PART_NAMESPACE] spark_catalog requires a single-part namespace, but got \iceberg_data.resultsal intentar crear una tabla con el nombre de tres partes likeiceberg_data.results.resultstable, asegúrese de que el catálogo está asociado con el motor spark. Si Spark no está configurado con un catálogo Iceberg, no tendrá los ajustes necesarios para reconocer o utilizar el formato de nombre de tres partes para el catálogo iceberg_data. -
Como el esquema de Spark no se selecciona en la carga útil, sino en la aplicación de Spark, si recibe el error
[SCHEMA_NOT_FOUND] The schema \<schema_name>al intentar conectarse al catálogo de Iceberg, asegúrese de que en la aplicación de Spark se proporcionan el catálogo y el nombre del esquema correctos. Asegúrese también de que el catálogo en el que busca el esquema está asociado al motor Spark.
API relacionadas
Para obtener información sobre las API relacionadas, consulte