Utilizzo di diversi formati di tabella

Si applica a: Motore a scintilla Glutine accelerato Motore a scintilla

L'argomento descrive la procedura per eseguire un'applicazione Spark che ingerisce i dati in diversi formati di tabella come Apache Hudi, Apache Iceberg o Delta Lake catalogo.

  1. Creare un archivio con il catalogo richiesto (il catalogo può essere Apache Hudi, Apache Iceberg o Delta Lake ) per memorizzare i dati utilizzati nell'applicazione Spark. Per creare l'archiviazione, vedi Aggiunta di una coppia archivio - catalogo.

  2. Associare la memoria al motore Native Spark. Per ulteriori informazioni, vedi Associazione di un catalogo con un motore.

  3. Crea Cloud Object Storage (COS) per memorizzare l'applicazione Spark. Per creare Cloud Object Storage e un bucket, consulta Creazione di un bucket di storage.

  4. Registra Cloud Object Storage in watsonx.data. Per ulteriori informazioni, vedi Aggiunta di una coppia archivio - catalogo.

  5. In base al catalogo selezionato, salva la seguente applicazione Spark (filePython ) sulla tua macchina locale. Qui, iceberg_demo.py, hudi_demo.py o delta_demo.py e carica l'applicazione Spark in COS, vedi Caricamento dei dati.

  6. Per inoltrare l'applicazione Spark con i dati che risiedono in Cloud Object Storage, specificare i valori dei parametri ed eseguire il comando curl dalla tabella seguente.

    • Apache Iceberg

    Il file di esempio dimostra le seguenti funzionalità:

    • Accesso alle tabelle da watsonx.data

    • Immissione dei dati in watsonx.data

    • Modifica dello schema in watsonx.data

    Esecuzione delle attività di manutenzione delle tabelle in watsonx.data.

    Devi inserire i dati nel bucket COS. Per ulteriori informazioni, vedere Inserimento di dati di esempio nel bucket COS.

    Applicazione Python: file Iceberg Python

    Comando Curl per inviare l'applicazione Python:

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

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

    Valori dei parametri :

    • <wxd_host_name>: il nome host dell'istanza di watsonx.data Cloud.

    • <instance_id>: L'ID dell'istanza dall' URL dell'istanza watsonx.data. Ad esempio, 1609968977179454.

    • <spark_engine_id>: l'ID motore del motore Spark nativo.

    • <token>: Il token al portatore. Per ulteriori informazioni sulla generazione del token, vedere Token IAM.

    • <user-authentication-string>: il valore deve essere una stringa codificata in base 64 dell'ID utente e della chiave API. Per ulteriori informazioni sul formato, consultare la nota.

    • Apache Hudi

    L'applicazione Python Spark dimostra le seguenti funzionalità:

    • Crea un database all'interno del catalogo Apache Hudi (creato per memorizzare i dati). Qui, <database_name>.

    • Crea una tabella all'interno del database '<database_name>, ovvero '<table_name>.

    • Inserisce i dati nel '<table_name> ed esegue una query SELECT.

    • Dopo l'uso, la tabella e lo schema vengono eliminati.

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

    Valori di parametro:

    • <database_name>: Specificare il nome del database che si desidera creare.
    • <table_name>: Specificare il nome della tabella che si desidera creare.
    • <data_storage_name>: Specificare il nome dello storage Apache Hudi creato.

    Comando Curl per inoltrare l'applicazione Python

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

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

    Valori di parametro:

    • <wxd_host_name>: il nome host dell'istanza di watsonx.data Cloud.

    • <instance_id>: L'ID dell'istanza dall' URL dell'istanza watsonx.data. Ad esempio, 1609968977179454.

    • <spark_engine_id>: l'ID motore del motore Spark nativo.

    • <token>: Il token al portatore. Per ulteriori informazioni sulla generazione del token, vedere Token IAM.

    • <user-authentication-string>: il valore deve essere una stringa codificata in base 64 dell'ID utente e della chiave API. Per ulteriori informazioni sul formato, consultare la nota.

    • Delta Lake

    L'applicazione Python Spark dimostra le seguenti funzionalità:

    • Crea un database all'interno del catalogo Delta Lake (creato per memorizzare i dati). Qui, <database_name>.

    • Crea una tabella all'interno del database '<database_name>, ovvero '<table_name>.

    • Inserisce i dati nel '<table_name> ed esegue una query SELECT.

    • Dopo l'uso, la tabella e lo schema vengono eliminati.

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

    Valori di parametro:

    • <database_name>: Specificare il nome del database che si desidera creare.
    • <table_name>: Specificare il nome della tabella che si desidera creare.
    • <data_storage_name>: Specificare il nome dello storage Apache Hudi creato.

    Comando Curl per inoltrare l'applicazione Python

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

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

    Valori parametro

    • <wxd_host_name>: il nome host dell'istanza di watsonx.data Cloud.
    • <instance_id> l'ID dell'istanza dall' URL dell'istanza del cluster watsonx.data. Ad esempio, 1609968977179454.
    • <spark_engine_id>: l'ID motore del motore Spark nativo.
    • <token> il token del portatore. Per ulteriori informazioni sulla generazione del token, vedere Token IAM.
    • <user-authentication-string>: il valore deve essere una stringa codificata in base 64 dell'ID utente e della chiave API. Per ulteriori informazioni sul formato, consultare la seguente nota.

    Il valore di <user-authentication-string> deve avere il formato echo -n 'ibmlhapikey_<username>:<user_apikey>' | base64. Qui, <user_id> è l'ID IBM Cloud dell'utente la cui chiave API viene utilizzata per accedere al bucket di dati. <IAM_APIKEY> qui è la chiave API dell'utente che accede al bucket dell'archivio oggetti. Per generare la chiave API, accedi alla console watsonx.data e passa a Profilo> Profilo e impostazioni> Chiavi API e genera una nuova chiave API. Se si genera una nuova chiave API, la vecchia chiave API diventa non valida.

  7. Dopo aver inoltrato l'applicazione Spark, si riceve un messaggio di conferma con l'ID applicazione e la versione Spark. Salvarla per riferimento.

  8. Accedi al cluster watsonx.data, accedi alla pagina dei dettagli del motore. Nella scheda Applicazioni, utilizzare l'ID applicazione per visualizzare l'applicazione e tenere traccia delle fasi. Per ulteriori informazioni, consultare Visualizza e gestisci applicazioni.

API correlate

Per informazioni sulle API correlate, vedere