Topik ini menjelaskan cara mengirim pekerjaan Spark, memeriksa status dan log pekerjaan Spark, menghentikan pekerjaan Spark, serta menampilkan riwayat pekerjaan Spark menggunakan Python SDK.
Prasyarat
-
Lingkungan Python telah terinstal dengan versi 3.7 atau lebih baru.
-
Kelompok sumber daya Job telah dibuat di kluster AnalyticDB for MySQL. Untuk informasi selengkapnya, lihat Buat kelompok sumber daya Job.
-
Kluster AnalyticDB for MySQL menggunakan Edisi Perusahaan, Edisi Dasar, atau Edisi Lakehouse.
-
Python SDK telah terinstal. Untuk informasi selengkapnya, lihat AnalyticDB MySQL SDK untuk Python.
-
Variabel lingkungan
ALIBABA_CLOUD_ACCESS_KEY_IDdanALIBABA_CLOUD_ACCESS_KEY_SECRETtelah dikonfigurasi. Untuk informasi selengkapnya, lihat Konfigurasikan variabel lingkungan di Linux, macOS, dan Windows. -
Jalur penyimpanan untuk log Spark telah dikonfigurasi.
CatatanAnda dapat mengonfigurasi jalur penyimpanan log Spark dengan salah satu metode berikut:
-
Pada konsol AnalyticDB for MySQL, buka halaman Spark Jar Development, lalu klik Log Configuration di pojok kanan atas untuk mengatur jalur penyimpanan log Spark.
-
Gunakan item konfigurasi
spark.app.log.rootPathuntuk menentukan jalur OSS guna menyimpan log eksekusi pekerjaan Spark.
-
Contoh
Kode contoh berikut menunjukkan cara mengirim pekerjaan Spark, memeriksa status dan log pekerjaan Spark, menghentikan pekerjaan Spark, serta menampilkan riwayat pekerjaan Spark.
from alibabacloud_adb20211201.models import SubmitSparkAppRequest, SubmitSparkAppResponse, GetSparkAppStateRequest, \
GetSparkAppStateResponse, GetSparkAppLogResponse, GetSparkAppLogRequest, KillSparkAppRequest, \
KillSparkAppResponse, ListSparkAppsRequest, ListSparkAppsResponse
from alibabacloud_tea_openapi.models import Config
from alibabacloud_adb20211201.client import Client
import os
def submit_spark_sql(client: Client, cluster_id, rg_name, sql):
"""
Mengirim pekerjaan Spark SQL
:param client: Klien Alibaba Cloud
:param cluster_id: ID kluster
:param rg_name: Nama kelompok sumber daya
:param sql: SQL
:return: ID pekerjaan Spark
:rtype: basestring
:exception ClientException
"""
# Inisialisasi permintaan
request = SubmitSparkAppRequest(
dbcluster_id=cluster_id,
resource_group_name=rg_name,
data=sql,
app_type="SQL",
agent_source="Python SDK",
agent_version="1.0.0"
)
# Kirim pekerjaan SQL dan dapatkan hasilnya
response: SubmitSparkAppResponse = client.submit_spark_app(request)
# Dapatkan ID pekerjaan Spark
print(response)
return response.body.data.app_id
def submit_spark_jar(client: Client, cluster_id: str, rg_name: str, json_conf: str):
"""
Mengirim pekerjaan Spark
:param client: Klien Alibaba Cloud
:param cluster_id: ID kluster
:param rg_name: Nama kelompok sumber daya
:param json_conf: Konfigurasi JSON
:return: ID pekerjaan Spark
:rtype: basestring
:exception ClientException
"""
# Inisialisasi permintaan
request = SubmitSparkAppRequest(
dbcluster_id=cluster_id,
resource_group_name=rg_name,
data=json_conf,
app_type="BATCH",
agent_source="Python SDK",
agent_version="1.0.0"
)
# Kirim pekerjaan SQL dan dapatkan hasilnya
response: SubmitSparkAppResponse = client.submit_spark_app(request)
# Dapatkan ID pekerjaan Spark
print(response)
return response.body.data.app_id
def get_status(client: Client, app_id):
"""
Menanyakan status pekerjaan Spark
:param client: Klien Alibaba Cloud
:param app_id: ID pekerjaan Spark
:return: Status pekerjaan Spark
:rtype: basestring
:exception ClientException
"""
# Inisialisasi permintaan
print(app_id)
request = GetSparkAppStateRequest(app_id=app_id)
# Dapatkan status pekerjaan Spark
response: GetSparkAppStateResponse = client.get_spark_app_state(request)
print(response)
return response.body.data.state
def get_log(client: Client, app_id):
"""
Menanyakan log pekerjaan Spark
:param client: Klien Alibaba Cloud
:param app_id: ID pekerjaan Spark
:return: Log pekerjaan Spark
:rtype: basestring
:exception ClientException
"""
# Inisialisasi permintaan
request = GetSparkAppLogRequest(app_id=app_id)
# Dapatkan log pekerjaan Spark
response: GetSparkAppLogResponse = client.get_spark_app_log(request)
print(response)
return response.body.data.log_content
def kill_app(client: Client, app_id):
"""
Menghentikan pekerjaan Spark
:param client: Klien Alibaba Cloud
:param app_id: ID pekerjaan Spark
:return: Status pekerjaan Spark
:exception ClientException
"""
# Inisialisasi permintaan
request = KillSparkAppRequest(app_id=app_id)
# Dapatkan status pekerjaan Spark
response: KillSparkAppResponse = client.kill_spark_app(request)
print(response)
return response.body.data.state
def list_apps(client: Client, cluster_id: str, page_number: int, page_size: int):
"""
Menanyakan riwayat pekerjaan Spark
:param client: Klien Alibaba Cloud
:param cluster_id: ID kluster
:param page_number: Nomor halaman, harus berupa bilangan bulat positif. Nilai default: 1
:param page_size: Jumlah entri per halaman
:return: Detail pekerjaan Spark
:exception ClientException
"""
# Inisialisasi permintaan
request = ListSparkAppsRequest(
dbcluster_id=cluster_id,
page_number=page_number,
page_size=page_size
)
# Dapatkan detail pekerjaan Spark
response: ListSparkAppsResponse = client.list_spark_apps(request)
print("Total App Number:", response.body.data.page_number)
for app_info in response.body.data.app_info_list:
print(app_info.app_id)
print(app_info.state)
print(app_info.detail)
if __name__ == '__main__':
# konfigurasi klien
config = Config(
# Dapatkan AccessKey ID dari variabel lingkungan ALIBABA_CLOUD_ACCESS_KEY_ID
access_key_id=os.environ['ALIBABA_CLOUD_ACCESS_KEY_ID'],
# Dapatkan AccessKey Secret dari variabel lingkungan ALIBABA_CLOUD_ACCESS_KEY_SECRET
access_key_secret=os.environ['ALIBABA_CLOUD_ACCESS_KEY_SECRET'],
# Titik akhir. cn-hangzhou adalah ID wilayah kluster.
endpoint="adb.cn-hangzhou.aliyuncs.com"
)
# buat klien baru
adb_client = Client(config)
sql_str = """
-- Berikut ini hanya contoh SparkSQL. Ubah kontennya dan jalankan program Spark Anda.
set spark.driver.resourceSpec=medium;
set spark.executor.instances=2;
set spark.executor.resourceSpec=medium;
set spark.app.name=Spark SQL Test;
-- Berikut adalah pernyataan SQL Anda
show databases;
"""
json_str = """
{
"comments": [
"-- Berikut ini hanya contoh SparkPi. Ubah kontennya dan jalankan program Spark Anda."
],
"args": [
"1000"
],
"file": "local:///tmp/spark-examples.jar",
"name": "SparkPi",
"className": "org.apache.spark.examples.SparkPi",
"conf": {
"spark.driver.resourceSpec": "medium",
"spark.executor.instances": 2,
"spark.executor.resourceSpec": "medium"
}
}
"""
"""
Mengirim pekerjaan Spark SQL
cluster_id: ID kluster
rg_name: Nama kelompok sumber daya
"""
sql_app_id = submit_spark_sql(client=adb_client, cluster_id="amv-bp1wo70f0k3c****", rg_name="test", sql=sql_str)
print(sql_app_id)
"""
Mengirim pekerjaan Spark
cluster_id: ID kluster
rg_name: Nama kelompok sumber daya
"""
json_app_id = submit_spark_jar(client=adb_client, cluster_id="amv-bp1wo70f0k3c****",
rg_name="test", json_conf=json_str)
print(json_app_id)
# Menanyakan status pekerjaan Spark
get_status(client=adb_client, app_id=sql_app_id)
get_status(client=adb_client, app_id=json_app_id)
"""
Menanyakan riwayat pekerjaan Spark
cluster_id: ID kluster
page_number: Nomor halaman, harus berupa bilangan bulat positif. Nilai default: 1
page_size: Jumlah entri per halaman
"""
list_apps(client=adb_client, cluster_id="amv-bp1wo70f0k3c****", page_size=10, page_number=1)
# Menghentikan pekerjaan Spark
kill_app(client=adb_client, app_id=json_app_id)