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.

  1. 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.

  2. 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.

  3. 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.

  4. Registre Cloud Object Storage en watsonx.data. Para obtener más información, consulte Adición de un par de almacenamiento-catálogo.

  5. 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.py o delta_demo.py y cargue la aplicación Spark en COS, consulte Carga de datos.

  6. 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 formato echo -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.

  7. 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.

  8. 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