Tous les produits
Search
Centre de documentation

AnalyticDB:Utiliser Airflow pour planifier des tâches Spark

Dernière mise à jour :Aug 10, 2026

Ce guide explique comment planifier des tâches Spark AnalyticDB for MySQL à l'aide d'Apache Airflow. Airflow orchestre les charges de travail sous forme de graphes acycliques dirigés (DAG). Deux méthodes permettent de connecter Airflow à AnalyticDB for MySQL :

Méthode Idéal pour
Spark Airflow Operator Intégration étroite avec AnalyticDB for MySQL ; utilise l'authentification par AccessKey ; prend en charge les tâches SQL et JAR
spark-submit Fournisseur Apache Airflow Spark standard ; adapté si vous utilisez déjà le package apache-airflow-providers-apache-spark

Prérequis

Avant de commencer, vérifiez que vous disposez des éléments suivants :

Planifier des tâches Spark SQL

AnalyticDB for MySQL prend en charge Spark SQL en mode batch et en mode interactif. La configuration varie selon le mode choisi.

Mode batch

Spark Airflow Operator

  1. Installez le plug-in 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. Créez une connexion Airflow. Dans l'interface web d'Airflow, accédez à Admin > Connections et ajoutez une connexion avec le code JSON suivant comme paramètre supplémentaire :

    Important

    Utilisez un utilisateur Resource Access Management (RAM) disposant des autorisations minimales requises. N'utilisez pas les identifiants de votre compte racine Alibaba Cloud.

    Paramètre Description
    auth_type Méthode d'authentification. Définissez la valeur AK pour utiliser l'authentification par paire AccessKey.
    access_key_id L'AccessKey ID de votre utilisateur RAM ayant accès à AnalyticDB for MySQL.
    access_key_secret L'AccessKey Secret de votre utilisateur RAM.
    region L'ID de région du cluster AnalyticDB for MySQL.
    {
      "auth_type": "AK",
      "access_key_id": "<your_access_key_ID>",
      "access_key_secret": "<your_access_key_secret>",
      "region": "<your_region>"
    }
  3. Créez un fichier DAG nommé spark_dags.py. L'exemple suivant utilise AnalyticDBSparkSQLOperator pour exécuter une requête SHOW DATABASES :

    Paramètres du DAG :

    Paramètre Obligatoire Description
    dag_id Oui Nom du DAG.
    default_args Oui Valeurs par défaut au niveau du cluster : cluster_id (ID du cluster), rg_name (nom du groupe de ressources de tâches), region (ID de région). Pour plus d'informations, consultez les paramètres DAG.

    Paramètres AnalyticDBSparkSQLOperator :

    Paramètre Obligatoire Description
    task_id Oui ID de la tâche.
    sql Oui Instruction Spark SQL. Pour plus d'informations, consultez les paramètres 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. Copiez spark_dags.py dans le répertoire dags_folder défini dans votre configuration Airflow.

  5. Déclenchez le DAG depuis l'interface web d'Airflow. Pour obtenir des instructions, consultez le tutoriel Airflow.

spark-submit

Remarque

Vous pouvez définir les paramètres spécifiques à AnalyticDB for MySQL (

clusterId

,

regionId

,

keyId

,

secretId

) dans le fichier

conf/spark-defaults.conf

ou en tant que paramètres Airflow. Pour la liste complète, consultez les

paramètres de configuration des applications Spark

.

  1. Installez le plug-in Airflow Spark :

    Important

    L'installation de apache-airflow-providers-apache-spark installe automatiquement PySpark. Pour supprimer PySpark, exécutez pip3 uninstall pyspark.

    pip3 install apache-airflow-providers-apache-spark
  2. Téléchargez le package spark-submit et configurez les paramètres.

  3. Ajoutez le binaire spark-submit au PATH d'Airflow avant de démarrer Airflow :

    Important

    Définissez le PATH avant de démarrer Airflow. Si Airflow démarre sans le chemin d'accès à spark-submit, il ne pourra pas localiser la commande lors de l'exécution.

    export PATH=$PATH:</your/adb/spark/path/bin>
  4. Créez un fichier DAG nommé 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. Copiez demo.py dans le dossier dags du répertoire d'installation d'Airflow.

  6. Déclenchez le DAG depuis l'interface web d'Airflow. Pour obtenir des instructions, consultez le tutoriel Airflow.

Mode interactif

Le mode interactif connecte Airflow au cluster AnalyticDB for MySQL via un endpoint JDBC utilisant HiveServer2 Thrift.

  1. Obtenez l'endpoint du groupe de ressources interactives Spark : cliquez sur Apply for Endpoint à côté de Public Endpoint pour demander un endpoint public dans les cas suivants :

    • L'outil client utilisé pour soumettre des tâches Spark SQL est déployé sur une machine locale ou un serveur externe.

    • L'outil client utilisé pour soumettre des tâches Spark SQL est déployé sur une instance ECS, et l'instance ECS et le cluster AnalyticDB for MySQL ne se trouvent pas dans le même VPC.

    1. Connectez-vous à la console AnalyticDB for MySQL. Dans le coin supérieur gauche, sélectionnez une région. Dans le volet de navigation de gauche, cliquez sur Clusters, puis cliquez sur l'ID de votre cluster.

    2. Dans le volet de navigation, choisissez Cluster Management > Resource Management, puis cliquez sur l'onglet Resource Groups.

    3. Localisez le groupe de ressources cible et cliquez sur Details dans la colonne Actions. Copiez l'endpoint interne ou public. Vous pouvez également copier la chaîne de connexion JDBC depuis le champ Port. image

  2. Installez les dépendances requises :

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

    Pour plus de détails sur les packages, consultez apache-airflow-providers-apache-hive et apache-airflow-providers-common-sql.

  3. Dans l'interface web d'Airflow, accédez à Admin > Connections.

  4. Cliquez sur image pour ajouter une connexion. Configurez les paramètres suivants :

    Paramètre Description
    Connection Id Nom de la connexion. Exemple : adb_spark_cluster.
    Connection Type Sélectionnez Hive Server 2 Thrift.
    Host Endpoint obtenu à l'étape 1. Remplacez default par le nom réel de la base de données et supprimez le suffixe resource_group=<resource group name>. Exemple : jdbc:hive2://amv-t4naxpqk****sparkwho.ads.aliyuncs.com:10000/adb_demo.
    Schema Nom de la base de données. Exemple : adb_demo.
    Login Nom du groupe de ressources et compte de base de données au format resource_group_name/database_account_name. Exemple : spark_interactive_prod/spark_user.
    Password Mot de passe du compte de base de données AnalyticDB for MySQL.
    Port 10000
    Extra {"auth_mechanism": "CUSTOM"}
  5. Créez un fichier DAG. L'exemple suivant exécute show databases selon une planification quotidienne :

    Paramètre Obligatoire Description
    task_id Oui ID de la tâche.
    conn_id Oui Nom de la connexion défini à l'étape 4.
    sql Oui Instruction 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',  # Connection created in step 4
        sql='show databases',
        dag=dag
    )
    
    jdbc_query
  6. Dans l'interface web d'Airflow, cliquez sur image à côté du DAG pour le déclencher.

Planifier des tâches Spark JAR

Spark Airflow Operator

  1. Installez le plug-in 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. Créez une connexion Airflow. Dans l'interface web d'Airflow, accédez à Admin > Connections et ajoutez une connexion avec le code JSON suivant :

    Important

    Utilisez un utilisateur RAM disposant des autorisations minimales requises. N'utilisez pas les identifiants de votre compte racine. Pour plus de détails, consultez les Comptes et autorisations.

    Paramètre Description
    auth_type Méthode d'authentification. Définissez la valeur AK.
    access_key_id AccessKey ID de votre utilisateur RAM.
    access_key_secret AccessKey Secret de votre utilisateur RAM.
    region ID de région du cluster AnalyticDB for MySQL.
    {
      "auth_type": "AK",
      "access_key_id": "<your_access_key_ID>",
      "access_key_secret": "<your_access_key_secret>",
      "region": "<your_region>"
    }
  3. Créez un fichier DAG nommé spark_dags.py. L'exemple suivant exécute deux tâches basées sur des fichiers JAR de manière séquentielle à l'aide de AnalyticDBSparkBatchOperator :

    Important

    Stockez tous les fichiers principaux des applications Spark dans Object Storage Service (OSS). Le bucket OSS et le cluster AnalyticDB for MySQL doivent se trouver dans la même région.

    Paramètres du DAG :

    Paramètre Obligatoire Description
    dag_id Oui Nom du DAG.
    default_args Oui Valeurs par défaut au niveau du cluster : cluster_id, rg_name, region. Pour plus d'informations, consultez les paramètres DAG.

    Paramètres AnalyticDBSparkBatchOperator :

    Paramètre Obligatoire Description
    task_id Oui ID de la tâche.
    file Oui Chemin absolu vers le fichier principal de l'application Spark — un package JAR (Java/Scala) ou un fichier point d'entrée Python. Doit être stocké dans OSS.
    class_name Obligatoire pour Java/Scala Classe du point d'entrée. Les applications Python n'en ont pas besoin. Pour plus d'informations, consultez les paramètres 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. Copiez spark_dags.py dans le répertoire dags_folder défini dans votre configuration Airflow.

  5. Déclenchez le DAG depuis l'interface web d'Airflow. Pour obtenir des instructions, consultez le tutoriel Airflow.

spark-submit

Remarque

Vous pouvez définir les paramètres spécifiques à AnalyticDB for MySQL (

clusterId

,

regionId

,

keyId

,

secretId

,

ossUploadPath

) dans le fichier

conf/spark-defaults.conf

ou en tant que paramètres Airflow. Pour la liste complète, consultez les

paramètres de configuration des applications Spark

.

  1. Installez le plug-in Airflow Spark :

    Important

    L'installation de apache-airflow-providers-apache-spark installe automatiquement PySpark. Pour supprimer PySpark, exécutez pip3 uninstall pyspark.

    pip3 install apache-airflow-providers-apache-spark
  2. Téléchargez le package spark-submit et configurez les paramètres.

  3. Ajoutez le binaire spark-submit au PATH d'Airflow avant de démarrer Airflow :

    Important

    Définissez le PATH avant de démarrer Airflow. Si Airflow démarre sans le chemin d'accès à spark-submit, il ne pourra pas localiser la commande lors de l'exécution.

    export PATH=$PATH:</your/adb/spark/path/bin>
  4. Créez un fichier DAG nommé demo.py. L'exemple suivant exécute deux tâches JAR de manière séquentielle :

    Paramètres du DAG :

    Paramètre Obligatoire Description
    dag_id Oui Nom du DAG.
    default_args Oui Valeurs par défaut au niveau du cluster : cluster_id, rg_name, region. Pour plus d'informations, consultez les paramètres DAG.

    Paramètres AnalyticDBSparkBatchOperator :

    Paramètre Obligatoire Description
    task_id Oui ID de la tâche.
    file Oui Chemin absolu vers le fichier principal de l'application Spark. Doit être stocké dans OSS.
    class_name Obligatoire pour Java/Scala Classe du point d'entrée. Les applications Python n'en ont pas besoin. Pour plus d'informations, consultez les paramètres 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. Copiez demo.py dans le dossier dags du répertoire d'installation d'Airflow.

  6. Déclenchez le DAG depuis l'interface web d'Airflow. Pour obtenir des instructions, consultez le tutoriel Airflow.