Spark 應用程式 REST API

IBM Analytics Engine 無伺服器方案提供 REST API 來提交及管理 Spark 應用程式。 支援以下操作:

  1. 取得必要的認證並設定許可權
  2. 提交 Spark 申請
  3. 擷取已提交 Spark 應用程式的狀態
  4. 擷取已提交 Spark 應用程式的詳細資料
  5. 停止執行中的 Spark 應用程式

如需可用 API 的說明,請參閱 IBM Analytics Engine 無伺服器方案的 REST API

本主題中的下列各節顯示每一個 Spark 應用程式管理 API 的範例。

必要的認證和許可權

您需要先取得鑑別認證,並在 Analytics Engine 無伺服器實例上設定正確的許可權,然後才能提交 Spark 應用程式。

  1. 您需要在佈建實例時記下之服務實例的 GUID。 如果您未記下 GUID,請參閱 擷取無伺服器實例的 GUID
  2. 您必須具有正確的許可權,才能執行必要的作業。 請參閱 使用者許可權
  3. Spark 應用程式 REST API 使用 IAM 型鑑別及授權。

提交 Spark 應用程式

Analytics Engine Serverless 提供 REST 介面來提交 Spark 應用程式。 傳遞至 REST API 的有效負載會對映至 spark-submit 指令支援的各種指令行引數。 如需詳細資料,請參閱 用於提交 Spark 應用程式的參數

提交 Spark 應用程式時,您需要參照應用程式檔案。 為了協助您快速開始使用並瞭解如何使用 AE 無伺服器 Spark API,本節以使用提交應用程式 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> 是提供給 IBM Cloud Object Storage 實例的任何名稱,您在 "application" 參數的 URL 中參照此名稱。 請參閱 瞭解 Object Storage 認證
    • 提交 Spark 應用程式大約需要一分鐘。 請確保在用戶端程式碼中設定足夠的逾時。
    • 記下回應中傳回的 "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 應用程式會在實例上執行,但在提交 Spark 應用程式時可以置換這些配置,方法是在有效負載中包括 "conf" 區段。

請參閱 Analytics Engine 無伺服器 中的 Spark 配置。

用於提交 Spark 應用程式的參數

下表列出 spark-submit 指令參數與其對等項目之間的對映,以傳遞至 Spark 應用程式提交 REST API 有效負載的 "application_details" 區段。

Spark-submit 指令參數與傳遞給有效負載的等效參數之間的映射
spark-submit 指令參數 Analytics Engine 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-> nameapplication_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 來停止處於下列狀態的應用程式: acceptedwaitingsubmittedrunning

提交應用程式時傳遞執行時期 Spark 版本

提交應用程式時,您可以使用有效負載 JSON Script 中 "application_details" 下的 "runtime" 區段來傳遞 Spark 執行時期版本。 透過 "runtime" 區段傳遞的 Spark 版本會置換在實例層次設定的預設執行時期 Spark 版本。 若要進一步瞭解預設執行時期版本,請參閱 預設 Spark 執行時期

在 Spark 3.4 中執行應用程式的 "runtime" 部分範例3.4:

{
    "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 Script 中使用 "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" 所設定的環境變數。

也可以使用 "spark.executorEnv.[EnvironmentVariableName]" 配置 (application_details> env) 來設定環境變數。 不過,它們只能供執行程式上執行的作業存取,而不能供驅動程式存取。

Shell中的環境變數名稱由大寫字母、數字和 ( '_' ) 並且不以數字開頭。

存取使用 "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 應用程式時,請遵循建議的 最佳作法