All Products
Search
Document Center

AnalyticDB:Mengembangkan aplikasi Spark menggunakan Python SDK

Last Updated:Aug 25, 2026

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_ID dan ALIBABA_CLOUD_ACCESS_KEY_SECRET telah dikonfigurasi. Untuk informasi selengkapnya, lihat Konfigurasikan variabel lingkungan di Linux, macOS, dan Windows.

  • Jalur penyimpanan untuk log Spark telah dikonfigurasi.

    Catatan

    Anda 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.rootPath untuk 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)