Utilisation de différents formats de table

S'applique à: Moteur d'allumage Gluten accéléré Moteur d'allumage

Cette rubrique décrit la procédure d'exécution d'une application Spark qui ingère des données dans différents formats de table tels que Apache Hudi, Apache Iceberg ou Delta Lake catalog.

  1. Créer un stockage avec le catalogue requis (le catalogue peut être Apache Hudi, Apache Iceberg ou Delta Lake ) pour stocker les données utilisées dans l'application Spark. Pour créer du stockage, voir Ajout d'une paire stock-catalogue.

  2. Associez le stockage au moteur Native Spark. Pour plus d'informations, voir Association d'un catalogue à un moteur.

  3. Créez Cloud Object Storage (COS) pour stocker l'application Spark. Pour créer Cloud Object Storage et un compartiment, voir Création d'un compartiment de stockage.

  4. Enregistrez Cloud Object Storage dans watsonx.data. Pour plus d'informations, voir Ajout d'une paire stockage-catalogue.

  5. En fonction du catalogue que vous sélectionnez, sauvegardez l'application Spark suivante (fichierPython ) sur votre machine locale. Ici, iceberg_demo.py, hudi_demo.py ou delta_demo.py et téléchargez l'application Spark dans COS, voir Téléchargement de données.

  6. Pour soumettre l'application Spark avec des données résidant dans Cloud Object Storage, spécifiez les valeurs de paramètre et exécutez la commande curl à partir du tableau suivant.

    • Apache Iceberg

    Le fichier d'exemple démontre les fonctionnalités suivantes :

    • Accès aux tableaux à partir de watsonx.data

    • Intégration des données dans watsonx.data

    • Modification du schéma dans watsonx.data

    Effectuer des activités de maintenance des tables dans watsonx.data

    Vous devez insérer les données dans le compartiment COS. Pour plus d'informations, voir Insérer des données échantillons dans le seau COS.

    Application Python: fichier Python Iceberg

    Commande Curl pour soumettre l'application Python:

    Exemple de V2 API

    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"  }
            }'
    

    Exemple de V3 API

    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"  }
            }'
    

    Valeurs des paramètres :

    • <wxd_host_name>: Le nom d'hôte de votre instance watsonx.data Cloud.

    • <instance_id>: L'ID de l'instance à partir de l' URL l'instance watsonx.data Par exemple, 1609968977179454.

    • <spark_engine_id>: l'identifiant du moteur Spark natif.

    • <token>: Le jeton du porteur. Pour plus d'informations sur la génération du jeton, voir Jeton IAM.

    • <user-authentication-string>: La valeur doit être une chaîne encodée en base 64 de l'identifiant de l'utilisateur et de la clé API. Pour plus d'informations sur le format, voir la note.

    • Apache Hudi

    L'application Python Spark présente les fonctionnalités suivantes :

    • Il crée une base de données dans le catalogue Apache Hudi (que vous avez créé pour stocker les données). Ici, " <database_name>.

    • Il crée une table dans la base de données " <database_name>, à savoir " <table_name>.

    • Il insère des données dans le " <table_name> et effectue une opération de requête SELECT.

    • Il supprime la table et le schéma après utilisation.

    Application 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()
    
    

    Valeurs des paramètres :

    • <database_name>: Indiquez le nom de la base de données que vous souhaitez créer.
    • <table_name>: Indiquez le nom de la table que vous souhaitez créer.
    • <data_storage_name>: Indiquez le nom du stockage Apache Hudi que vous avez créé.

    Commande Curl pour soumettre l'application Python

    Exemple de V2 API

    
    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"    }}
    

    Exemple de V3 API

    
    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"    }}
    

    Valeurs des paramètres :

    • <wxd_host_name>: Le nom d'hôte de votre instance watsonx.data Cloud.

    • <instance_id>: L'ID de l'instance à partir de l' URL l'instance watsonx.data Par exemple, 1609968977179454.

    • <spark_engine_id>: l'identifiant du moteur Spark natif.

    • <token>: Le jeton du porteur. Pour plus d'informations sur la génération du jeton, voir Jeton IAM.

    • <user-authentication-string>: La valeur doit être une chaîne encodée en base 64 de l'identifiant de l'utilisateur et de la clé API. Pour plus d'informations sur le format, voir la note.

    • Delta Lake

    L'application Python Spark présente les fonctionnalités suivantes :

    • Il crée une base de données dans le catalogue Delta Lake (que vous avez créé pour stocker les données). Ici, " <database_name>.

    • Il crée une table dans la base de données " <database_name>, à savoir " <table_name>.

    • Il insère des données dans le " <table_name> et effectue une opération de requête SELECT.

    • Il supprime la table et le schéma après utilisation.

    Application 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()
    
    

    Valeurs des paramètres :

    • <database_name>: Indiquez le nom de la base de données que vous souhaitez créer.
    • <table_name>: Indiquez le nom de la table que vous souhaitez créer.
    • <data_storage_name>: Indiquez le nom du stockage Apache Hudi que vous avez créé.

    Commande Curl pour soumettre l'application Python

    Exemple de V2 API

    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"        }    }
    

    Exemple de V3 API

    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"        }    }
    

    Valeurs des paramètres

    • <wxd_host_name>: Le nom d'hôte de votre instance watsonx.data Cloud.
    • <instance_id> iD de l'instance : L'ID de l'instance à partir de l' URL l'instance du cluster watsonx.data Par exemple, 1609968977179454.
    • <spark_engine_id>: l'identifiant du moteur Spark natif.
    • <token> le jeton du porteur : Le jeton du porteur. Pour plus d'informations sur la génération du jeton, voir Jeton IAM.
    • <user-authentication-string>: La valeur doit être une chaîne encodée en base 64 de l'identifiant de l'utilisateur et de la clé API. Pour plus d'informations sur le format, voir la remarque suivante.

    La valeur de <user-authentication-string> doit être au format echo -n 'ibmlhapikey_<username>:<user_apikey>' | base64. Ici, <user_id> est l'ID IBM Cloud de l'utilisateur dont la clé d'API est utilisée pour accéder au compartiment de données. <IAM_APIKEY> est la clé d'API de l'utilisateur qui accède au compartiment de magasin d'objets. Pour générer une clé d'API, connectez-vous à la console watsonx.data et accédez à Profil > Profil et paramètres > Clés d'API et générez une nouvelle clé d'API. Si vous générez une nouvelle clé d'API, votre ancienne clé d'API devient non valide.

  7. Une fois que vous avez soumis l'application Spark, vous recevez un message de confirmation avec l'ID de l'application et la version de Spark. Sauvegardez-le pour référence.

  8. Connectez-vous au cluster watsonx.data et accédez à la page des détails du moteur. Dans l'onglet Applications, utilisez l'ID application pour répertorier l'application et suivre les étapes. Pour plus d'informations, voir Affichage et gestion des applications.

API connexes

Pour plus d'informations sur l'API correspondante, voir