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 :
Un cluster AnalyticDB for MySQL Enterprise Edition, Basic Edition ou Data Lakehouse Edition
Un groupe de ressources de tâches ou un groupe de ressources interactives Spark créé pour le cluster
Python 3.7 ou version ultérieure
L'adresse IP du serveur Airflow ajoutée à la liste d'autorisation du cluster
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
-
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 -
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 :
ImportantUtilisez 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_typeMéthode d'authentification. Définissez la valeur AKpour utiliser l'authentification par paire AccessKey.access_key_idL'AccessKey ID de votre utilisateur RAM ayant accès à AnalyticDB for MySQL. access_key_secretL'AccessKey Secret de votre utilisateur RAM. regionL'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>" } -
Créez un fichier DAG nommé
spark_dags.py. L'exemple suivant utiliseAnalyticDBSparkSQLOperatorpour exécuter une requêteSHOW DATABASES:Paramètres du DAG :
Paramètre Obligatoire Description dag_idOui Nom du DAG. default_argsOui 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_idOui ID de la tâche. sqlOui 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 Copiez
spark_dags.pydans le répertoiredags_folderdéfini dans votre configuration Airflow.Déclenchez le DAG depuis l'interface web d'Airflow. Pour obtenir des instructions, consultez le tutoriel Airflow.
spark-submit
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
.
-
Installez le plug-in Airflow Spark :
ImportantL'installation de
apache-airflow-providers-apache-sparkinstalle automatiquement PySpark. Pour supprimer PySpark, exécutezpip3 uninstall pyspark.pip3 install apache-airflow-providers-apache-spark Téléchargez le package spark-submit et configurez les paramètres.
-
Ajoutez le binaire spark-submit au
PATHd'Airflow avant de démarrer Airflow :ImportantDéfinissez le
PATHavant 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> -
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 Copiez
demo.pydans le dossierdagsdu répertoire d'installation d'Airflow.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.
-
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.
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.
Dans le volet de navigation, choisissez Cluster Management > Resource Management, puis cliquez sur l'onglet Resource Groups.
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.

-
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.
Dans l'interface web d'Airflow, accédez à Admin > Connections.
-
Cliquez sur
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 defaultpar le nom réel de la base de données et supprimez le suffixeresource_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 10000Extra {"auth_mechanism": "CUSTOM"} -
Créez un fichier DAG. L'exemple suivant exécute
show databasesselon une planification quotidienne :Paramètre Obligatoire Description task_idOui ID de la tâche. conn_idOui Nom de la connexion défini à l'étape 4. sqlOui 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 Dans l'interface web d'Airflow, cliquez sur
à côté du DAG pour le déclencher.
Planifier des tâches Spark JAR
Spark Airflow Operator
-
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 -
Créez une connexion Airflow. Dans l'interface web d'Airflow, accédez à Admin > Connections et ajoutez une connexion avec le code JSON suivant :
ImportantUtilisez 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_typeMéthode d'authentification. Définissez la valeur AK.access_key_idAccessKey ID de votre utilisateur RAM. access_key_secretAccessKey Secret de votre utilisateur RAM. regionID 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>" } -
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 deAnalyticDBSparkBatchOperator:ImportantStockez 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_idOui Nom du DAG. default_argsOui 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_idOui ID de la tâche. fileOui 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_nameObligatoire 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 Copiez
spark_dags.pydans le répertoiredags_folderdéfini dans votre configuration Airflow.Déclenchez le DAG depuis l'interface web d'Airflow. Pour obtenir des instructions, consultez le tutoriel Airflow.
spark-submit
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
.
-
Installez le plug-in Airflow Spark :
ImportantL'installation de
apache-airflow-providers-apache-sparkinstalle automatiquement PySpark. Pour supprimer PySpark, exécutezpip3 uninstall pyspark.pip3 install apache-airflow-providers-apache-spark Téléchargez le package spark-submit et configurez les paramètres.
-
Ajoutez le binaire spark-submit au
PATHd'Airflow avant de démarrer Airflow :ImportantDéfinissez le
PATHavant 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> -
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_idOui Nom du DAG. default_argsOui 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_idOui ID de la tâche. fileOui Chemin absolu vers le fichier principal de l'application Spark. Doit être stocké dans OSS. class_nameObligatoire 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 Copiez
demo.pydans le dossierdagsdu répertoire d'installation d'Airflow.Déclenchez le DAG depuis l'interface web d'Airflow. Pour obtenir des instructions, consultez le tutoriel Airflow.