Presentazione dell'applicazione Spark utilizzando il motore Spark nativo
Si applica a: Motore a scintilla Glutine accelerato Motore a scintilla
Questo argomento fornisce la procedura per inviare un'applicazione Spark utilizzando il motore Spark nativo in watsonx.data su IBM Cloud.
Prerequisiti
-
Creare un object storage: per memorizzare l'applicazione Spark e il relativo output, creare un bucket di storage. Per creare Cloud Object Storage e un bucket, vedere Creazione di un bucket di archiviazione. Mantenere uno storage separato per l'applicazione e i dati. Registra solo i secchi di dati con watsonx.data.
-
Registrare il Cloud Object Storage: registrare il bucket Cloud Object Storage in watsonx.data. Per registrare il bucket Cloud Object Storage, vedere Aggiungi una coppia di cataloghi di bucket.
È possibile creare diversi Cloud Object Storage bucket per memorizzare il codice dell'applicazione e l'output. Registrare il bucket dei dati, che memorizza i dati di input, e le tabelle watsonx.data. Non è necessario registrare il bucket di archiviazione, che mantiene il codice dell'applicazione con watsonx.data.
-
Associare il magazzino al motore Spark. Per informazioni su come associarsi al motore Spark, vedere Associare un catalogo a un motore.
Archivi supportati
-
Azure Stoccaggio dei laghi di dati (ADLS)
Azure Data Lake Storage (ADLS) Gen1 è deprecato e sarà rimosso in una prossima versione. È necessario passare all'ADLS Gen2 poiché l'ADLS Gen1 non sarà più disponibile.
-
Amazon S3
-
Google Cloud Storage (GCS)
-
COS (Cloud Object Storage)
Invio di un'applicazione Spark senza accesso al catalogo watsonx.data
È possibile inviare un'applicazione Spark eseguendo un comando CURL. Completate i seguenti passaggi per presentare una domanda su Python.
Eseguire il seguente comando curl per inviare l'applicazione del conteggio delle parole.
Esempio di API V2
curl --request POST --url https://<region>.lakehouse.cloud.ibm.com/lakehouse/api/v2/spark_engines/<spark_engine_id>/applications --header 'Authorization: Bearer <token>' --header 'Content-Type: application/json' --header 'AuthInstanceID: <crn_instance>' --data '{
"application_details": {
"application": "/opt/ibm/spark/examples/src/main/python/wordcount.py",
"arguments": [
"/opt/ibm/spark/examples/src/main/resources/people.txt"
]
}
}'
Esempio di API V3
curl --request POST --url https://<region>.lakehouse.cloud.ibm.com/lakehouse/api/v3/spark_engines/<spark_engine_id>/applications --header 'Authorization: Bearer <token>' --header 'Content-Type: application/json' --header 'AuthInstanceID: <crn_instance>' --data '{
"application_details": {
"application": "/opt/ibm/spark/examples/src/main/python/wordcount.py",
"arguments": [
"/opt/ibm/spark/examples/src/main/resources/people.txt"
]
}
}'
Parametri:
<crn_instance>: Il CRN dell'istanza watsonx.data.<region>: L'area in cui l'istanza di Spark viene fornita.<spark_engine_id>: l'ID del motore Spark.<token>il token del portatore. Per ulteriori informazioni sulla generazione del token, vedere Generazione di un token al portatore.
Presentazione di un'applicazione Spark accedendo al catalogo watsonx.data
Per accedere ai dati di un catalogo associato al motore Spark ed eseguire alcune operazioni di base su tale catalogo, procedere come segue:
Esegui il seguente comando curl:
Esempio di API V2
curl --request POST --url https://<region>.lakehouse.cloud.ibm.com/lakehouse/api/v2/spark_engines/<spark_engine_id>/applications --header 'Authorization: Bearer <token>' --header 'Content-Type: application/json' --header 'AuthInstanceID: <crn_instance>' --data '{
"application_details": {
"conf": {
"spark.hadoop.wxd.apiKey": "Basic <encoded-api-key>"
"spark.eventLog.logBlockUpdates.enabled":"true"
},
"application": "<storage>://<application-bucket-name>/iceberg.py"
}
}'
Esempio di API V3
curl --request POST --url https://<region>.lakehouse.cloud.ibm.com/lakehouse/api/v3/spark_engines/<spark_engine_id>/applications --header 'Authorization: Bearer <token>' --header 'Content-Type: application/json' --header 'AuthInstanceID: <crn_instance>' --data '{
"application_details": {
"conf": {
"spark.hadoop.wxd.apiKey": "Basic <encoded-api-key>"
"spark.eventLog.logBlockUpdates.enabled":"true"
},
"application": "<storage>://<application-bucket-name>/iceberg.py"
}
}'
Valori di parametro:
<encoded-api-key>: il valore deve essere nel formatoecho -n"ibmlhapikey_<user_id>:<user’s api key>" | base64. Qui, <user_id> è l'ID IBM Cloud dell'utente la cui chiave api è usata per accedere al bucket dei dati. Il<IAM_APIKEY>qui è la chiave API dell'utente che accede al bucket Object store. Per generare una chiave API, accedere alla console watsonx.data e navigare in Profilo > Profilo e impostazioni > Chiavi API e generare una nuova chiave API.<storage>il valore dipende dal tipo di memorizzazione scelto. Deve esseres3aper l' Amazon S3, o Cloud object Storage (COS),abfssper ADLS, egsper GCS storage.<application_bucket_name>: il nome dell'archivio oggetti contenente il codice dell'applicazione. È necessario passare le credenziali di questo archivio se non è registrato con watsonx.data.
Esempio di applicazione Python per il catalogo Iceberg Operazioni
Di seguito è riportato un esempio di applicazione Python per eseguire operazioni di base sui dati memorizzati in un catalogo Iceberg:
Poiché lo schema spark non è selezionato nel payload, ma è selezionato all'interno dell'applicazione spark, se si riceve l'errore [SCHEMA_NOT_FOUND] The schema \<schema_name> durante il tentativo di connessione al catalogo
iceberg, assicurarsi che il catalogo e il nome dello schema corretti siano forniti nell'applicazione spark. Assicuratevi anche che il catalogo in cui cercate lo schema sia associato al motore Spark.
from pyspark.sql import SparkSession
import os
from datetime import datetime
def init_spark():
spark = SparkSession.builder.appName("lh-hms-cloud").enableHiveSupport().getOrCreate()
return spark
def create_database(spark,bucket_name,catalog):
spark.sql(f"create database if not exists {catalog}.<db_name> LOCATION 's3a://{bucket_name}/'")
def list_databases(spark,catalog):
spark.sql(f"show databases from {catalog}").show()
def basic_iceberg_table_operations(spark,catalog):
spark.sql(f"create table if not exists {catalog}.<db_name>.<table_name>(id INTEGER, name
VARCHAR(10), age INTEGER, salary DECIMAL(10, 2)) using iceberg").show()
spark.sql(f"insert into {catalog}.<db_name>.<table_name>
values(1,'Alan',23,3400.00),(2,'Ben',30,5500.00),(3,'Chen',35,6500.00)")
spark.sql(f"select * from {catalog}.<db_name>.<table_name>").show()
def clean_database(spark,catalog):
spark.sql(f'drop table if exists {catalog}.<db_name>.<table_name> purge')
spark.sql(f'drop database if exists {catalog}.<db_name> cascade')
def main():
try:
spark = init_spark()
create_database(spark,"<wxd-data-bucket-name>","<wxd-data-bucket-catalog-name>")
list_databases(spark,"<wxd-data-bucket-catalog-name>")
basic_iceberg_table_operations(spark,"<wxd-data-bucket-catalog-name>")
finally:
clean_database(spark,"<wxd-data-bucket-catalog-name>")
spark.stop()
if __name__ == '__main__':
main()
Risoluzione dei problemi
-
Se si riceve l'errore
[REQUIRES_SINGLE_PART_NAMESPACE] spark_catalog requires a single-part namespace, but got \iceberg_data.resultsmentre si cerca di creare una tabella con il nome di tre parti likeiceberg_data.results.resultstable, assicurarsi che il catalogo sia associato al motore spark. Se Spark non è configurato con un catalogo Iceberg, non avrà le impostazioni necessarie per riconoscere o utilizzare il formato dei nomi in tre parti per il catalogo iceberg_data. -
Poiché lo schema spark non è selezionato nel payload, ma è selezionato all'interno dell'applicazione spark, se si riceve l'errore
[SCHEMA_NOT_FOUND] The schema \<schema_name>durante il tentativo di connessione al catalogo iceberg, assicurarsi che il catalogo e il nome dello schema corretti siano forniti nell'applicazione spark. Assicuratevi anche che il catalogo in cui cercate lo schema sia associato al motore Spark.
API correlate
Per informazioni sulle API correlate, vedere