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:
-
Kluster AnalyticDB for MySQL Edisi Perusahaan, Edisi Dasar, atau Edisi Data Lakehouse
-
kelompok sumber daya pekerjaan atau kelompok sumber daya Spark interaktif yang dibuat untuk kluster tersebut
-
Python 3.7 atau versi yang lebih baru
-
Alamat IP server Airflow ditambahkan ke daftar putih kluster
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
-
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 -
Buat koneksi Airflow. Di UI web Airflow, buka Admin > Connections dan tambahkan koneksi dengan JSON berikut sebagai connection extra:
PentingGunakan pengguna Resource Access Management (RAM) dengan izin minimum yang diperlukan. Jangan gunakan kredensial akun root Alibaba Cloud Anda.
Parameter Deskripsi auth_typeMetode autentikasi. Atur ke AKuntuk menggunakan autentikasi pasangan AccessKey.access_key_idID AccessKey pengguna RAM Anda yang memiliki akses ke AnalyticDB for MySQL. access_key_secretRahasia AccessKey pengguna RAM Anda. regionID 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>" } -
Buat file DAG bernama
spark_dags.py. Contoh berikut menggunakanAnalyticDBSparkSQLOperatoruntuk menjalankan kueriSHOW DATABASES:Parameter DAG:
Parameter Wajib Deskripsi dag_idYa Nama DAG. default_argsYa 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_idYa ID pekerjaan. sqlYa 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 -
Salin
spark_dags.pyke direktoridags_folderyang ditentukan dalam konfigurasi Airflow Anda. -
Picu DAG dari UI web Airflow. Untuk panduan, lihat tutorial Airflow.
spark-submit
clusterId, regionId, keyId, secretId) di file conf/spark-defaults.conf atau sebagai parameter Airflow. Untuk daftar lengkapnya, lihat parameter konfigurasi aplikasi Spark.-
Instal plugin Airflow Spark:
PentingMenginstal
apache-airflow-providers-apache-sparksecara otomatis menginstal PySpark. Untuk menghapus PySpark, jalankanpip3 uninstall pyspark.pip3 install apache-airflow-providers-apache-spark -
Tambahkan biner spark-submit ke
PATHAirflow sebelum memulai Airflow:PentingAtur
PATHsebelum 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> -
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 -
Salin
demo.pyke folderdagsdi direktori instalasi Airflow. -
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.
-
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.
-
Masuk ke Konsol AnalyticDB for MySQL. Di pojok kiri atas, pilih wilayah. Di panel navigasi sebelah kiri, klik Clusters, lalu klik ID kluster Anda.
-
Di panel navigasi, pilih Cluster Management > Resource Management, lalu klik tab Resource Groups.
-
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.

-
-
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.
-
Di UI web Airflow, buka Admin > Connections.
-
Klik
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 defaultdengan nama database aktual dan hapus sufiksresource_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 10000Extra {"auth_mechanism": "CUSTOM"} -
Buat file DAG. Contoh berikut menjalankan
show databasessesuai jadwal harian:Parameter Wajib Deskripsi task_idYa ID pekerjaan. conn_idYa Nama koneksi dari langkah 4. sqlYa 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 -
Di UI web Airflow, klik
di samping DAG untuk memicunya.
Jadwalkan pekerjaan Spark JAR
Spark Airflow Operator
-
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 -
Buat koneksi Airflow. Di UI web Airflow, buka Admin > Connections dan tambahkan koneksi dengan JSON berikut:
PentingGunakan pengguna RAM dengan izin minimum yang diperlukan. Jangan gunakan kredensial akun root Anda. Untuk detailnya, lihat Akun dan izin.
Parameter Deskripsi auth_typeMetode autentikasi. Atur ke AK.access_key_idID AccessKey pengguna RAM Anda. access_key_secretRahasia AccessKey pengguna RAM Anda. regionID 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>" } -
Buat file DAG bernama
spark_dags.py. Contoh berikut menjalankan dua tugas berbasis JAR secara berurutan menggunakanAnalyticDBSparkBatchOperator:PentingSimpan 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_idYa Nama DAG. default_argsYa Nilai default tingkat kluster: cluster_id,rg_name,region. Untuk informasi selengkapnya, lihat parameter DAG.Parameter AnalyticDBSparkBatchOperator:
Parameter Wajib Deskripsi task_idYa ID pekerjaan. fileYa Path absolut ke file utama aplikasi Spark — paket JAR (Java/Scala) atau file entry-point Python. Harus disimpan di OSS. class_nameWajib 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 -
Salin
spark_dags.pyke direktoridags_folderyang ditentukan dalam konfigurasi Airflow Anda. -
Picu DAG dari UI web Airflow. Untuk panduan, lihat tutorial Airflow.
spark-submit
clusterId, regionId, keyId, secretId, ossUploadPath) di file conf/spark-defaults.conf atau sebagai parameter Airflow. Untuk daftar lengkapnya, lihat parameter konfigurasi aplikasi Spark.-
Instal plugin Airflow Spark:
PentingMenginstal
apache-airflow-providers-apache-sparksecara otomatis menginstal PySpark. Untuk menghapus PySpark, jalankanpip3 uninstall pyspark.pip3 install apache-airflow-providers-apache-spark -
Tambahkan biner spark-submit ke
PATHAirflow sebelum memulai Airflow:PentingAtur
PATHsebelum 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> -
Buat file DAG bernama
demo.py. Contoh berikut menjalankan dua tugas JAR secara berurutan:Parameter DAG:
Parameter Wajib Deskripsi dag_idYa Nama DAG. default_argsYa Nilai default tingkat kluster: cluster_id,rg_name,region. Untuk informasi selengkapnya, lihat parameter DAG.Parameter AnalyticDBSparkBatchOperator:
Parameter Wajib Deskripsi task_idYa ID pekerjaan. fileYa Path absolut ke file utama aplikasi Spark. Harus disimpan di OSS. class_nameWajib 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 -
Salin
demo.pyke folderdagsdi direktori instalasi Airflow. -
Picu DAG dari UI web Airflow. Untuk panduan, lihat tutorial Airflow.