Einreichen einer Spark-Anwendung mit der nativen Spark-Engine

Gilt für: Funkenmotor Gluten beschleunigter Funkenmotor

Dieses Thema beschreibt das Verfahren zum Einreichen eines Spark-Antrags mithilfe der nativen Spark-Engine in watsonx.data in IBM Cloud.

Voraussetzungen

  • Erstellen eines Objektspeichers: Um die Spark-Anwendung und die zugehörige Ausgabe zu speichern, erstellen Sie einen Speicherbereich. Um Cloud Object Storage und einen Bucket zu erstellen, siehe Erstellen eines Speicher-Buckets. Getrennter Speicher für Anwendung und Daten. Registrieren Sie nur Dateneimer mit watsonx.data.

  • Registrieren Sie den Cloud Object Storage: Registrieren Sie den Cloud Object Storage Bucket in watsonx.data. Um den Cloud Object Storage Bucket zu registrieren, siehe Bucket-Katalogpaar hinzufügen.

    Sie können verschiedene Cloud Object Storage Buckets erstellen, um Anwendungscode und die Ausgabe zu speichern. Registrieren Sie den Datenbereich, in dem die Eingabedaten gespeichert sind, und die Tabellen watsonx.data. Sie müssen den Speicherbereich, der den Anwendungscode verwaltet, nicht mit watsonx.data registrieren.

  • Verknüpfen Sie den Speicher mit der Spark-Engine. Informationen über die Verknüpfung mit der Spark-Engine finden Sie unter Verknüpfung eines Katalogs mit einer Engine.

Unterstützte Speichermedien

  • Azure Datensee-Speicher (ADLS)

    Azure Data Lake Storage (ADLS) Gen1 ist veraltet und wird in einer der nächsten Versionen entfernt werden. Sie müssen auf ADLS Gen2 umsteigen, da ADLS Gen1 nicht mehr zur Verfügung stehen wird.

  • Amazon S3

  • Google Cloud Storage (GCS)

  • Cloud Object Storage (COS)

Übermittlung einer Spark-Anwendung ohne Zugriff auf den Katalog watsonx.data

Sie können eine Spark-Anwendung einreichen, indem Sie einen CURL-Befehl ausführen. Führen Sie die folgenden Schritte aus, um eine Python-Bewerbung einzureichen.

Führen Sie den folgenden curl-Befehl aus, um die Anwendung zur Wortzählung zu übermitteln.

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

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

Parameter:

  • <crn_instance>: Die CRN der Instanz watsonx.data.
  • <region>: Die Region, in der die Spark-Instanz bereitgestellt wird.
  • <spark_engine_id>: Die Engine-ID der Spark-Engine.
  • <token>: Der Überbringer des Tokens. Weitere Informationen zur Erzeugung des Tokens finden Sie unter Erzeugung eines Inhaber-Tokens.

Übermittlung einer Spark-Anwendung durch Zugriff auf den Katalog watsonx.data

Gehen Sie wie folgt vor, um auf Daten aus einem Katalog zuzugreifen, der mit der Spark-Engine verbunden ist, und einige grundlegende Operationen mit diesem Katalog durchzuführen:

Führen Sie folgenden curl-Befehl aus:

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

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

Parameterwerte:

  • <encoded-api-key>: Der Wert muss das Format echo -n"ibmlhapikey_<user_id>:<user’s api key>" | base64 haben. Hier ist <user_id> die IBM Cloud ID des Benutzers, dessen api-Schlüssel für den Zugriff auf den Daten-Bucket verwendet wird. Das <IAM_APIKEY> hier ist der API-Schlüssel des Benutzers, der auf den Objektspeicherbereich zugreift. Um einen API-Schlüssel zu generieren, melden Sie sich in der Konsole watsonx.data an und navigieren Sie zu Profil > Profil und Einstellungen > API-Schlüssel und generieren Sie einen neuen API-Schlüssel.
  • <storage>: Der Wert hängt von der gewählten Speicherart ab. Es muss s3a für Amazon S3 oder Cloud Object Storage (COS), abfss für ADLS und gs für GCS-Speicher sein.
  • <application_bucket_name>: Der Name des Objektspeichers, der Ihren Anwendungscode enthält. Sie müssen die Anmeldeinformationen dieses Speichers übergeben, wenn er nicht bei watsonx.data registriert ist.

Beispielanwendung Python für Iceberg-Katalog Operationen

Es folgt eine Python-Beispielanwendung zur Durchführung grundlegender Operationen mit Daten, die in einem Iceberg-Katalog gespeichert sind:

Da das Spark-Schema nicht in der Nutzlast, sondern in der Spark-Anwendung ausgewählt wird, müssen Sie, wenn Sie beim Versuch, eine Verbindung zum Iceberg-Katalog herzustellen, die Fehlermeldung [SCHEMA_NOT_FOUND] The schema \<schema_name> erhalten, sicherstellen, dass der richtige Katalog- und Schemaname in der Spark-Anwendung angegeben wird. Stellen Sie außerdem sicher, dass der Katalog, unter dem Sie das Schema suchen, mit der Spark-Engine verbunden ist.

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

Fehlerbehebung

  • Wenn Sie die Fehlermeldung [REQUIRES_SINGLE_PART_NAMESPACE] spark_catalog requires a single-part namespace, but got \iceberg_data.results erhalten, während Sie versuchen, eine Tabelle mit dem Namen likeiceberg_data.results.resultstable zu erstellen, stellen Sie sicher, dass der Katalog mit der Spark-Engine verbunden ist. Wenn Spark nicht mit einem Iceberg-Katalog konfiguriert ist, verfügt es nicht über die notwendigen Einstellungen, um das dreiteilige Namensformat für den iceberg_data-Katalog zu erkennen oder zu verwenden.

  • Da das Spark-Schema nicht in der Nutzlast, sondern in der Spark-Anwendung ausgewählt wird, müssen Sie, wenn Sie beim Versuch, eine Verbindung zum Iceberg-Katalog herzustellen, die Fehlermeldung [SCHEMA_NOT_FOUND] The schema \<schema_name> erhalten, sicherstellen, dass der richtige Katalog- und Schemaname in der Spark-Anwendung angegeben wird. Stellen Sie außerdem sicher, dass der Katalog, unter dem Sie das Schema suchen, mit der Spark-Engine verbunden ist.

Verwandte APIs

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