SparkアプリケーションREST API

IBM Analytics Engineサーバーレス・プランは、Sparkアプリケーションをサブミットして管理するためのREST APIを提供します。 以下の操作がサポートされている:

  1. 必要な資格情報を取得し、権限を設定します
  2. Sparkアプリケーションのサブミット
  3. サブミットされたSparkアプリケーションの状態の取得
  4. サブミットされたSpark アプリケーションの詳細を取得します.
  5. 実行中のSparkアプリケーションの停止

使用可能なAPIの説明については、サーバーレス・プラン用のIBM Analytics EngineREST APIを参照してください。

このトピックの以下のセクションでは、各Spark アプリケーション管理APIのサンプルを示します。

必要な資格情報と権限

Sparkアプリケーションをサブミットする前に、認証資格情報を取得し、分析エンジンサーバーレス・インスタンスに正しい許可を設定する必要があります。

  1. インスタンスをプロビジョンしたときに書き留めたサービス・インスタンスのGUIDが必要です。 GUIDをメモしていない場合は、サーバーレス・インスタンスのGUIDの取得を参照してください。
  2. 必要な操作を実行するには、適切な権限を持っている必要があります。 ユーザー権限を参照してください。
  3. SparkアプリケーションREST APIは、IAMベースの認証と許可を使用します。

Sparkアプリケーションのサブミット

Analytics Engineサーバーレスには、SparkアプリケーションをサブミットするためのRESTインターフェースが用意されています。 REST APIに渡されるペイロードは、spark-submitコマンドでサポートされるさまざまなコマンド行引数にマップされます。 詳しくは、Sparkアプリケーションをサブミットするためのパラメーターを参照してください。

Sparkアプリケーションをサブミットするときには、アプリケーション・ファイルを参照する必要があります。 速やかに開始し、AEサーバーレスSparkAPIの使用方法を学習するために、このセクションでは、サブミット・アプリケーションAPIペイロードで参照される事前にバンドルされたSparkアプリケーション・ファイルを使用する例から始めます。 後続のセクションでは、Object Storageバケットに保管されているアプリケーションを実行する方法について説明します。

事前バンドル・ファイルの参照

提供されているサンプル・アプリケーションは、ジョブ・ペイロード内の.pyワード・カウント・アプリケーション・ファイルとデータ・ファイルを参照する方法を示しています。

事前にバンドルされたサンプル・アプリケーション・ファイルの使用を速やかに開始する方法を学習するには、以下のようにします:

  1. まだ生成していない場合、IAMトークンを生成します。 IAMアクセス・トークンの取得を参照してください。
  2. トークンを変数にエクスポートします:
    export token=<token generated>
    
  3. ペイロード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"]
        }
    }
    
  4. 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アプリケーションをサブミットするには、以下のようにします:

  1. アプリケーション・ファイル用のバケットを作成します。 バケットの作成について詳しくは、バケット操作を参照してください。

  2. 新しく作成したバケットにアプリケーション・ファイルを追加します。 バケットへのアプリケーション・ファイルの追加については、オブジェクトのアップロードを参照してください。

  3. まだ生成していない場合、IAMトークンを生成します。 IAMアクセス・トークンの取得を参照してください。

  4. トークンを変数にエクスポートします:

    export token=<token generated>
    
  5. ペイロード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"をメモします。 この値は、アプリケーションの状態の取得、アプリケーションの詳細の取得、アプリケーションの削除などの操作を実行するために必要です。
  6. 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"
    }
    
  7. インスタンスのフォワード・ロギングが有効になっている場合は、 IBM Log Analysisに転送されたプラットフォーム・ログでアプリケーション出力を表示できます。 詳しくは、ログの構成および表示を参照してください。

アプリケーションへのSpark構成の引き渡し

ペイロードの"conf"セクションを使用して、Sparkアプリケーション構成を渡すことができます。 インスタンス・レベルでSpark構成を指定した場合、それらの構成はインスタンス上で実行されるSparkアプリケーションに継承されますが、ペイロードに"conf"セクションを含めることで、Sparkアプリケーションのサブミット時にオーバーライドできます。

Spark構成inAnalytics Engineサーバーレスを参照してください。

Sparkアプリケーションをサブミットするためのパラメーター

以下の表に、spark-submitコマンド・パラメーターと、Sparkアプリケーション実行依頼REST APIペイロードの"application_details"セクションに渡される同等のパラメーターとの間のマッピングをリストします。

spark-submitコマンドのパラメータと、ペイロードに渡される同等のパラメータとのマッピング
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を使用して、以下の状態:acceptedwaitingsubmitted、および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 アプリケーションを管理する場合は、推奨されている ベスト・プラクティス に従ってください。