네이티브 스파크 엔진을 사용하여 스파크 애플리케이션 제출하기

적용됩니다: 스파크 엔진 글루텐 가속 스파크 엔진

이 주제는 IBM Cloud watsonx.data 네이티브 Spark 엔진을 사용하여 Spark 애플리케이션을 제출하는 절차를 제공합니다.

전제조건

  • 개체 저장소 만들기: Spark 애플리케이션과 관련 출력을 저장하려면 저장소 버킷을 만드세요. Cloud Object Storage와 버킷을 만들려면 스토리지 버킷 만들기 를 참조하세요. 애플리케이션과 데이터를 위한 별도의 스토리지를 유지하세요. watsonx.data로 데이터 버킷만 등록하세요.

  • Cloud Object Storage: Cloud Object Storage 버킷을 watsonx.data에 등록합니다. Cloud Object Storage 버킷을 등록하려면 버킷 카탈로그 쌍 추가하기 를 참조하세요.

    애플리케이션 코드와 출력을 저장하기 위해 다른 Cloud Object Storage 버킷을 만들 수 있습니다. 입력 데이터를 저장하는 데이터 버킷과 watsonx.data 테이블을 등록합니다. 애플리케이션 코드를 유지하는 스토리지 버킷을 watsonx.data로 등록할 필요는 없습니다.

  • 스토리지를 스파크 엔진과 연결합니다. Spark 엔진과 연동하는 방법에 대한 자세한 내용은 카탈로그와 엔진 연동하기를 참조하세요.

지원되는 스토리지

  • Azure 데이터 레이크 스토리지(ADLS)

    Azure 데이터 레이크 스토리지(ADLS)( Gen1 )는 더 이상 사용되지 않으며 다음 릴리스에서 제거될 예정입니다. ADLS Gen1 을 더 이상 사용할 수 없으므로 ADLS Gen2 로 전환해야 합니다.

  • Amazon S3

  • Google Cloud Storage (GCS)

  • Cloud Object Storage(COS)

watsonx.data 카탈로그에 액세스하지 않고 Spark 애플리케이션 제출하기

CURL 명령을 실행하여 Spark 애플리케이션을 제출할 수 있습니다. 다음 단계를 완료하여 Python 신청서를 제출하세요.

다음 curl 명령을 실행하여 단어 수 신청서를 제출합니다.

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

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

매개변수:

  • <crn_instance> watsonx.data 인스턴스의 CRN입니다.
  • <region>: Spark 인스턴스가 프로비저닝되는 리전입니다.
  • <spark_engine_id>: 스파크 엔진의 엔진 ID입니다.
  • <token> 무기명 토큰입니다. 토큰 생성에 대한 자세한 내용은 무기명 토큰 생성하기 를 참조하세요.

watsonx.data 카탈로그에 액세스하여 Spark 애플리케이션 제출하기

스파크 엔진과 연결된 카탈로그의 데이터에 액세스하고 해당 카탈로그에서 몇 가지 기본 작업을 수행하려면 다음과 같이 하세요:

다음 curl 명령을 실행하십시오.

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

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

매개변수 값:

  • <encoded-api-key>: 값은 echo -n"ibmlhapikey_<user_id>:<user’s api key>" | base64 형식이어야 합니다. 여기서 <user_id>는 데이터 버킷에 액세스하는 데 API 키를 사용하는 사용자의 IBM Cloud ID입니다. 여기서 <IAM_APIKEY> 은 오브젝트 스토어 버킷에 액세스하는 사용자의 API 키입니다. API 키를 생성하려면 watsonx.data 콘솔에 로그인하고 프로필 > 프로필 및 설정 > API 키로 이동하여 새 API 키를 생성하세요.
  • <storage>: 선택한 스토리지 유형에 따라 값이 달라집니다. Amazon S3 의 경우 s3a, Cloud Object Storage(COS)의 경우 abfss, ADLS의 경우 gs, GCS 스토리지의 경우 이어야 합니다.
  • <application_bucket_name>: 애플리케이션 코드가 포함된 개체 저장소의 이름입니다. 이 저장소가 watsonx.data에 등록되어 있지 않은 경우 이 저장소의 자격 증명을 전달해야 합니다.

Iceberg 카탈로그 작업을 위한 Python 애플리케이션 샘플

다음은 Iceberg 카탈로그에 저장된 데이터에 대한 기본 연산을 수행하는 샘플 Python 애플리케이션입니다:

스파크 스키마는 페이로드에서 선택되지 않고 스파크 애플리케이션 내에서 선택되므로 빙산 카탈로그에 연결하려고 할 때 [SCHEMA_NOT_FOUND] The schema \<schema_name> 오류가 발생하면 스파크 애플리케이션에서 올바른 카탈로그와 스키마 이름이 제공되었는지 확인하세요. 또한 스키마를 찾고자 하는 카탈로그가 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()

문제점 해결

  • 오류가 발생하면 [REQUIRES_SINGLE_PART_NAMESPACE] spark_catalog requires a single-part namespace, but got \iceberg_data.results 세 개의 파트 이름으로 테이블을 만들려고 하는 동안 likeiceberg_data.results.resultstable, 카탈로그가 스파크 엔진과 연결되어 있는지 확인합니다. Spark가 Iceberg 카탈로그로 구성되지 않은 경우, iceberg_data 카탈로그의 세 부분으로 구성된 이름 형식을 인식하거나 사용하는 데 필요한 설정이 없습니다.

  • 스파크 스키마는 페이로드에서 선택되지 않고 스파크 애플리케이션 내에서 선택되므로 빙산 카탈로그에 연결하려고 할 때 [SCHEMA_NOT_FOUND] The schema \<schema_name> 오류가 발생하면 스파크 애플리케이션에서 올바른 카탈로그와 스키마 이름이 제공되었는지 확인하세요. 또한 스키마를 찾고자 하는 카탈로그가 Spark 엔진과 연결되어 있는지 확인하세요.

관련 API

관련 API에 대한 자세한 내용은 다음을 참조하세요