すべてのプロダクト
Search
ドキュメントセンター

AnalyticDB:Airflow を使用した Spark ジョブのスケジューリング

最終更新日:Sep 17, 2026

このガイドでは、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 パッケージを使用している場合に適しています

前提条件

開始する前に、以下をご確認ください。

Spark SQL ジョブのスケジューリング

AnalyticDB for MySQL は、バッチモードインタラクティブモードでの Spark SQL をサポートしています。設定はモードによって異なります。

バッチモード

Spark Airflow オペレーター

  1. 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. Airflow コネクションを作成します。Airflow Web UI で、[Admin] > [Connections] に移動し、コネクションの extra として次の JSON を使用してコネクションを追加します。

    重要

    最小権限の Resource Access Management (RAM) ユーザーを使用してください。Alibaba Cloud のルートアカウントの認証情報は使用しないでください。

    パラメーター 説明
    auth_type 認証方法。アクセスキーペア認証を使用するには AK に設定します。
    access_key_id AnalyticDB for MySQL へのアクセス権を持つ RAM ユーザーのアクセスキー ID。
    access_key_secret RAM ユーザーのアクセスキーシークレット。
    region AnalyticDB for MySQL クラスターのリージョン ID。
    {
      "auth_type": "AK",
      "access_key_id": "<your_access_key_ID>",
      "access_key_secret": "<your_access_key_secret>",
      "region": "<your_region>"
    }
  3. 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
  4. spark_dags.py を Airflow の設定で定義された dags_folder ディレクトリにコピーします。

  5. Airflow Web UI から DAG をトリガーします。手順については、「Airflow チュートリアル」をご参照ください。

spark-submit

説明 AnalyticDB for MySQL 固有のパラメーター (clusterIdregionIdkeyIdsecretId) は、conf/spark-defaults.conf ファイルまたは Airflow パラメーターとして設定できます。完全なリストについては、「Spark アプリケーション設定パラメーター」をご参照ください。
  1. Airflow Spark プラグインをインストールします。

    重要

    apache-airflow-providers-apache-spark をインストールすると、PySpark が自動的にインストールされます。PySpark を削除するには、pip3 uninstall pyspark を実行します。

    pip3 install apache-airflow-providers-apache-spark
  2. spark-submit パッケージをダウンロードしてパラメーターを設定します

  3. Airflow を起動する前に、spark-submit バイナリを Airflow の PATH に追加します。

    重要

    PATH は Airflow を起動する前に設定してください。spark-submit のパスなしで Airflow を起動した場合、実行時にコマンドが見つかりません。

    export PATH=$PATH:</your/adb/spark/path/bin>
  4. 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
  5. demo.py を Airflow インストールディレクトリの dags フォルダーにコピーします。

  6. Airflow Web UI から DAG をトリガーします。手順については、「Airflow チュートリアル」をご参照ください。

インタラクティブモード

インタラクティブモードでは、HiveServer2 Thrift を使用して、JDBC エンドポイント経由で Airflow を AnalyticDB for MySQL クラスターに接続します。

  1. Spark インタラクティブリソースグループのエンドポイントを取得します。次の場合には、[Public Endpoint] の横にある [Apply for Endpoint] をクリックしてパブリックエンドポイントをリクエストする必要があります。

    • Spark SQL ジョブを送信するために使用するクライアントツールがローカルマシンまたは外部サーバーにデプロイされている場合。

    • Spark SQL ジョブを送信するために使用するクライアントツールが ECS インスタンスにデプロイされており、ECS インスタンスと AnalyticDB for MySQL クラスターが同じ VPC にない場合。

    1. AnalyticDB for MySQL コンソールにログインします。画面左上でリージョンを選択します。左側メニューで、[Clusters] をクリックし、クラスター ID をクリックします。

    2. ナビゲーションペインで、[Cluster Management] > [Resource Management] を選択し、[Resource Groups] タブをクリックします。

    3. 対象のリソースグループを見つけ、[Actions] 列の [Details] をクリックします。内部またはパブリックエンドポイントをコピーします。[Port] フィールドから JDBC 接続文字列をコピーすることもできます。image

  2. 必要な依存関係をインストールします。

    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」をご参照ください。

  3. Airflow Web UI で、[Admin] > [Connections] に移動します。

  4. image をクリックしてコネクションを追加します。次のパラメーターを設定します。

    パラメーター 説明
    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 10000
    Extra {"auth_mechanism": "CUSTOM"}
  5. 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
  6. Airflow Web UI で、DAG の横にあるimageをクリックしてトリガーします。

Spark JAR ジョブのスケジューリング

Spark Airflow オペレーター

  1. 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. Airflow コネクションを作成します。Airflow Web UI で、[Admin] > [Connections] に移動し、次の JSON を使用してコネクションを追加します。

    重要

    最小権限の RAM ユーザーを使用してください。ルートアカウントの認証情報は使用しないでください。詳細については、「アカウントと権限」をご参照ください。

    パラメーター 説明
    auth_type 認証方法。AK に設定します。
    access_key_id RAM ユーザーのアクセスキー ID。
    access_key_secret RAM ユーザーのアクセスキーシークレット。
    region AnalyticDB for MySQL クラスターのリージョン ID。
    {
      "auth_type": "AK",
      "access_key_id": "<your_access_key_ID>",
      "access_key_secret": "<your_access_key_secret>",
      "region": "<your_region>"
    }
  3. 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_name Java/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
  4. spark_dags.py を Airflow の設定で定義された dags_folder ディレクトリにコピーします。

  5. Airflow Web UI から DAG をトリガーします。手順については、「Airflow チュートリアル」をご参照ください。

spark-submit

説明 AnalyticDB for MySQL 固有のパラメーター (clusterIdregionIdkeyIdsecretIdossUploadPath) は、conf/spark-defaults.conf ファイルまたは Airflow パラメーターとして設定できます。完全なリストについては、「Spark アプリケーション設定パラメーター」をご参照ください。
  1. Airflow Spark プラグインをインストールします。

    重要

    apache-airflow-providers-apache-spark をインストールすると、PySpark が自動的にインストールされます。PySpark を削除するには、pip3 uninstall pyspark を実行します。

    pip3 install apache-airflow-providers-apache-spark
  2. spark-submit パッケージをダウンロードしてパラメーターを設定します

  3. Airflow を起動する前に、spark-submit バイナリを Airflow の PATH に追加します。

    重要

    PATH は Airflow を起動する前に設定してください。spark-submit のパスなしで Airflow を起動した場合、実行時にコマンドが見つかりません。

    export PATH=$PATH:</your/adb/spark/path/bin>
  4. demo.py という名前の DAG ファイルを作成します。次の例では、2 つの JAR タスクを順番に実行します:

    DAG パラメーター:

    Parameter Required Description
    dag_id Yes The DAG name.
    default_args Yes Cluster-level defaults: cluster_id, rg_name, region. For more information, see DAG parameters.

    AnalyticDBSparkBatchOperator パラメーター:

    Parameter Required Description
    task_id Yes The job ID.
    file Yes Absolute path to the Spark application's main file. Must be stored in OSS.
    class_name Required 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
  5. demo_jar.py を Airflow インストールディレクトリの dags フォルダーにコピーします。

  6. Airflow Web UI から DAG をトリガーします。手順については、「Airflow チュートリアル」をご参照ください。