Airflow est un outil d'ordonnancement open source populaire qui fournit des utilitaires en ligne de commande et une interface web pour orchestrer des charges de travail sous forme de DAG. Utilisez Airflow pour orchestrer des jobs ETL et des flux de données en temps réel dans AnalyticDB for MySQL, afin d'automatiser le traitement des données et d'améliorer l'efficacité.
Prérequis
Un cluster Enterprise Edition, Basic Edition ou Data Lakehouse Edition AnalyticDB for MySQL est créé.
Airflow est installé. Pour plus d'informations, consultez la documentation Airflow.
L'adresse IP du serveur Airflow est ajoutée à la liste d'autorisation du cluster AnalyticDB for MySQL. Pour plus d'informations, consultez la rubrique Configurer une liste d'autorisation d'adresses IP.
Procédure
-
Vérifiez si le fournisseur apache-airflow-providers-mysql est installé.
Dans l'interface utilisateur Airflow, cliquez sur .
Sur la page Providers, vérifiez la présence de apache-airflow-providers-mysql dans la liste.
-
(Facultatif) Si le fournisseur apache-airflow-providers-mysql ne figure pas dans la liste, exécutez la commande suivante pour l'installer :
pip install apache-airflow-providers-mysqlImportantSi l'erreur
OSError: mysql_config not foundse produit, exécutez la commandeyum install mysql-develpour installer les fichiers de développement MySQL. Ensuite, relancez la commande d'installation de apache-airflow-providers-mysql.
-
Créez une connexion.
Dans l'interface utilisateur Airflow, cliquez sur .
-
Cliquez sur l'icône
. Sur la page Add Connection, configurez les paramètres suivants.Parameter
Description
Connection ID
ID unique de la connexion.
Connection type
Sélectionnez MySQL.
Host
Endpoint du cluster AnalyticDB for MySQL. Vous trouverez cet endpoint sur la page Cluster Information de la console.
Login
Nom d'utilisateur du compte AnalyticDB for MySQL.
Password
Mot de passe du compte AnalyticDB for MySQL.
Port
Port du cluster AnalyticDB for MySQL. La valeur est fixée à 3306.
RemarqueLes autres paramètres sont facultatifs. Configurez-les selon vos besoins.
-
Accédez au répertoire d'installation d'Airflow et vérifiez le paramètre dags_folder dans le fichier
airflow.cfg.-
Accédez au répertoire d'installation d'Airflow.
cd /root/airflow -
Vérifiez le paramètre dags_folder dans le fichier
airflow.cfg.cat airflow.cfg -
(Facultatif) Si le dossier spécifié par le paramètre dags_folder n'existe pas, exécutez la commande
mkdirpour le créer.RemarquePar exemple, si le chemin dags_folder est
/root/airflow/dagsmais que le dossierdagsn'existe pas dans le répertoire/root/airflow, créez-le.
-
-
Créez un fichier DAG, tel que
mysql_dags.py:from airflow import DAG from airflow.providers.mysql.operators.mysql import MySqlOperator from airflow.utils.dates import days_ago default_args = { 'owner': 'airflow', } dag = DAG( 'example_mysql', default_args=default_args, start_date=days_ago(2), tags=['example'], ) mysql_test = MySqlOperator( task_id='mysql_test', mysql_conn_id='test', sql='SHOW DATABASES;', dag=dag, ) mysql_test_task = MySqlOperator( task_id='mysql_test_task', mysql_conn_id='test', sql='SELECT * FROM test;', dag=dag, ) mysql_test >> mysql_test_task if __name__ == "__main__": dag.cli()Le tableau suivant décrit les paramètres clés.
mysql_conn_id: ID de connexion créé à l'étape 2.sql: instruction SQL à exécuter.
Pour plus d'informations sur ces paramètres, consultez la documentation Airflow.
-
Dans l'interface utilisateur Airflow, localisez votre DAG et cliquez sur l'icône
dans la colonne Actions pour l'exécuter.Une fois l'exécution du DAG terminée, cliquez sur le cercle vert dans la colonne Runs pour afficher les détails de l'exécution.
Un « 1 » dans le cercle vert à côté de example_mysql indique une exécution réussie du DAG.
Sur la page Task Instances, vous pouvez consulter trois enregistrements de tâches. Si le champ State affiche success, que le Dag ID est example_mysql, que le Task ID est mysql_test_task et que l'Operator est MySqlOperator, la tâche a été exécutée avec succès.
ImportantPar défaut, Airflow utilise le fuseau horaire UTC (Coordinated Universal Time). Cela signifie que l'heure d'exécution affichée accuse un retard de 8 heures par rapport à l'heure normale de Chine (UTC+8).