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.
-
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.
-
Associare la memoria al motore Native Spark. Per ulteriori informazioni, vedi Associazione di un catalogo con un motore.
-
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.
-
Registra Cloud Object Storage in watsonx.data. Per ulteriori informazioni, vedi Aggiunta di una coppia archivio - catalogo.
-
In base al catalogo selezionato, salva la seguente applicazione Spark (filePython ) sulla tua macchina locale. Qui,
iceberg_demo.py,hudi_demo.pyodelta_demo.pye carica l'applicazione Spark in COS, vedi Caricamento dei dati. -
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 formatoecho -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. -
Dopo aver inoltrato l'applicazione Spark, si riceve un messaggio di conferma con l'ID applicazione e la versione Spark. Salvarla per riferimento.
-
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