Cómo trabajar con distintos formatos de tabla
Se aplica a: motor Spark Motor Spark acelerado por gluten
El tema describe el procedimiento para ejecutar una aplicación Spark que ingiere datos en diferentes formatos de tabla como Apache Hudi, Apache Iceberg o Delta Lake catálogo.
-
Cree un almacenamiento con el catálogo requerido (el catálogo puede ser Apache Hudi, Apache Iceberg o Delta Lake ) para almacenar los datos utilizados en la aplicación Spark. Para crear almacenamiento, consulte Adición de un par de almacenamiento-catálogo.
-
Asocie el almacenamiento con el motor de Spark nativo. Para obtener más información, consulte Asociación de un catálogo con un motor.
-
Cree Cloud Object Storage (COS) para almacenar la aplicación Spark. Para crear Cloud Object Storage y un grupo, consulte Creación de un grupo de almacenamiento.
-
Registre Cloud Object Storage en watsonx.data. Para obtener más información, consulte Adición de un par de almacenamiento-catálogo.
-
En función del catálogo que seleccione, guarde la siguiente aplicación Spark (archivoPython ) en la máquina local. Aquí,
iceberg_demo.py,hudi_demo.pyodelta_demo.pyy cargue la aplicación Spark en COS, consulte Carga de datos. -
Para enviar la aplicación Spark con datos que residen en Cloud Object Storage, especifique los valores de parámetro y ejecute el mandato curl desde la tabla siguiente.
- Apache Iceberg
El archivo de ejemplo demuestra las siguientes funcionalidades:
-
Acceso a tablas desde watsonx.data
-
Ingesta de datos en watsonx.data
-
Modificación del esquema en watsonx.data
Realización de actividades de mantenimiento de tablas en watsonx.data.
Debe insertar los datos en el grupo de COS. Para obtener más información, consulte Inserción de datos de muestra en el cubo COS.
Aplicación Python: Archivo Iceberg Python
Comando Curl para enviar la aplicación Python:
Ejemplo de API V2
curl --request POST \ --url https://<wxd_host_name>/lakehouse/api/v2/spark_engines/<spark_engine_id>/applications \ --header 'Authorization: Bearer <token>' \ --header 'Content-Type: application/json' \ --header 'LhInstanceId: <instance_id>' \ --data '{ "application_details": { "conf": { "spark.hadoop.wxd.apiKey":"Basic <user-authentication-string>" }, "application": "s3a://<application-bucket-name>/iceberg.py" } }'Ejemplo de API V3
curl --request POST \ --url https://<wxd_host_name>/lakehouse/api/v3/spark_engines/<spark_engine_id>/applications \ --header 'Authorization: Bearer <token>' \ --header 'Content-Type: application/json' \ --header 'LhInstanceId: <instance_id>' \ --data '{ "application_details": { "conf": { "spark.hadoop.wxd.apiKey":"Basic <user-authentication-string>" }, "application": "s3a://<application-bucket-name>/iceberg.py" } }'Valores de los parámetros :
-
<wxd_host_name>: El nombre de host de tu instancia de watsonx.data Cloud. -
<instance_id>: El ID de instancia de la URL de instancia de watsonx.data. Por ejemplo, 1609968977179454. -
<spark_engine_id>: El ID del motor Spark nativo. -
<token>: El token portador. Para más información sobre la generación del token, véase Token IAM. -
<user-authentication-string>: El valor debe ser una cadena codificada en base 64 de ID de usuario y clave API . Para más información sobre el formato, consulte la nota. -
Apache Hudi
La aplicación Python Spark demuestra la siguiente funcionalidad:
-
Crea una base de datos dentro del catálogo Apache Hudi (que creó para almacenar datos). Aquí, '
<database_name>. -
Crea una tabla dentro de la base de datos '
<database_name>, a saber, '<table_name>. -
Inserta los datos en el '
<table_name>' y realiza la operación de consulta SELECT. -
Elimina la tabla y el esquema después de su uso.
Aplicación Python:
from pyspark.sql import SparkSession def init_spark(): spark = SparkSession.builder.appName("CreateHudiTableInCOS").enableHiveSupport().getOrCreate() return spark def main(): try: spark = init_spark() spark.sql("show databases").show() spark.sql("create database if not exists spark_catalog.<database_name> LOCATION 's3a://<data_storage_name>/'").show() spark.sql("create table if not exists spark_catalog.<database_name>.<table_name> (id bigint, name string, location string) USING HUDI OPTIONS ('primaryKey' 'id', hoodie.write.markers.type= 'direct', hoodie.embed.timeline.server= 'false')").show() spark.sql("insert into <database_name>.<table_name> VALUES (1, 'Sam','Kochi'), (2, 'Tom','Bangalore'), (3, 'Bob','Chennai'), (4, 'Alex','Bangalore')").show() spark.sql("select * from spark_catalog.<database_name>.<table_name>").show() spark.sql("drop table spark_catalog.<database_name>.<table_name>").show() spark.sql("drop schema spark_catalog.<database_name> CASCADE").show() finally: spark.stop() if __name__ == '__main__': main()Valores de parámetros:
<database_name>: Especifique el nombre de la base de datos que desea crear.<table_name>: Especifique el nombre de la tabla que desea crear.<data_storage_name>: Especifique el nombre del almacenamiento Apache Hudi que ha creado.
Mandato Curl para enviar la aplicación Python
Ejemplo de API V2
curl --request POST --url https://<wxd_host_name>/lakehouse/api/v2/spark_engines/<spark_engine_id>/applications --header 'Authorization: Bearer <token>' --header 'Content-Type: application/json' --header 'LhInstanceId: <instance_id>' --data '{ "application_details": { "conf": { "spark.sql.catalog.spark_catalog.type": "hive", "spark.sql.catalog.spark_catalog": "org.apache.spark.sql.hudi.catalog.HoodieCatalog", "spark.hadoop.wxd.apiKey":"Basic <user-authentication-string>" }, "application": "s3a://<data_storage_name>/hudi_demo.py" }}Ejemplo de API V3
curl --request POST --url https://<wxd_host_name>/lakehouse/api/v3/spark_engines/<spark_engine_id>/applications --header 'Authorization: Bearer <token>' --header 'Content-Type: application/json' --header 'LhInstanceId: <instance_id>' --data '{ "application_details": { "conf": { "spark.sql.catalog.spark_catalog.type": "hive", "spark.sql.catalog.spark_catalog": "org.apache.spark.sql.hudi.catalog.HoodieCatalog", "spark.hadoop.wxd.apiKey":"Basic <user-authentication-string>" }, "application": "s3a://<data_storage_name>/hudi_demo.py" }}Valores de parámetros:
-
<wxd_host_name>: El nombre de host de tu instancia de watsonx.data Cloud. -
<instance_id>: El ID de instancia de la URL de instancia de watsonx.data. Por ejemplo, 1609968977179454. -
<spark_engine_id>: El ID del motor Spark nativo. -
<token>: El token portador. Para más información sobre la generación del token, véase Token IAM. -
<user-authentication-string>: El valor debe ser una cadena codificada en base 64 de ID de usuario y clave API . Para más información sobre el formato, consulte la nota. -
Delta Lake
La aplicación Python Spark demuestra la siguiente funcionalidad:
-
Crea una base de datos dentro del catálogo Delta Lake (que creó para almacenar datos). Aquí, '
<database_name>. -
Crea una tabla dentro de la base de datos '
<database_name>, a saber, '<table_name>. -
Inserta los datos en el '
<table_name>' y realiza la operación de consulta SELECT. -
Elimina la tabla y el esquema después de su uso.
Aplicación Python:
from pyspark.sql import SparkSession import os def init_spark(): spark = SparkSession.builder.appName("lh-hms-cloud").enableHiveSupport().getOrCreate() return spark def main(): spark = init_spark() spark.sql("show databases").show() spark.sql("create database if not exists spark_catalog.<database_name> LOCATION 's3a://<data_storage_name>/'").show() spark.sql("create table if not exists spark_catalog.<database_name>.<table_name> (id bigint, name string, location string) USING DELTA").show() spark.sql("insert into spark_catalog.<database_name>.<table_name> VALUES (1, 'Sam','Kochi'), (2, 'Tom','Bangalore'), (3, 'Bob','Chennai'), (4, 'Alex','Bangalore')").show() spark.sql("select * from spark_catalog.<database_name>.<table_name>").show() spark.sql("drop table spark_catalog.<database_name>.<table_name>").show() spark.sql("drop schema spark_catalog.<database_name> CASCADE").show() spark.stop() if __name__ == '__main__': main()Valores de parámetros:
<database_name>: Especifique el nombre de la base de datos que desea crear.<table_name>: Especifique el nombre de la tabla que desea crear.<data_storage_name>: Especifique el nombre del almacenamiento Apache Hudi que ha creado.
Mandato Curl para enviar la aplicación Python
Ejemplo de API V2
curl --request POST --url https://<wxd_host_name>/lakehouse/api/v2/spark_engines/<spark_engine_id>/applications --header 'Authorization: Bearer <token>' --header 'Content-Type: application/json' --header 'LhInstanceId: <instance_id>' --data '{ "application_details": { "conf": { "spark.sql.catalog.spark_catalog" : "org.apache.spark.sql.delta.catalog.DeltaCatalog", "spark.sql.catalog.spark_catalog.type" : "hive", "spark.hadoop.wxd.apiKey":"<user-authentication-string>" }, "application": "s3a://<database_name>/delta_demo.py" } }Ejemplo de API V3
curl --request POST --url https://<wxd_host_name>/lakehouse/api/v3/spark_engines/<spark_engine_id>/applications --header 'Authorization: Bearer <token>' --header 'Content-Type: application/json' --header 'LhInstanceId: <instance_id>' --data '{ "application_details": { "conf": { "spark.sql.catalog.spark_catalog" : "org.apache.spark.sql.delta.catalog.DeltaCatalog", "spark.sql.catalog.spark_catalog.type" : "hive", "spark.hadoop.wxd.apiKey":"<user-authentication-string>" }, "application": "s3a://<database_name>/delta_demo.py" } }Valores de parámetros
<wxd_host_name>: El nombre de host de tu instancia de watsonx.data Cloud.<instance_id>el ID de instancia de la URL de instancia del clúster watsonx.data. Por ejemplo, 1609968977179454.<spark_engine_id>: El ID del motor Spark nativo.<token>: La ficha de portador. Para más información sobre la generación del token, véase Token IAM.<user-authentication-string>: El valor debe ser una cadena codificada en base 64 de ID de usuario y clave API . Para obtener más información sobre el formato, consulte la nota siguiente.
El valor de
<user-authentication-string>debe tener el formatoecho -n 'ibmlhapikey_<username>:<user_apikey>' | base64. Aquí,<user_id>es el ID de IBM Cloud del usuario cuya clave de API se utiliza para acceder al grupo de datos.<IAM_APIKEY>aquí es la clave de API del usuario que accede al grupo del almacén de objetos. Para generar una clave de API, inicie sesión en la consola de watsonx.data y vaya a Perfil > Perfil y valores > Claves de API y genere una nueva clave de API. Si genera una nueva clave de API, la clave de API antigua deja de ser válida. -
Después de enviar la aplicación Spark, recibirá un mensaje de confirmación con el ID de aplicación y la versión de Spark. Guárdelo como referencia.
-
Inicie sesión en el clúster watsonx.data y acceda a la página de detalles del motor. En la ficha Aplicaciones, utilice el ID de aplicación para listar la aplicación y realizar un seguimiento de las etapas. Para obtener más información, consulte Ver y gestionar aplicaciones.
API relacionadas
Para obtener información sobre las API relacionadas, consulte