Mit verschiedenen Tabellenformaten arbeiten

Gilt für: Funkenmotor Gluten beschleunigter Funkenmotor

Das Thema beschreibt das Verfahren zum Ausführen einer Spark-Anwendung, die Daten in verschiedene Tabellenformate wie Apache Hudi, Apache Iceberg oder Delta Lake catalog einspeist.

  1. Erstellen Sie einen Speicher mit dem erforderlichen Katalog (der Katalog kann Apache Hudi, Apache Iceberg oder Delta Lake sein), um die in der Spark-Anwendung verwendeten Daten zu speichern. Informationen zum Erstellen von Speicher finden Sie unter Speicher-Katalog-Paar hinzufügen.

  2. Ordnen Sie den Speicher der nativen Spark-Engine zu. Weitere Informationen finden Sie unter Katalog einer Engine zuordnen.

  3. Erstellen Sie Cloud Object Storage (COS), um die Spark-Anwendung zu speichern. Informationen zum Erstellen von Cloud Object Storage und eines Buckets finden Sie unter Speicherbucket erstellen.

  4. Registrieren Sie Cloud Object Storage in watsonx.data. Weitere Informationen finden Sie unter Speicher-/Katalogpaar hinzufügen.

  5. Speichern Sie auf der Basis des ausgewählten Katalogs die folgende Spark-Anwendung (Python-Datei) auf Ihrer lokalen Maschine. Informationen zu iceberg_demo.py, hudi_demo.py oder delta_demo.py und zum Hochladen der Spark-Anwendung in COS finden Sie unter Daten hochladen.

  6. Um die Spark-Anwendung mit Daten zu übergeben, die sich in Cloud Object Storagebefinden, geben Sie die Parameterwerte an und führen Sie den curl-Befehl aus der folgenden Tabelle aus.

    • Apache Iceberg

    Die Beispieldatei demonstriert die folgenden Funktionalitäten:

    • Zugriff auf Tabellen aus watsonx.data

    • Einspeisung von Daten in watsonx.data

    • Ändern des Schemas in watsonx.data

    Durchführung von Tabellenpflegeaktivitäten in watsonx.data.

    Sie müssen die Daten in das COS-Bucket einfügen. Weitere Informationen finden Sie unter Einfügen von Beispieldaten in den COS-Eimer.

    Python: Iceberg Python

    Curl zur Übermittlung der Python:

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

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

    Parameterwerte :

    • <wxd_host_name>: Der Hostname Ihrer watsonx.data Cloud-Instanz.

    • <instance_id>: Die Instanz-ID aus der watsonx.data URL. Zum Beispiel: 1609968977179454.

    • <spark_engine_id>: Die Engine-ID der nativen Spark-Engine.

    • <token>: Der Überbringer des Tokens. Weitere Informationen zur Erstellung des Tokens finden Sie unter IAM-Token.

    • <user-authentication-string>: Der Wert muss eine mit Base 64 kodierte Zeichenfolge aus Benutzer-ID und API-Schlüssel sein. Weitere Informationen über das Format finden Sie in der Anmerkung.

    • Apache Hudi

    Die Python Spark-Anwendung demonstriert die folgenden Funktionen:

    • Es erstellt eine Datenbank innerhalb des Apache Hudi Katalogs (den Sie zum Speichern von Daten erstellt haben). Hier, ' <database_name>.

    • Es wird eine Tabelle innerhalb der Datenbank " <database_name> erstellt, nämlich " <table_name>.

    • Er fügt Daten in den ' <table_name> ein und führt eine SELECT-Abfrage durch.

    • Es löscht die Tabelle und das Schema nach der Verwendung.

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

    Parameterwerte:

    • <database_name>: Geben Sie den Namen der Datenbank an, die Sie erstellen möchten.
    • <table_name>: Geben Sie den Namen der Tabelle an, die Sie erstellen möchten.
    • <data_storage_name>: Geben Sie den Namen des Apache Hudi Speichers an, den Sie erstellt haben.

    Curl-Befehl zum Übergeben der Python-Anwendung

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

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

    Parameterwerte:

    • <wxd_host_name>: Der Hostname Ihrer watsonx.data Cloud-Instanz.

    • <instance_id>: Die Instanz-ID aus der watsonx.data URL. Zum Beispiel: 1609968977179454.

    • <spark_engine_id>: Die Engine-ID der nativen Spark-Engine.

    • <token>: Der Überbringer des Tokens. Weitere Informationen zur Erstellung des Tokens finden Sie unter IAM-Token.

    • <user-authentication-string>: Der Wert muss eine mit Base 64 kodierte Zeichenfolge aus Benutzer-ID und API-Schlüssel sein. Weitere Informationen über das Format finden Sie in der Anmerkung.

    • Delta Lake

    Die Python Spark-Anwendung demonstriert die folgenden Funktionen:

    • Es erstellt eine Datenbank innerhalb des Delta Lake Katalogs (den Sie zum Speichern von Daten erstellt haben). Hier, ' <database_name>.

    • Es wird eine Tabelle innerhalb der Datenbank " <database_name> erstellt, nämlich " <table_name>.

    • Er fügt Daten in den ' <table_name> ein und führt eine SELECT-Abfrage durch.

    • Es löscht die Tabelle und das Schema nach der Verwendung.

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

    Parameterwerte:

    • <database_name>: Geben Sie den Namen der Datenbank an, die Sie erstellen möchten.
    • <table_name>: Geben Sie den Namen der Tabelle an, die Sie erstellen möchten.
    • <data_storage_name>: Geben Sie den Namen des Apache Hudi Speichers an, den Sie erstellt haben.

    Curl-Befehl zum Übergeben der Python-Anwendung

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

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

    Parameterwerte

    • <wxd_host_name>: Der Hostname Ihrer watsonx.data Cloud-Instanz.
    • <instance_id> instanz-ID: Die Instanz-ID aus der URL der watsonx.data. Zum Beispiel: 1609968977179454.
    • <spark_engine_id>: Die Engine-ID der nativen Spark-Engine.
    • <token>: Der Überbringer des Tokens. Weitere Informationen zur Erstellung des Tokens finden Sie unter IAM-Token.
    • <user-authentication-string>: Der Wert muss eine mit Base 64 kodierte Zeichenfolge aus Benutzer-ID und API-Schlüssel sein. Weitere Informationen zum Format finden Sie im folgenden Hinweis.

    Der Wert von <user-authentication-string> muss das Format echo -n 'ibmlhapikey_<username>:<user_apikey>' | base64 haben. Hier ist <user_id> die IBM Cloud-ID des Benutzers, dessen API-Schlüssel für den Zugriff auf das Datenbucket verwendet wird. Die <IAM_APIKEY> hier ist der API-Schlüssel des Benutzers, der auf das Objektspeicherbucket zugreift. Um einen API-Schlüssel zu generieren, melden Sie sich an der watsonx.data-Konsole an und navigieren Sie zu Profil > Profil und Einstellungen > API-Schlüssel und generieren Sie einen neuen API-Schlüssel. Wenn Sie einen neuen API-Schlüssel generieren, wird Ihr alter API-Schlüssel ungültig.

  7. Nachdem Sie die Spark-Anwendung übergeben haben, erhalten Sie eine Bestätigungsnachricht mit der Anwendungs-ID und der Spark-Version. Speichern Sie sie zu Referenzzwecken.

  8. Melden Sie sich beim Cluster watsonx.data an und rufen Sie die Detailseite der Engine auf. Verwenden Sie auf der Registerkarte "Anwendungen" die Anwendungs-ID, um die Anwendung aufzulisten und die Phasen zu verfolgen. Weitere Informationen finden Sie unter Anwendungen anzeigen und verwalten.

Verwandte APIs

Für Informationen über die zugehörige API siehe