All Products
Search
Document Center

AnalyticDB:Gunakan Airflow untuk menjadwalkan pekerjaan Spark

Last Updated:Sep 17, 2026

Panduan ini menjelaskan cara menjadwalkan pekerjaan Spark pada AnalyticDB for MySQL menggunakan Apache Airflow. Airflow mengoordinasikan beban kerja dalam bentuk directed acyclic graphs (DAGs). Anda dapat menghubungkan Airflow ke AnalyticDB for MySQL dengan dua metode:

Metode Paling cocok untuk
Spark Airflow Operator Integrasi lebih erat dengan AnalyticDB for MySQL; menggunakan autentikasi AccessKey; mendukung pekerjaan SQL dan JAR
spark-submit Provider Spark Apache Airflow standar; cocok jika Anda sudah menggunakan paket apache-airflow-providers-apache-spark

Prasyarat

Sebelum memulai, pastikan Anda telah memiliki:

Jadwalkan pekerjaan Spark SQL

AnalyticDB for MySQL mendukung Spark SQL dalam mode batch dan mode interaktif. Konfigurasinya berbeda antara keduanya.

Mode batch

Spark Airflow Operator

  1. Instal plugin Airflow Spark:

    pip install https://help-static-aliyun-doc.aliyuncs.com/file-manage-files/zh-CN/20230608/qvjf/adb_spark_airflow-0.0.1-py3-none-any.whl
  2. Buat koneksi Airflow. Di UI web Airflow, buka Admin > Connections dan tambahkan koneksi dengan JSON berikut sebagai connection extra:

    Penting

    Gunakan pengguna Resource Access Management (RAM) dengan izin minimum yang diperlukan. Jangan gunakan kredensial akun root Alibaba Cloud Anda.

    Parameter Deskripsi
    auth_type Metode autentikasi. Atur ke AK untuk menggunakan autentikasi pasangan AccessKey.
    access_key_id ID AccessKey pengguna RAM Anda yang memiliki akses ke AnalyticDB for MySQL.
    access_key_secret Rahasia AccessKey pengguna RAM Anda.
    region ID wilayah kluster AnalyticDB for MySQL.
    {
      "auth_type": "AK",
      "access_key_id": "<your_access_key_ID>",
      "access_key_secret": "<your_access_key_secret>",
      "region": "<your_region>"
    }
  3. Buat file DAG bernama spark_dags.py. Contoh berikut menggunakan AnalyticDBSparkSQLOperator untuk menjalankan kueri SHOW DATABASES:

    Parameter DAG:

    Parameter Wajib Deskripsi
    dag_id Ya Nama DAG.
    default_args Ya Nilai default tingkat kluster: cluster_id (ID kluster), rg_name (nama kelompok sumber daya pekerjaan), region (ID wilayah). Untuk informasi selengkapnya, lihat parameter DAG.

    Parameter AnalyticDBSparkSQLOperator:

    Parameter Wajib Deskripsi
    task_id Ya ID pekerjaan.
    sql Ya Pernyataan SQL Spark. Untuk informasi selengkapnya, lihat parameter Airflow.
    from datetime import datetime
    
    from airflow.models.dag import DAG
    from airflow_alibaba.providers.alibaba.cloud.operators.analyticdb_spark import AnalyticDBSparkSQLOperator
    
    with DAG(
        dag_id="my_dag_name",
        default_args={"cluster_id": "<your_cluster_ID>", "rg_name": "<your_resource_group>", "region": "<your_region>"},
    ) as dag:
    
        spark_sql = AnalyticDBSparkSQLOperator(
            task_id="task2",
            sql="SHOW DATABASES;"
        )
    
        spark_sql
  4. Salin spark_dags.py ke direktori dags_folder yang ditentukan dalam konfigurasi Airflow Anda.

  5. Picu DAG dari UI web Airflow. Untuk panduan, lihat tutorial Airflow.

spark-submit

Catatan Anda dapat mengatur parameter khusus AnalyticDB for MySQL (clusterId, regionId, keyId, secretId) di file conf/spark-defaults.conf atau sebagai parameter Airflow. Untuk daftar lengkapnya, lihat parameter konfigurasi aplikasi Spark.
  1. Instal plugin Airflow Spark:

    Penting

    Menginstal apache-airflow-providers-apache-spark secara otomatis menginstal PySpark. Untuk menghapus PySpark, jalankan pip3 uninstall pyspark.

    pip3 install apache-airflow-providers-apache-spark
  2. Unduh paket spark-submit dan konfigurasikan parameternya.

  3. Tambahkan biner spark-submit ke PATH Airflow sebelum memulai Airflow:

    Penting

    Atur PATH sebelum memulai Airflow. Jika Airflow dimulai tanpa path spark-submit, perintah tersebut tidak dapat ditemukan saat waktu proses.

    export PATH=$PATH:</your/adb/spark/path/bin>
  4. Buat file DAG bernama demo.py:

    from airflow.models import DAG
    from airflow.providers.apache.spark.operators.spark_sql import SparkSqlOperator
    from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
    from airflow.utils.dates import days_ago
    args = {
        'owner': 'Aliyun ADB Spark',
    }
    with DAG(
        dag_id='example_spark_operator',
        default_args=args,
        schedule_interval=None,
        start_date=days_ago(2),
        tags=['example'],
    ) as dag:
        adb_spark_conf = {
            "spark.driver.resourceSpec": "medium",
            "spark.executor.resourceSpec": "medium"
        }
        # Submit a Spark application from an OSS path
        submit_job = SparkSubmitOperator(
            conf=adb_spark_conf,
            application="oss://<bucket_name>/jar/pi.py",
            task_id="submit_job",
            verbose=True
        )
        # Run a Spark SQL query
        sql_job = SparkSqlOperator(
            conn_id="spark_default",
            sql="SELECT * FROM yourdb.yourtable",
            conf=",".join([k + "=" + v for k, v in adb_spark_conf.items()]),
            task_id="sql_job",
            verbose=True
        )
        submit_job >> sql_job
  5. Salin demo.py ke folder dags di direktori instalasi Airflow.

  6. Picu DAG dari UI web Airflow. Untuk panduan, lihat tutorial Airflow.

Mode interaktif

Mode interaktif menghubungkan Airflow ke kluster AnalyticDB for MySQL melalui titik akhir JDBC menggunakan HiveServer2 Thrift.

  1. Dapatkan titik akhir kelompok sumber daya Spark interaktif: Anda harus mengklik Apply for Endpoint di samping Public Endpoint untuk meminta titik akhir publik dalam kasus berikut:

    • Tool client yang digunakan untuk mengirim pekerjaan Spark SQL dideploy di mesin lokal atau server eksternal.

    • Tool client yang digunakan untuk mengirim pekerjaan Spark SQL dideploy di instance ECS, dan instance ECS serta kluster AnalyticDB for MySQL tidak berada dalam VPC yang sama.

    1. Masuk ke Konsol AnalyticDB for MySQL. Di pojok kiri atas, pilih wilayah. Di panel navigasi sebelah kiri, klik Clusters, lalu klik ID kluster Anda.

    2. Di panel navigasi, pilih Cluster Management > Resource Management, lalu klik tab Resource Groups.

    3. Temukan kelompok sumber daya target dan klik Details di kolom Actions. Salin titik akhir internal atau publik. Anda juga dapat menyalin string koneksi JDBC dari bidang Port. image

  2. Instal dependensi yang diperlukan:

    pip install apache-airflow-providers-apache-hive "apache-airflow-providers-common-sql==1.21.0"

    Untuk detail paket, lihat apache-airflow-providers-apache-hive dan apache-airflow-providers-common-sql.

  3. Di UI web Airflow, buka Admin > Connections.

  4. Klik image untuk menambahkan koneksi. Konfigurasikan parameter berikut:

    Parameter Deskripsi
    Connection Id Nama koneksi. Contoh: adb_spark_cluster.
    Connection Type Pilih Hive Server 2 Thrift.
    Host Titik akhir dari langkah 1. Ganti default dengan nama database aktual dan hapus sufiks resource_group=<resource group name>. Contoh: jdbc:hive2://amv-t4naxpqk****sparkwho.ads.aliyuncs.com:10000/adb_demo.
    Schema Nama database. Contoh: adb_demo.
    Login Nama kelompok sumber daya dan akun database dalam format resource_group_name/database_account_name. Contoh: spark_interactive_prod/spark_user.
    Password Password akun database AnalyticDB for MySQL.
    Port 10000
    Extra {"auth_mechanism": "CUSTOM"}
  5. Buat file DAG. Contoh berikut menjalankan show databases sesuai jadwal harian:

    Parameter Wajib Deskripsi
    task_id Ya ID pekerjaan.
    conn_id Ya Nama koneksi dari langkah 4.
    sql Ya Pernyataan SQL Spark.
    from airflow import DAG
    from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
    from datetime import datetime
    
    default_args = {
        'owner': 'airflow',
        'start_date': datetime(2025, 2, 10),
        'retries': 1,
    }
    
    dag = DAG(
        'adb_spark_sql_test',
        default_args=default_args,
        schedule_interval='@daily',
    )
    
    jdbc_query = SQLExecuteQueryOperator(
        task_id='execute_spark_sql_query',
        conn_id='adb_spark_cluster',  # Connection created in step 4
        sql='show databases',
        dag=dag
    )
    
    jdbc_query
  6. Di UI web Airflow, klik image di samping DAG untuk memicunya.

Jadwalkan pekerjaan Spark JAR

Spark Airflow Operator

  1. Instal plugin Airflow Spark:

    pip install https://help-static-aliyun-doc.aliyuncs.com/file-manage-files/zh-CN/20230608/qvjf/adb_spark_airflow-0.0.1-py3-none-any.whl
  2. Buat koneksi Airflow. Di UI web Airflow, buka Admin > Connections dan tambahkan koneksi dengan JSON berikut:

    Penting

    Gunakan pengguna RAM dengan izin minimum yang diperlukan. Jangan gunakan kredensial akun root Anda. Untuk detailnya, lihat Akun dan izin.

    Parameter Deskripsi
    auth_type Metode autentikasi. Atur ke AK.
    access_key_id ID AccessKey pengguna RAM Anda.
    access_key_secret Rahasia AccessKey pengguna RAM Anda.
    region ID wilayah kluster AnalyticDB for MySQL.
    {
      "auth_type": "AK",
      "access_key_id": "<your_access_key_ID>",
      "access_key_secret": "<your_access_key_secret>",
      "region": "<your_region>"
    }
  3. Buat file DAG bernama spark_dags.py. Contoh berikut menjalankan dua tugas berbasis JAR secara berurutan menggunakan AnalyticDBSparkBatchOperator:

    Penting

    Simpan semua file utama aplikasi Spark di Object Storage Service (OSS). Bucket OSS dan kluster AnalyticDB for MySQL harus berada dalam wilayah yang sama.

    Parameter DAG:

    Parameter Wajib Deskripsi
    dag_id Ya Nama DAG.
    default_args Ya Nilai default tingkat kluster: cluster_id, rg_name, region. Untuk informasi selengkapnya, lihat parameter DAG.

    Parameter AnalyticDBSparkBatchOperator:

    Parameter Wajib Deskripsi
    task_id Ya ID pekerjaan.
    file Ya Path absolut ke file utama aplikasi Spark — paket JAR (Java/Scala) atau file entry-point Python. Harus disimpan di OSS.
    class_name Wajib untuk Java/Scala Kelas titik masuk. Aplikasi Python tidak memerlukan ini. Untuk informasi selengkapnya, lihat parameter AnalyticDBSparkBatchOperator.
    from datetime import datetime
    
    from airflow.models.dag import DAG
    from airflow_alibaba.providers.alibaba.cloud.operators.analyticdb_spark import AnalyticDBSparkBatchOperator
    from airflow_alibaba.providers.alibaba.cloud.operators.analyticdb_spark import AnalyticDBSparkSQLOperator
    
    with DAG(
        dag_id=DAG_ID,
        default_args={"cluster_id": "your cluster", "rg_name": "your resource group", "region": "your region"},
    ) as dag:
        spark_pi = AnalyticDBSparkBatchOperator(
            task_id="task1",
            file="local:///tmp/spark-examples.jar",
            class_name="org.apache.spark.examples.SparkPi",
        )
    
        spark_lr = AnalyticDBSparkBatchOperator(
            task_id="task2",
            file="local:///tmp/spark-examples.jar",
            class_name="org.apache.spark.examples.SparkLR",
        )
    
        spark_pi >> spark_lr
  4. Salin spark_dags.py ke direktori dags_folder yang ditentukan dalam konfigurasi Airflow Anda.

  5. Picu DAG dari UI web Airflow. Untuk panduan, lihat tutorial Airflow.

spark-submit

Catatan Anda dapat mengatur parameter khusus AnalyticDB for MySQL (clusterId, regionId, keyId, secretId, ossUploadPath) di file conf/spark-defaults.conf atau sebagai parameter Airflow. Untuk daftar lengkapnya, lihat parameter konfigurasi aplikasi Spark.
  1. Instal plugin Airflow Spark:

    Penting

    Menginstal apache-airflow-providers-apache-spark secara otomatis menginstal PySpark. Untuk menghapus PySpark, jalankan pip3 uninstall pyspark.

    pip3 install apache-airflow-providers-apache-spark
  2. Unduh paket spark-submit dan konfigurasikan parameternya.

  3. Tambahkan biner spark-submit ke PATH Airflow sebelum memulai Airflow:

    Penting

    Atur PATH sebelum memulai Airflow. Jika Airflow dimulai tanpa path spark-submit, perintah tersebut tidak dapat ditemukan saat waktu proses.

    export PATH=$PATH:</your/adb/spark/path/bin>
  4. Buat file DAG bernama demo.py. Contoh berikut menjalankan dua tugas JAR secara berurutan:

    Parameter DAG:

    Parameter Wajib Deskripsi
    dag_id Ya Nama DAG.
    default_args Ya Nilai default tingkat kluster: cluster_id, rg_name, region. Untuk informasi selengkapnya, lihat parameter DAG.

    Parameter AnalyticDBSparkBatchOperator:

    Parameter Wajib Deskripsi
    task_id Ya ID pekerjaan.
    file Ya Path absolut ke file utama aplikasi Spark. Harus disimpan di OSS.
    class_name Wajib untuk Java/Scala Kelas titik masuk. Aplikasi Python tidak memerlukan ini. Untuk informasi selengkapnya, lihat parameter AnalyticDBSparkBatchOperator.
    from datetime import datetime
    
    from airflow.models.dag import DAG
    from airflow_alibaba.providers.alibaba.cloud.operators.analyticdb_spark import AnalyticDBSparkBatchOperator
    from airflow_alibaba.providers.alibaba.cloud.operators.analyticdb_spark import AnalyticDBSparkSQLOperator
    
    with DAG(
        dag_id=DAG_ID,
        start_date=datetime(2021, 1, 1),
        schedule=None,
        default_args={"cluster_id": "your cluster", "rg_name": "your resource group", "region": "your region"},
        max_active_runs=1,
        catchup=False,
    ) as dag:
        spark_pi = AnalyticDBSparkBatchOperator(
            task_id="task1",
            file="local:///tmp/spark-examples.jar",
            class_name="org.apache.spark.examples.SparkPi",
        )
    
        spark_lr = AnalyticDBSparkBatchOperator(
            task_id="task2",
            file="local:///tmp/spark-examples.jar",
            class_name="org.apache.spark.examples.SparkLR",
        )
    
        spark_pi >> spark_lr
  5. Salin demo.py ke folder dags di direktori instalasi Airflow.

  6. Picu DAG dari UI web Airflow. Untuk panduan, lihat tutorial Airflow.