このガイドでは、Apache Airflow を使用して AnalyticDB for MySQL Spark ジョブをスケジュールする方法について説明します。Airflow は、有向非巡回グラフ (DAG) としてワークロードをオーケストレーションします。Airflow を AnalyticDB for MySQL に接続するには、次の 2 つの方法があります。
| 方法 | 推奨されるケース |
|---|---|
| Spark Airflow オペレーター | AnalyticDB for MySQL との緊密な統合、アクセスキー認証の利用、SQL ジョブと JAR ジョブの両方に対応 |
| spark-submit | 標準の Apache Airflow Spark プロバイダー。すでに apache-airflow-providers-apache-spark パッケージを使用している場合に適しています |
前提条件
開始する前に、以下をご確認ください。
-
AnalyticDB for MySQL Enterprise Edition、Basic Edition、または Data Lakehouse Edition のクラスター
-
クラスター用に作成されたジョブリソースグループまたは Spark インタラクティブリソースグループ
-
Python 3.7 以降
-
Airflow サーバーの IP アドレスがクラスターのホワイトリストに追加されていること
Spark SQL ジョブのスケジューリング
AnalyticDB for MySQL は、バッチモードとインタラクティブモードでの Spark SQL をサポートしています。設定はモードによって異なります。
バッチモード
Spark Airflow オペレーター
-
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 -
Airflow コネクションを作成します。Airflow Web UI で、[Admin] > [Connections] に移動し、コネクションの extra として次の JSON を使用してコネクションを追加します。
重要最小権限の Resource Access Management (RAM) ユーザーを使用してください。Alibaba Cloud のルートアカウントの認証情報は使用しないでください。
パラメーター 説明 auth_type認証方法。アクセスキーペア認証を使用するには AKに設定します。access_key_idAnalyticDB for MySQL へのアクセス権を持つ RAM ユーザーのアクセスキー ID。 access_key_secretRAM ユーザーのアクセスキーシークレット。 regionAnalyticDB for MySQL クラスターのリージョン ID。 { "auth_type": "AK", "access_key_id": "<your_access_key_ID>", "access_key_secret": "<your_access_key_secret>", "region": "<your_region>" } -
spark_dags.pyという名前の DAG ファイルを作成します。次の例では、AnalyticDBSparkSQLOperatorを使用してSHOW DATABASESクエリを実行します。DAG パラメーター:
パラメーター 必須 説明 dag_idはい DAG 名。 default_argsはい クラスターレベルのデフォルト値: cluster_id(クラスター ID)、rg_name(ジョブリソースグループ名)、region(リージョン ID)。詳細については、「DAG パラメーター」をご参照ください。AnalyticDBSparkSQLOperator パラメーター:
パラメーター 必須 説明 task_idはい ジョブ ID。 sqlはい Spark SQL ステートメント。詳細については、「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 -
spark_dags.pyを Airflow の設定で定義されたdags_folderディレクトリにコピーします。 -
Airflow Web UI から DAG をトリガーします。手順については、「Airflow チュートリアル」をご参照ください。
spark-submit
clusterId、regionId、keyId、secretId) は、conf/spark-defaults.conf ファイルまたは Airflow パラメーターとして設定できます。完全なリストについては、「Spark アプリケーション設定パラメーター」をご参照ください。-
Airflow Spark プラグインをインストールします。
重要apache-airflow-providers-apache-sparkをインストールすると、PySpark が自動的にインストールされます。PySpark を削除するには、pip3 uninstall pysparkを実行します。pip3 install apache-airflow-providers-apache-spark -
Airflow を起動する前に、spark-submit バイナリを Airflow の
PATHに追加します。重要PATHは Airflow を起動する前に設定してください。spark-submit のパスなしで Airflow を起動した場合、実行時にコマンドが見つかりません。export PATH=$PATH:</your/adb/spark/path/bin> -
demo.pyという名前の DAG ファイルを作成します。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" } # OSS パスから Spark アプリケーションを送信します submit_job = SparkSubmitOperator( conf=adb_spark_conf, application="oss://<bucket_name>/jar/pi.py", task_id="submit_job", verbose=True ) # Spark SQL クエリを実行します 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 -
demo.pyを Airflow インストールディレクトリのdagsフォルダーにコピーします。 -
Airflow Web UI から DAG をトリガーします。手順については、「Airflow チュートリアル」をご参照ください。
インタラクティブモード
インタラクティブモードでは、HiveServer2 Thrift を使用して、JDBC エンドポイント経由で Airflow を AnalyticDB for MySQL クラスターに接続します。
-
Spark インタラクティブリソースグループのエンドポイントを取得します。次の場合には、[Public Endpoint] の横にある [Apply for Endpoint] をクリックしてパブリックエンドポイントをリクエストする必要があります。
-
Spark SQL ジョブを送信するために使用するクライアントツールがローカルマシンまたは外部サーバーにデプロイされている場合。
-
Spark SQL ジョブを送信するために使用するクライアントツールが ECS インスタンスにデプロイされており、ECS インスタンスと AnalyticDB for MySQL クラスターが同じ VPC にない場合。
-
AnalyticDB for MySQL コンソールにログインします。画面左上でリージョンを選択します。左側メニューで、[Clusters] をクリックし、クラスター ID をクリックします。
-
ナビゲーションペインで、[Cluster Management] > [Resource Management] を選択し、[Resource Groups] タブをクリックします。
-
対象のリソースグループを見つけ、[Actions] 列の [Details] をクリックします。内部またはパブリックエンドポイントをコピーします。[Port] フィールドから JDBC 接続文字列をコピーすることもできます。

-
-
必要な依存関係をインストールします。
pip install apache-airflow-providers-apache-hive "apache-airflow-providers-common-sql==1.21.0"パッケージの詳細については、「apache-airflow-providers-apache-hive」および「apache-airflow-providers-common-sql」をご参照ください。
-
Airflow Web UI で、[Admin] > [Connections] に移動します。
-
をクリックしてコネクションを追加します。次のパラメーターを設定します。パラメーター 説明 Connection Id コネクション名。例: adb_spark_cluster。Connection Type [Hive Server 2 Thrift] を選択します。 Host ステップ 1 で取得したエンドポイント。 defaultを実際のデータベース名に置き換え、resource_group=<resource group name>サフィックスを削除します。例:jdbc:hive2://amv-t4naxpqk****sparkwho.ads.aliyuncs.com:10000/adb_demo。Schema データベース名。例: adb_demo。Login リソースグループ名とデータベースアカウントを resource_group_name/database_account_name形式で指定します。例:spark_interactive_prod/spark_user。Password AnalyticDB for MySQL データベースアカウントのパスワード。 Port 10000Extra {"auth_mechanism": "CUSTOM"} -
DAG ファイルを作成します。次の例では、
show databasesを日次スケジュールで実行します。パラメーター 必須 説明 task_idはい ジョブ ID。 conn_idはい ステップ 4 で作成したコネクション名。 sqlはい Spark SQL ステートメント。 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', # ステップ 4 で作成したコネクション sql='show databases', dag=dag ) jdbc_query -
Airflow Web UI で、DAG の横にある
をクリックしてトリガーします。
Spark JAR ジョブのスケジューリング
Spark Airflow オペレーター
-
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 -
Airflow コネクションを作成します。Airflow Web UI で、[Admin] > [Connections] に移動し、次の JSON を使用してコネクションを追加します。
重要最小権限の RAM ユーザーを使用してください。ルートアカウントの認証情報は使用しないでください。詳細については、「アカウントと権限」をご参照ください。
パラメーター 説明 auth_type認証方法。 AKに設定します。access_key_idRAM ユーザーのアクセスキー ID。 access_key_secretRAM ユーザーのアクセスキーシークレット。 regionAnalyticDB for MySQL クラスターのリージョン ID。 { "auth_type": "AK", "access_key_id": "<your_access_key_ID>", "access_key_secret": "<your_access_key_secret>", "region": "<your_region>" } -
spark_dags.pyという名前の DAG ファイルを作成します。次の例では、AnalyticDBSparkBatchOperatorを使用して、2 つの JAR ベースのタスクを順番に実行します。重要すべての Spark アプリケーションのメインファイルは Object Storage Service (OSS) に保存してください。OSS バケットと AnalyticDB for MySQL クラスターは同じリージョンにある必要があります。
DAG パラメーター:
パラメーター 必須 説明 dag_idはい DAG 名。 default_argsはい クラスターレベルのデフォルト値: cluster_id(クラスター ID)、rg_name(ジョブリソースグループ名)、region(リージョン ID)。詳細については、「DAG パラメーター」をご参照ください。AnalyticDBSparkBatchOperator パラメーター:
パラメーター 必須 説明 task_idはい ジョブ ID。 fileはい Spark アプリケーションのメインファイルの絶対パス — JAR パッケージ (Java/Scala) または Python エントリポイントファイル。OSS に保存する必要があります。 class_nameJava/Scala の場合は必須 エントリポイントクラス。Python アプリケーションでは不要です。詳細については、「AnalyticDBSparkBatchOperator パラメーター」をご参照ください。 from datetime import datetime from airflow.models.dag import DAG from airflow_alibaba.providers.alibaba.cloud.operators.analyticdb_spark import AnalyticDBSparkBatchOperator with DAG( dag_id="my_jar_dag", default_args={"cluster_id": "<your_cluster_ID>", "rg_name": "<your_resource_group>", "region": "<your_region>"}, start_date=datetime(2023, 1, 1), schedule_interval=None, catchup=False, ) as dag: spark_pi = AnalyticDBSparkBatchOperator( task_id="task1", file="oss://<your-bucket-name>/path/to/spark-examples.jar", class_name="org.apache.spark.examples.SparkPi", ) spark_lr = AnalyticDBSparkBatchOperator( task_id="task2", file="oss://<your-bucket-name>/path/to/spark-examples.jar", class_name="org.apache.spark.examples.SparkLR", ) spark_pi >> spark_lr -
spark_dags.pyを Airflow の設定で定義されたdags_folderディレクトリにコピーします。 -
Airflow Web UI から DAG をトリガーします。手順については、「Airflow チュートリアル」をご参照ください。
spark-submit
clusterId、regionId、keyId、secretId、ossUploadPath) は、conf/spark-defaults.conf ファイルまたは Airflow パラメーターとして設定できます。完全なリストについては、「Spark アプリケーション設定パラメーター」をご参照ください。-
Airflow Spark プラグインをインストールします。
重要apache-airflow-providers-apache-sparkをインストールすると、PySpark が自動的にインストールされます。PySpark を削除するには、pip3 uninstall pysparkを実行します。pip3 install apache-airflow-providers-apache-spark -
Airflow を起動する前に、spark-submit バイナリを Airflow の
PATHに追加します。重要PATHは Airflow を起動する前に設定してください。spark-submit のパスなしで Airflow を起動した場合、実行時にコマンドが見つかりません。export PATH=$PATH:</your/adb/spark/path/bin> -
demo.pyという名前の DAG ファイルを作成します。次の例では、2 つの JAR タスクを順番に実行します:DAG パラメーター:
Parameter Required Description dag_idYes The DAG name. default_argsYes Cluster-level defaults: cluster_id,rg_name,region. For more information, see DAG parameters.AnalyticDBSparkBatchOperator パラメーター:
Parameter Required Description task_idYes The job ID. fileYes Absolute path to the Spark application's main file. Must be stored in OSS. class_nameRequired for Java/Scala The entry point class. Python applications do not require this. For more information, see AnalyticDBSparkBatchOperator parameters. from airflow.models.dag import DAG 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_submit_jar_operator', default_args=args, schedule_interval=None, start_date=days_ago(2), tags=['example'], ) as dag: # Spark 設定 adb_spark_conf = { "spark.driver.resourceSpec": "medium", "spark.executor.resourceSpec": "medium" } # OSS パスから Spark JAR アプリケーションを送信 submit_jar_task = SparkSubmitOperator( task_id='submit_spark_pi_jar', application='oss://<your-bucket-name>/path/to/spark-examples.jar', # JAR ファイルの OSS パス java_class='org.apache.spark.examples.SparkPi', # メインクラス application_args=['10'], # メインクラスへの引数 conf=adb_spark_conf, conn_id='spark_default', # Airflow の Spark コネクション ID verbose=True ) submit_jar_task -
demo_jar.pyを Airflow インストールディレクトリのdagsフォルダーにコピーします。 -
Airflow Web UI から DAG をトリガーします。手順については、「Airflow チュートリアル」をご参照ください。