ネイティブSparkエンジンを使用したSparkアプリケーションの提出
適用範囲 : スパークエンジン グルテン加速スパークエンジン
このトピックでは IBM Cloudのネイティブ Spark エンジンを使用して、Spark アプリケーションwatsonx.dataに提出する手順を説明します。
前提条件
-
オブジェクトストレージの作成: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エンジンに関連付ける。 Sparkエンジンと関連付ける方法については、 カタログとエンジンの関連付けを 参照してください。
対応ストレージ
-
Azure データレイク・ストレージ(ADLS)
Azure Data Lake Storage (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アプリケーションを提出する
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> は、IBM Cloud IDであり、そのapiキーでデータバケットにアクセスします。 ここでの<IAM_APIKEY>は、Object store バケットにアクセスするユーザの API キーです。 API キーを生成するには、watsonx.data コンソールにログインし、[Profile] > [Profile and Settings] > [API Keys] に移動して新しい API キーを生成します。<storage>この値は、選択したストレージタイプによって異なります。s3aは Amazon S3 またはクラウドオブジェクトストレージ(COS)、abfssはADLS、gsはGCSストレージ用でなければなりません。<application_bucket_name>: アプリケーションコードが格納されているオブジェクトストレージの名前。 watsonx.dataに登録されていない場合は、このストレージの認証情報を渡す必要があります。
Iceberg カタログ操作のためのサンプル Python アプリケーション
以下はPythonアプリケーションのサンプルで、Icebergカタログに格納されたデータに対して基本的な操作を行います:
sparkスキーマはペイロードでは選択されず、sparkアプリケーション内で選択されるため、icebergカタログに接続しようとして [SCHEMA_NOT_FOUND] The schema \<schema_name> エラーが表示された場合は、sparkアプリケーションで正しいカタログとスキーマ名が提供されていることを確認してください。 また、スキーマを探すカタログが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()
トラブルシューティング
-
likeiceberg_data.results.resultstable という 3 つのパート名のテーブルを作成しようとして
[REQUIRES_SINGLE_PART_NAMESPACE] spark_catalog requires a single-part namespace, but got \iceberg_data.resultsというエラーが表示された場合は、カタログがスパーク エンジンに関連付けられていることを確認してください。 SparkにIcebergカタログが設定されていない場合、iceberg_dataカタログの3つの部分からなる名前形式を認識したり使用したりするのに必要な設定がありません。 -
sparkスキーマはペイロードでは選択されず、sparkアプリケーション内で選択されるため、icebergカタログに接続しようとして
[SCHEMA_NOT_FOUND] The schema \<schema_name>エラーが表示された場合は、sparkアプリケーションで正しいカタログとスキーマ名が提供されていることを確認してください。 また、スキーマを探すカタログがSparkエンジンに関連付けられていることを確認してください。
関連API
関連APIについては