다른 테이블 형식에 대한 작업
에 적용됩니다: 스파크 엔진 글루텐 가속 스파크 엔진
이 주제에서는 Apache Hudi, Apache Iceberg 또는 Delta Lake 카탈로그와 같은 다양한 테이블 형식으로 데이터를 수집하는 Spark 애플리케이션을 실행하는 절차에 대해 설명합니다.
-
필요한 카탈로그(카탈로그는 Apache Hudi, Apache Iceberg 또는 Delta Lake )가 있는 저장소를 만들어 Spark 애플리케이션에서 사용되는 데이터를 저장합니다. 스토리지를 작성하려면 스토리지-카탈로그 쌍 추가 를 참조하십시오.
-
스토리지를 원시 Spark 엔진과 연관시키십시오. 자세한 정보는 카탈로그를 엔진과 연관 을 참조하십시오.
-
Cloud Object Storage (COS) 를 작성하여 Spark 애플리케이션을 저장하십시오. Cloud Object Storage 및 버킷을 작성하려면 스토리지 버킷 작성 을 참조하십시오.
-
watsonx.data에서 Cloud Object Storage 를 등록하십시오. 자세한 정보는 스토리지-카탈로그 쌍 추가 를 참조하십시오.
-
선택한 카탈로그에 따라 다음 Spark 애플리케이션 (Python 파일) 을 로컬 시스템에 저장하십시오. 여기서
iceberg_demo.py,hudi_demo.py또는delta_demo.py이며 Spark 애플리케이션을 COS에 업로드하십시오. 데이터 업로드 를 참조하십시오. -
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키가 올바르지 않게 됩니다. -
Spark 애플리케이션을 제출하면 애플리케이션 ID및 Spark 버전과 함께 확인 메시지가 수신됩니다. 참조를 위해 저장하십시오.
-
watsonx.data 클러스터에 로그인하고 엔진 세부사항 페이지에 액세스하십시오. 애플리케이션 탭에서 애플리케이션 ID를 사용하여 애플리케이션을 나열하고 스테이지를 추적하십시오. 자세한 정보는 애플리케이션 보기 및 관리 를 참조하십시오.
관련 API
관련 API에 대한 자세한 내용은 다음을 참조하세요