SparkアプリケーションREST API
IBM Analytics Engineサーバーレス・プランは、Sparkアプリケーションをサブミットして管理するためのREST APIを提供します。 以下の操作がサポートされている:
- 必要な資格情報を取得し、権限を設定します。
- Sparkアプリケーションのサブミット。
- サブミットされたSparkアプリケーションの状態の取得。
- サブミットされたSpark アプリケーションの詳細を取得します.
- 実行中のSparkアプリケーションの停止。
使用可能なAPIの説明については、サーバーレス・プラン用のIBM Analytics EngineREST APIを参照してください。
このトピックの以下のセクションでは、各Spark アプリケーション管理APIのサンプルを示します。
必要な資格情報と権限
Sparkアプリケーションをサブミットする前に、認証資格情報を取得し、分析エンジンサーバーレス・インスタンスに正しい許可を設定する必要があります。
- インスタンスをプロビジョンしたときに書き留めたサービス・インスタンスのGUIDが必要です。 GUIDをメモしていない場合は、サーバーレス・インスタンスのGUIDの取得を参照してください。
- 必要な操作を実行するには、適切な権限を持っている必要があります。 ユーザー権限を参照してください。
- SparkアプリケーションREST APIは、IAMベースの認証と許可を使用します。
Sparkアプリケーションのサブミット
Analytics Engineサーバーレスには、SparkアプリケーションをサブミットするためのRESTインターフェースが用意されています。 REST APIに渡されるペイロードは、spark-submitコマンドでサポートされるさまざまなコマンド行引数にマップされます。 詳しくは、Sparkアプリケーションをサブミットするためのパラメーターを参照してください。
Sparkアプリケーションをサブミットするときには、アプリケーション・ファイルを参照する必要があります。 速やかに開始し、AEサーバーレスSparkAPIの使用方法を学習するために、このセクションでは、サブミット・アプリケーションAPIペイロードで参照される事前にバンドルされたSparkアプリケーション・ファイルを使用する例から始めます。 後続のセクションでは、Object Storageバケットに保管されているアプリケーションを実行する方法について説明します。
事前バンドル・ファイルの参照
提供されているサンプル・アプリケーションは、ジョブ・ペイロード内の.pyワード・カウント・アプリケーション・ファイルとデータ・ファイルを参照する方法を示しています。
事前にバンドルされたサンプル・アプリケーション・ファイルの使用を速やかに開始する方法を学習するには、以下のようにします:
- まだ生成していない場合、IAMトークンを生成します。 IAMアクセス・トークンの取得を参照してください。
- トークンを変数にエクスポートします:
export token=<token generated> - ペイロードJSONファイルを準備します。 例えば、submit-spark-quick-start-app.jsonのようにします:
{ "application_details": { "application": "/opt/ibm/spark/examples/src/main/python/wordcount.py", "arguments": ["/opt/ibm/spark/examples/src/main/resources/people.txt"] } } - Sparkアプリケーションをサブミットします:
curl -X POST https://api.us-south.ae.cloud.ibm.com/v3/analytics_engines/<instance_id>/spark_applications --header "Authorization: Bearer $token" -H "content-type: application/json" -d @submit-spark-quick-start-app.json
Object Storageバケットからのファイルの参照
Object Storage バケットから Spark アプリケーション・ファイルを参照するには、バケットを作成し、そのファイルをバケットに追加して、ペイロード JSON ファイルからそのファイルを参照する必要があります。
ペイロードJSONファイルのIBM Cloud Object Storageインスタンスへのエンドポイントは、プライベートエンドポイントでなければなりません。 ダイレクトエンドポイントは、パブリックエンドポイントよりも優れたパフォーマンスを提供し、送受信帯域幅の料金は発生しません。
Sparkアプリケーションをサブミットするには、以下のようにします:
-
アプリケーション・ファイル用のバケットを作成します。 バケットの作成について詳しくは、バケット操作を参照してください。
-
新しく作成したバケットにアプリケーション・ファイルを追加します。 バケットへのアプリケーション・ファイルの追加については、オブジェクトのアップロードを参照してください。
-
まだ生成していない場合、IAMトークンを生成します。 IAMアクセス・トークンの取得を参照してください。
-
トークンを変数にエクスポートします:
export token=<token generated> -
ペイロードJSONファイルを準備します。 例えば、
submit-spark-app.jsonは以下のようになります:{ "application_details": { "application": "cos://<application-bucket-name>.<cos-reference-name>/my_spark_application.py", "arguments": ["arg1", "arg2"], "conf": { "spark.hadoop.fs.cos.<cos-reference-name>.endpoint": "https://s3.direct.us-south.cloud-object-storage.appdomain.cloud", "spark.hadoop.fs.cos.<cos-reference-name>.access.key": "<access_key>", "spark.hadoop.fs.cos.<cos-reference-name>.secret.key": "<secret_key>", "spark.app.name": "MySparkApp" } } }注:
- ペイロードの
"conf"セクションを介してSparkアプリケーション構成値を渡すことができます。 詳しくは、Sparkアプリケーションをサブミットするためのパラメーターを参照してください。 - サンプルペイロードの'
"conf"セクションの<cos-reference-name>は、'"application"IBM Cloud Object Storageインスタンスに与えられた任意の名前です。 Object Storage 資格情報について を参照してください。 - Sparkアプリケーションのサブミットには、約1分かかる場合があります。 クライアント・コードに十分なタイムアウトを設定してください。
- 応答で返された
"id"をメモします。 この値は、アプリケーションの状態の取得、アプリケーションの詳細の取得、アプリケーションの削除などの操作を実行するために必要です。
- ペイロードの
-
Sparkアプリケーションをサブミットします:
curl -X POST https://api.us-south.ae.cloud.ibm.com/v3/analytics_engines/<instance_id>/spark_applications --header "Authorization: Bearer $token" -H "content-type: application/json" -d @submit-spark-app.json応答例:
{ "id": "87e63712-a823-4aa1-9f6e-7291d4e5a113", "state": "accepted" } -
インスタンスのフォワード・ロギングが有効になっている場合は、 IBM Log Analysisに転送されたプラットフォーム・ログでアプリケーション出力を表示できます。 詳しくは、ログの構成および表示を参照してください。
アプリケーションへのSpark構成の引き渡し
ペイロードの"conf"セクションを使用して、Sparkアプリケーション構成を渡すことができます。 インスタンス・レベルでSpark構成を指定した場合、それらの構成はインスタンス上で実行されるSparkアプリケーションに継承されますが、ペイロードに"conf"セクションを含めることで、Sparkアプリケーションのサブミット時にオーバーライドできます。
Spark構成inAnalytics Engineサーバーレスを参照してください。
Sparkアプリケーションをサブミットするためのパラメーター
以下の表に、spark-submitコマンド・パラメーターと、Sparkアプリケーション実行依頼REST APIペイロードの"application_details"セクションに渡される同等のパラメーターとの間のマッピングをリストします。
| spark-submitコマンド・パラメーター | 分析エンジンSparkサブミットREST APIへのペイロード |
|---|---|
<application binary passed as spark-submit command parameter> |
application_details -> application |
<application-arguments> |
application_details -> arguments |
class |
application_details -> class |
jars |
application_details -> jars |
name |
application_details->nameまたはapplication_details->conf->spark.app.name |
packages |
application_details -> packages |
repositories |
application_details -> repositories |
files |
application_details -> files |
archives |
application_details -> archives |
driver-cores |
application_details -> conf -> spark.driver.cores |
driver-memory |
application_details -> conf -> spark.driver.memory |
driver-java-options |
application_details -> conf -> spark.driver.defaultJavaOptions |
driver-library-path |
application_details -> conf -> spark.driver.extraLibraryPath |
driver-class-path |
application_details -> conf -> spark.driver.extraClassPath |
executor-cores |
application_details -> conf -> spark.executor.cores |
executor-memory |
application_details -> conf -> spark.executor.memory |
num-executors |
application_details -> conf -> ae.spark.executor.count |
pyFiles |
application_details -> conf -> spark.submit.pyFiles |
<environment-variables> |
application_details -> env -> {"key1" : "value1", "key2" : "value2", ..... "} |
サブミットされたアプリケーションの状態の取得
サブミットされたアプリケーションの状態を取得するには、次のように入力します:
curl -X GET https://api.us-south.ae.cloud.ibm.com/v3/analytics_engines/<instance_id>/spark_applications/<application_id>/state --header "Authorization: Bearer $token"
応答例:
{
"id": "a9a6f328-56d8-4923-8042-97652fff2af3",
"state": "finished",
"start_time": "2020-11-25T14:14:31.311+0000",
"finish_time": "2020-11-25T14:30:43.625+0000"
}
サブミットされたアプリケーションの詳細の取得
サブミットされたアプリケーションの詳細を取得するには、次のように入力します:
curl -X GET https://api.us-south.ae.cloud.ibm.com/v3/analytics_engines/<instance_id>/spark_applications/<application_id> --header "Authorization: Bearer $token"
応答例:
{
"id": "ecd608d5-xxxx-xxxx-xxxx-08e27456xxxx",
"spark_application_id": "null",
"application_details": {
"application": "cos://sbn-test-bucket-serverless-1.mycosservice/my_spark_application.py",
"conf": {
"spark.hadoop.fs.cos.mycosservice.endpoint": "https://s3.direct.us-south.cloud-object-storage.appdomain.cloud",
"spark.hadoop.fs.cos.mycosservice.access.key": "xxxx",
"spark.app.name": "MySparkApp",
"spark.hadoop.fs.cos.mycosservice.secret.key": "xxxx"
},
"arguments": [
"arg1",
"arg2"
]
},
"state": "failed",
"submission_time": "2021-11-30T18:29:21+0000"
}
サブミットされたアプリケーションの停止
サブミットされたアプリケーションを停止するには、以下を実行します:
curl -X DELETE https://api.us-south.ae.cloud.ibm.com/v3/analytics_engines/<instance_id>/spark_applications/<application_id> --header "Authorization: Bearer $token"
削除が成功した場合は、204 – No Contentを返します。 アプリケーションの状態はSTOPPEDに設定されます。
このAPIはべき等です。 既に完了または停止したアプリケーションを停止しようとしても、204が返されます。
このAPIを使用して、以下の状態:accepted、waiting、submitted、およびrunningのアプリケーションを停止できます。
アプリケーションのサブミット時にランタイム Spark バージョンを渡す
ペイロード JSON スクリプトの "application_details" の下にある "runtime" セクションを使用して、アプリケーションのサブミット時に Spark ランタイム・バージョンを渡すことができます。 "runtime" セクションを通じて渡される Spark バージョンは、インスタンス・レベルで設定されたデフォルトのランタイム
Spark バージョンをオーバーライドします。 デフォルトのランタイム・バージョンについて詳しくは、 デフォルトの Spark ランタイム を参照してください。
Spark3.4:でアプリケーションを実行するための'"runtime" セクションの例:
{
"application_details": {
"application": "/opt/ibm/spark/examples/src/main/python/wordcount.py",
"arguments": [
"/opt/ibm/spark/examples/src/main/resources/people.txt"
],
"runtime": {
"spark_version": "3.4"
}
}
}
環境変数の使用
アプリケーションをサブミットするときに、ペイロード JSON スクリプトの "application_details" の下にある "env" セクションを使用して、環境固有の情報 (使用するデータ・セットや秘密値など、アプリケーションの結果を決定する情報) を渡すことができます。
ペイロード内の "env" セクションの例:
{
"application_details": {
"application": "cos://<application-bucket-name>.<cos-reference-name>/my_spark_application.py",
"arguments": ["arg1", "arg2"],
"conf": {
"spark.hadoop.fs.cos.<cos-reference-name>.endpoint": "https://s3.direct.us-south.cloud-object-storage.appdomain.cloud",
"spark.hadoop.fs.cos.<cos-reference-name>.access.key": "<access_key>",
"spark.hadoop.fs.cos.<cos-reference-name>.secret.key": "<secret_key>",
"spark.app.name": "MySparkApp"
},
"env": {
"key1": "value1",
"key2": "value2",
"key3": "value3"
}
}
}
ここで説明するように、 "application_details" > "env" を使用して設定された環境変数は、executor とドライバー・コードの両方からアクセスできます。
環境変数は、 "spark.executorEnv.[EnvironmentVariableName]" 構成 (application_details> env) を使用して設定することもできます。 ただし、これらにアクセスできるのは、ドライバーではなく、executor で実行されているタスクのみです。
シェルの環境変数名は、大文字、数字、そして
"os.getenv" 呼び出しを使用して渡される環境変数にアクセスする pyspark アプリケーションの例。
from pyspark.sql.types import IntegerType
import os
def init_spark():
spark = SparkSession.builder.appName("spark-env-test").getOrCreate()
sc = spark.sparkContext
return spark,sc
def returnExecutorEnv(x):
# Attempt to access environment variable from a task running on executor
return os.getenv("TESTENV1")
def main():
spark,sc = init_spark()
# dummy dataframe
df=spark.createDataFrame([("1","one")])
df.show()
df.rdd.map(lambda x: (x[0],returnExecutorEnv(x[0]))).toDF().show()
# Attempt to access environment variable on driver
print (os.getenv("TESTENV1"))
spark.stop()
if __name__ == '__main__':
main()
デフォルト以外の言語バージョンで Spark アプリケーションを実行する
Spark ランタイムは、以下の言語で作成された Spark アプリケーションをサポートします。
- Scala
- Python
- R
Spark ランタイム・バージョンには、デフォルトのランタイム言語バージョンが付属しています。 IBM は、新しい言語バージョンのサポートを拡張し、既存の言語バージョンを削除して、ランタイムにセキュリティーの脆弱性がないようにします。 また、システムは、新しい言語バージョンが存在する場合にワークロードを移行するための整定時間も提供します。 アプリケーションの言語バージョンを指す環境変数を渡すことにより、言語バージョンを使用してワークロードをテストできます。
サンプル Python コード:
{
"application_details": {
"application": "/opt/ibm/spark/examples/src/main/python/wordcount.py",
"arguments": [
"/opt/ibm/spark/examples/src/main/resources/people.txt"
],
"env": {
"RUNTIME_PYTHON_ENV": "python310"
}
}
}
サンプル R コード:
{
"application_details": {
"env": {
"RUNTIME_R_ENV": "r42"
},
"application": "/opt/ibm/spark/examples/src/main/r/dataframe.R"
}
}
詳細情報
Spark アプリケーションを管理する場合は、推奨されている ベスト・プラクティス に従ってください。