다른 테이블 형식에 대한 작업

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

이 주제에서는 Apache Hudi, Apache Iceberg 또는 Delta Lake 카탈로그와 같은 다양한 테이블 형식으로 데이터를 수집하는 Spark 애플리케이션을 실행하는 절차에 대해 설명합니다.

  1. 필요한 카탈로그(카탈로그는 Apache Hudi, Apache Iceberg 또는 Delta Lake )가 있는 저장소를 만들어 Spark 애플리케이션에서 사용되는 데이터를 저장합니다. 스토리지를 작성하려면 스토리지-카탈로그 쌍 추가 를 참조하십시오.

  2. 스토리지를 원시 Spark 엔진과 연관시키십시오. 자세한 정보는 카탈로그를 엔진과 연관 을 참조하십시오.

  3. Cloud Object Storage (COS) 를 작성하여 Spark 애플리케이션을 저장하십시오. Cloud Object Storage 및 버킷을 작성하려면 스토리지 버킷 작성 을 참조하십시오.

  4. watsonx.data에서 Cloud Object Storage 를 등록하십시오. 자세한 정보는 스토리지-카탈로그 쌍 추가 를 참조하십시오.

  5. 선택한 카탈로그에 따라 다음 Spark 애플리케이션 (Python 파일) 을 로컬 시스템에 저장하십시오. 여기서 iceberg_demo.py, hudi_demo.py 또는 delta_demo.py 이며 Spark 애플리케이션을 COS에 업로드하십시오. 데이터 업로드 를 참조하십시오.

  6. Cloud Object Storage에 있는 데이터를 사용하여 Spark 애플리케이션을 제출하려면 매개변수 값을 지정하고 다음 표에서 curl 명령을 실행하십시오.

    • Apache Iceberg

    샘플 파일은 다음과 같은 기능을 보여줍니다:

    • watsonx.data 테이블에 액세스하기

    • watsonx.data 데이터 수집하기

    • watsonx.data 스키마 수정하기

    watsonx.data 테이블 유지 관리 활동 수행.

    COS 버킷에 데이터를 삽입해야 합니다. 자세한 내용은 COS 버킷에 샘플 데이터 삽입하기를 참조하세요.

    Python 애플리케이션: Iceberg Python 파일

    Curl 명령을 사용하여 Python 애플리케이션을 제출합니다:

    샘플 V2 API

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

    샘플 V3 API

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

    매개변수 값 :

    • <wxd_host_name>': watsonx.data 클라우드 인스턴스의 호스트 이름입니다.

    • <instance_id>: watsonx.data 인스턴스 URL 인스턴스 ID입니다. 예: 1609968977179454.

    • <spark_engine_id>': 네이티브 스파크 엔진의 엔진 ID입니다.

    • <token>: 무기명 토큰입니다. 토큰 생성에 대한 자세한 내용은 IAM 토큰을 참조하세요.

    • <user-authentication-string>: 사용자 ID와 API 키의 베이스 64 인코딩 문자열이어야 합니다. 형식에 대한 자세한 내용은 참고 사항을 참조하세요.

    • Apache Hudi

    Python 스파크 애플리케이션은 다음과 같은 기능을 보여줍니다:

    • Apache Hudi 카탈로그(데이터를 저장하기 위해 만든 카탈로그) 내에 데이터베이스를 만듭니다. 여기서는 ' <database_name>.

    • ' <database_name> 데이터베이스 내에 테이블, 즉 ' <table_name>'를 생성합니다.

    • ' <table_name> '에 데이터를 삽입하고 SELECT 쿼리 연산을 수행합니다.

    • 사용 후 테이블과 스키마를 삭제합니다.

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

    매개변수 값:

    • <database_name>: 만들려는 데이터베이스의 이름을 지정합니다.
    • <table_name>: 만들려는 테이블의 이름을 지정합니다.
    • <data_storage_name>: 생성한 Apache Hudi 스토리지의 이름을 지정합니다.

    Python 애플리케이션을 제출하기 위한 Curl 명령

    샘플 V2 API

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

    샘플 V3 API

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

    매개변수 값:

    • <wxd_host_name>': watsonx.data 클라우드 인스턴스의 호스트 이름입니다.

    • <instance_id>: watsonx.data 인스턴스 URL 인스턴스 ID입니다. 예: 1609968977179454.

    • <spark_engine_id>': 네이티브 스파크 엔진의 엔진 ID입니다.

    • <token>: 무기명 토큰입니다. 토큰 생성에 대한 자세한 내용은 IAM 토큰을 참조하세요.

    • <user-authentication-string>: 사용자 ID와 API 키의 베이스 64 인코딩 문자열이어야 합니다. 형식에 대한 자세한 내용은 참고 사항을 참조하세요.

    • Delta Lake

    Python 스파크 애플리케이션은 다음과 같은 기능을 보여줍니다:

    • Delta Lake 카탈로그(데이터를 저장하기 위해 만든 카탈로그) 내에 데이터베이스를 만듭니다. 여기서는 ' <database_name>.

    • ' <database_name> 데이터베이스 내에 테이블, 즉 ' <table_name>'를 생성합니다.

    • ' <table_name> '에 데이터를 삽입하고 SELECT 쿼리 연산을 수행합니다.

    • 사용 후 테이블과 스키마를 삭제합니다.

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

    매개변수 값:

    • <database_name>: 만들려는 데이터베이스의 이름을 지정합니다.
    • <table_name>: 만들려는 테이블의 이름을 지정합니다.
    • <data_storage_name>: 생성한 Apache Hudi 스토리지의 이름을 지정합니다.

    Python 애플리케이션을 제출하기 위한 Curl 명령

    샘플 V2 API

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

    샘플 V3 API

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

    매개변수 값

    • <wxd_host_name>': watsonx.data 클라우드 인스턴스의 호스트 이름입니다.
    • <instance_id> watsonx.data 클러스터 인스턴스 URL 인스턴스 ID입니다. 예: 1609968977179454.
    • <spark_engine_id> ': 네이티브 스파크 엔진의 엔진 ID입니다.
    • <token> 무기명 토큰입니다. 토큰 생성에 대한 자세한 내용은 IAM 토큰을 참조하세요.
    • <user-authentication-string>: 사용자 ID와 API 키의 베이스 64 인코딩 문자열이어야 합니다. 형식에 대한 자세한 정보는 다음 참고를 참조하십시오.

    <user-authentication-string> 의 값은 echo -n 'ibmlhapikey_<username>:<user_apikey>' | base64 형식이어야 합니다. 여기서 <user_id> 는 apikey가 데이터 버킷에 액세스하는 데 사용되는 사용자의 IBM Cloud ID입니다. 여기서 <IAM_APIKEY> 는 오브젝트 저장소 버킷에 액세스하는 사용자의 API키입니다. API키를 생성하려면 watsonx.data 콘솔에 로그인하고 프로파일 > 프로파일 및 설정 > API키로 이동하여 새 API키를 생성하십시오. 새 API키를 생성하면 이전 API키가 올바르지 않게 됩니다.

  7. Spark 애플리케이션을 제출하면 애플리케이션 ID및 Spark 버전과 함께 확인 메시지가 수신됩니다. 참조를 위해 저장하십시오.

  8. watsonx.data 클러스터에 로그인하고 엔진 세부사항 페이지에 액세스하십시오. 애플리케이션 탭에서 애플리케이션 ID를 사용하여 애플리케이션을 나열하고 스테이지를 추적하십시오. 자세한 정보는 애플리케이션 보기 및 관리 를 참조하십시오.

관련 API

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