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.
-
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.
-
Ordnen Sie den Speicher der nativen Spark-Engine zu. Weitere Informationen finden Sie unter Katalog einer Engine zuordnen.
-
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.
-
Registrieren Sie Cloud Object Storage in watsonx.data. Weitere Informationen finden Sie unter Speicher-/Katalogpaar hinzufügen.
-
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.pyoderdelta_demo.pyund zum Hochladen der Spark-Anwendung in COS finden Sie unter Daten hochladen. -
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 Formatecho -n 'ibmlhapikey_<username>:<user_apikey>' | base64haben. 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. -
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.
-
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