Soumettre une application Spark en utilisant le moteur Spark natif

S'applique à: Moteur d'allumage Gluten accéléré Moteur d'allumage

Cette rubrique fournit la procédure pour soumettre une application Spark en utilisant le moteur Spark natif dans watsonx.data sur IBM Cloud

Prérequis

  • Créer un stockage d'objets : Pour stocker l'application Spark et les résultats associés, créer un seau de stockage. Pour créer Cloud Object Storage et un seau, voir Création d'un seau de stockage. Maintenir un stockage séparé pour l'application et les données. N'enregistrez que les seaux de données avec watsonx.data.

  • Enregistrer le Cloud Object Storage: Enregistrer le Cloud Object Storage bucket dans watsonx.data. Pour enregistrer un seau Cloud Object Storage, voir Adding bucket catalog pair.

    Vous pouvez créer différents Cloud Object Storage buckets pour stocker le code de l'application et la sortie. Enregistrez le seau de données, qui stocke les données d'entrée, et les tables watsonx.data. Il n'est pas nécessaire d'enregistrer le seau de stockage, qui conserve le code de l'application avec watsonx.data.

  • Associer le stockage au moteur Spark. Pour plus d'informations sur l'association avec le moteur Spark, voir Associer un catalogue à un moteur.

Stockages pris en charge

  • Azure Stockage dans un lac de données (ADLS)

    Azure Data Lake Storage (ADLS) Gen1 est obsolète et sera supprimé dans une prochaine version. Vous devez passer à ADLS Gen2 car ADLS Gen1 ne sera plus disponible.

  • Amazon S3

  • Google Cloud Storage (GCS)

  • Cloud Object Storage (COS)

Soumettre une application Spark sans accéder au catalogue watsonx.data

Vous pouvez soumettre une application Spark en exécutant une commande CURL. Effectuez les étapes suivantes pour soumettre une demande Python.

Exécutez la commande curl suivante pour soumettre l'application de comptage de mots.

Exemple de V2 API

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

Exemple de V3 API

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

Paramètres :

  • <crn_instance>: Le CRN de l'instance watsonx.data.
  • <region>: La région où l'instance Spark est provisionnée.
  • <spark_engine_id>: L'identifiant du moteur Spark.
  • <token> le jeton du porteur : Le jeton du porteur. Pour plus d'informations sur la génération du jeton, voir Génération d'un jeton de porteur.

Soumission d'une application Spark en accédant au catalogue watsonx.data

Pour accéder aux données d'un catalogue associé au moteur Spark et effectuer quelques opérations de base sur ce catalogue, procédez comme suit :

Exécutez la commande curl suivante :

Exemple de V2 API

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

Exemple de V3 API

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

Valeurs des paramètres :

  • <encoded-api-key>: La valeur doit être au format echo -n"ibmlhapikey_<user_id>:<user’s api key>" | base64. Ici, <user_id> est l'ID IBM Cloud de l'utilisateur dont la clé api est utilisée pour accéder au seau de données. La <IAM_APIKEY> ici est la clé API de l'utilisateur qui accède au seau du magasin d'objets. Pour générer une clé API, connectez-vous à la console watsonx.data et naviguez vers Profil > Profil et paramètres > Clés API et générez une nouvelle clé API.
  • <storage> la valeur dépend du type de stockage choisi. Il doit s'agir d' s3a, pour le stockage d'objets dans le cloud ( Amazon S3, ou Cloud object Storage, COS), d' abfss, pour ADLS, et d' gs, pour le stockage GCS.
  • <application_bucket_name>: Le nom de l'espace de stockage d'objets contenant le code de votre application. Vous devez transmettre les informations d'identification de ce stockage s'il n'est pas enregistré auprès de watsonx.data.

Exemple d'application Python pour les opérations du catalogue Iceberg

Voici l'exemple d'application Python permettant d'effectuer des opérations de base sur les données stockées dans un catalogue Iceberg :

Comme le schéma spark n'est pas sélectionné dans le payload, mais dans l'application spark, si vous recevez l'erreur [SCHEMA_NOT_FOUND] The schema \<schema_name> en essayant de vous connecter au catalogue iceberg, assurez-vous que le catalogue et le nom du schéma sont corrects dans l'application spark. Assurez-vous également que le catalogue sous lequel vous recherchez le schéma est associé au moteur 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()

Traitement des incidents

  • Si vous recevez l'erreur [REQUIRES_SINGLE_PART_NAMESPACE] spark_catalog requires a single-part namespace, but got \iceberg_data.results alors que vous essayez de créer une table avec le nom de trois parties likeiceberg_data.results.resultstable, assurez-vous que le catalogue est associé au moteur d'étincelles. Si Spark n'est pas configuré avec un catalogue Iceberg, il n'aura pas les paramètres nécessaires pour reconnaître ou utiliser le format de nom en trois parties pour le catalogue iceberg_data.

  • Comme le schéma spark n'est pas sélectionné dans le payload, mais dans l'application spark, si vous recevez l'erreur [SCHEMA_NOT_FOUND] The schema \<schema_name> en essayant de vous connecter au catalogue iceberg, assurez-vous que le catalogue et le nom du schéma sont corrects dans l'application spark. Assurez-vous également que le catalogue sous lequel vous recherchez le schéma est associé au moteur Spark.

API connexes

Pour plus d'informations sur l'API correspondante, voir