Orchestrer les Lakeflow Jobs avec Apache Airflow
Utilisez Apache Airflow pour orchestrer les pipelines de données avec Databricks via le fournisseur Databricks open source. Installez et configurez Airflow localement, puis déployez et exécutez un Job Databricks à partir d'un DAG Airflow.
Le support Apache Airflow décrit sur cette page utilise plusieurs packages open source. Cela inclut le fournisseur Databricks pour Airflow (y compris les opérateurs Databricks Airflow). Ces packages ne sont pas directement pris en charge par Databricks. Pour des informations sur le fournisseur Databricks pour Airflow, consultez apache-airflow-providers-databricks sur Apache.org.
Orchestration de Job dans une pipeline de données
Le développement et le déploiement d'un pipeline de traitement de données nécessitent souvent la gestion de dépendances complexes entre les tâches. Par exemple, un pipeline peut lire des données à partir d'une source, nettoyer les données, transformer les données nettoyées et écrire les données transformées vers une cible. Vous avez également besoin d'un support pour les tests, la planification et le dépannage des erreurs lorsque vous opérationnalisez un pipeline.
Les systèmes de workflow relèvent ces défis en vous permettant de définir les dépendances entre les tâches, de planifier l'exécution des pipelines et de surveiller les workflows. Apache Airflow est une solution open source pour la gestion et la planification des pipelines de données. Airflow représente les pipelines de données sous forme de graphes orientés acycliques (DAG) d'opérations. Vous définissez un flux de travail dans un fichier Python, et Airflow gère la planification et l'exécution. La connexion Airflow Databricks vous permet de tirer parti du moteur Spark optimisé offert par Databricks avec les fonctionnalités de planification d'Airflow.
Exigences
- L'intégration entre Airflow et Databricks nécessite Airflow version 2.5.0 et ultérieure. Les exemples de cet article sont testés avec Airflow version 2.6.1.
- Airflow nécessite Python 3.8, 3.9, 3.10 ou 3.11. Les exemples de cet article sont testés avec Python 3.8.
- Les instructions de cet article pour installer et exécuter Airflow nécessitent pipenv pour créer un environnement virtuel Python.
Opérateurs Airflow pour Databricks
Un DAG Airflow est composé de tâches, où chaque tâche exécute un opérateur Airflow. Les opérateurs Airflow prenant en charge l’intégration à Databricks sont implémentés dans le fournisseur Databricks.
Le fournisseur Databricks inclut des opérateurs pour exécuter plusieurs tâches sur un Workspace Databricks, notamment l'importation de données dans une table, l'exécution de requêtes SQL et l'utilisation de dossiers Git Databricks.
Le fournisseur Databricks implémente deux opérateurs pour déclencher des jobs :
- Le DatabricksRunNowOperator nécessite un Job Databricks existant et utilise le POST /api/2.1/jobs/run-now Requête API pour Trigger une exécution. Databricks vous recommande d'utiliser le
DatabricksRunNowOperatorcar il réduit la duplication des définitions de Jobs, et les exécutions de Jobs déclenchées avec cet opérateur peuvent être trouvées dans l'interface utilisateur des Jobs. - Le DatabricksSubmitRunOperator ne nécessite pas qu'un Job existe dans Databricks et utilise la POST /api/2.1/jobs/runs/submit Requête API pour soumettre la spécification du Job et Trigger un run.
Pour créer un nouveau Job Databricks ou Reset un Job existant, le fournisseur Databricks implémente le DatabricksCreateJobsOperator. Le DatabricksCreateJobsOperator utilise l'API POST /api/2.1/jobs/create et POST /api/2.1/jobs/reset Requêtes API. Vous pouvez utiliser le DatabricksCreateJobsOperator avec le DatabricksRunNowOperator pour créer et exécuter un job.
L'utilisation des opérateurs Databricks pour Trigger un Job nécessite de fournir des identifiants dans la configuration de connexion Databricks. Voir Créez un jeton d'accès personnel Databricks pour Airflow.
Les opérateurs Databricks Airflow écrivent l'URL de la page d'exécution du Job dans les logs Airflow toutes les polling_period_seconds (la valeur par default est de 30 secondes). Pour plus d'informations, consultez la page du package apache-airflow-providers-databricks sur le site web d'Airflow.
Installez l'intégration Airflow Databricks localement
Pour installer Airflow et le fournisseur Databricks localement pour les tests et le développement, utilisez les étapes suivantes. Pour les autres options d'installation d'Airflow, y compris la création d'une installation de production, consultez l'installation dans la documentation Airflow.
Ouvrez un terminal et exécutez les commandes suivantes :
mkdir airflow
cd airflow
pipenv --python 3.8
pipenv shell
export AIRFLOW_HOME=$(pwd)
pipenv install apache-airflow
pipenv install apache-airflow-providers-databricks
mkdir dags
airflow db init
airflow users create --username admin --firstname <firstname> --lastname <lastname> --role Admin --email <email>
Remplacez <firstname>, <lastname> et <email> par votre nom d'utilisateur et votre e-mail. Il vous est demandé de saisir un mot de passe pour l'utilisateur administrateur. Veillez à enregistrer ce mot de passe, car il est requis pour se connecter à l'interface utilisateur d'Airflow.
Ce script effectue les étapes suivantes :
- Crée un répertoire nommé
airflowet se place dans ce répertoire. - Utilise
pipenvpour créer et générer un environnement virtuel Python. Databricks recommande d'utiliser un environnement virtuel Python pour isoler les versions de package et les dépendances de code à cet environnement. Cette isolation aide à réduire les incohérences inattendues de versions de package et les collisions de dépendances de code. - Initialise une variable d'environnement nommée
AIRFLOW_HOMEdéfinie sur le chemin du répertoireairflow. - Installe Airflow et les packages du fournisseur Airflow Databricks.
- Crée un répertoire
airflow/dags. Airflow utilise le répertoiredagspour stocker les définitions de DAG. - Initialise une base de données SQLite qu'Airflow utilise pour suivre les métadonnées. Dans un déploiement Airflow de production, vous configurerez Airflow avec une base de données standard. La base de données SQLite et la configuration par default de votre déploiement Airflow sont initialisées dans le répertoire
airflow. - Crée un utilisateur admin pour Airflow.
Pour confirmer l'installation du fournisseur Databricks, exécutez la commande suivante dans le répertoire d'installation d'Airflow :
airflow providers list
Start le serveur web et l'ordonnanceur Airflow
Le serveur web Airflow est requis pour afficher l'interface utilisateur Airflow. Pour start le serveur web, ouvrez un terminal dans le répertoire d'installation d'Airflow et exécutez les commandes suivantes :
Si le serveur web Airflow ne start pas en raison d'un conflit de ports, vous pouvez modifier le port default dans la configuration Airflow.
pipenv shell
export AIRFLOW_HOME=$(pwd)
airflow webserver
L'ordonnanceur est le composant Airflow qui planifie les DAGs. Pour start le planificateur, ouvrez un nouveau terminal dans le répertoire d'installation d'Airflow et exécutez les commandes suivantes :
pipenv shell
export AIRFLOW_HOME=$(pwd)
airflow scheduler
Tester l'installation d'Airflow
Pour vérifier l'installation d'Airflow, vous pouvez exécuter l'un des DAGs d'exemple inclus avec Airflow :
- Dans une fenêtre de navigateur, ouvrez
http://localhost:8080/home. Connexion à l'interface utilisateur d'Airflow avec le nom d'utilisateur et le mot de passe que vous avez créés lors de l'installation d'Airflow. La page DAGs d'Airflow apparaît. - Cliquez sur le bouton **Pause/Reprendre DAG** pour réactiver l'un des DAG d'exemple, par exemple,
example_python_operatorle. - Trigger l'exemple de DAG en cliquant sur le bouton Déclencher le DAG .
- Cliquez sur le nom du DAG pour afficher les détails, y compris l'état d'exécution du DAG.
Créez un jeton d'accès personnel Databricks pour Airflow
Airflow se connecte à Databricks à l'aide d'un jeton d'accès personnel Databricks (PAT). Pour créer un PAT, suivez les étapes de Créer des jetons d'accès personnels pour les utilisateurs du workspace.
En tant que bonne pratique de sécurité lorsque vous vous authentifiez avec des outils, des systèmes, des scripts et des applications automatisés, Databricks vous recommande d'utiliser des jetons OAuth.
Si vous utilisez l'authentification par jeton d'accès personnel, Databricks recommande d'utiliser des jetons d'accès personnels appartenant aux Service Principal plutôt qu'aux utilisateurs du Workspace. Pour créer des jetons pour les Service Principals, consultez Gérer les jetons pour un Service Principal.
Vous pouvez également vous authentifier auprès de Databricks à l’aide de l’OAuth Databricks pour les Service Principal. Voir Databricks Connection dans la documentation Airflow.
Configurez une connexion Databricks
Votre installation Airflow contient une connexion default pour Databricks. Pour mettre à jour la connexion afin de vous connecter à votre Workspace à l'aide du jeton d'accès personnel que vous avez créé ci-dessus :
- Dans une fenêtre de navigateur, ouvrez
http://localhost:8080/connection/list/. Si vous êtes invité à vous connecter, entrez votre nom d’utilisateur et mot de passe d’administrateur. - Sous Conn ID , localisez databricks_default et cliquez sur le bouton Edit record .
- Remplacez la valeur du champ **Hôte** par le nom d'instance du Workspace de votre déploiement Databricks, par
https://adb-123456789.cloud.databricks.comexemple,. - Dans le champ Mot de passe , saisissez votre jeton d’accès personnel Databricks.
- Cliquez sur Enregistrer .
Exemple : Créez un DAG Airflow pour exécuter un Job Databricks.
L'exemple suivant montre comment créer un déploiement Airflow simple qui s'exécute sur votre machine locale et déploie un DAG d'exemple pour Trigger des exécutions dans Databricks. Dans cet exemple, vous allez :
- Créez un nouveau Notebook et ajoutez du code pour afficher un message de bienvenue basé sur un parameter configuré.
- Créez un job Databricks avec une seule tâche qui exécute le notebook.
- Configurez une connexion Airflow à votre workspace Databricks.
- Créez un DAG Airflow pour Trigger le Job Notebook. Vous définissez le DAG dans un script Python en utilisant
DatabricksRunNowOperator. - Utilisez l'interface utilisateur d'Airflow pour Trigger le DAG et afficher l'état de l'exécution.
Créer un Notebook
Cet exemple utilise un Notebook contenant deux cellules :
- La première cellule contient un widget de texte des utilitaires Databricks définissant une variable nommée
greetingdéfinie sur la valeur par defaultworld. - La deuxième cellule imprime la valeur de la variable
greetingpréfixée parhello.
Pour créer le Notebook :
-
Accédez à votre Workspace Databricks, cliquez
sur **Nouveau** dans la barre latérale, et sélectionnez **Notebook**.
-
Donnez un nom à votre Notebook, tel que Hello Airflow, et assurez-vous que la langue default est définie sur Python.
-
Copiez le code Python suivant et collez-le dans la première cellule du notebook.
Pythondbutils.widgets.text("greeting", "world", "Greeting")
greeting = dbutils.widgets.get("greeting") -
Ajoutez une nouvelle cellule sous la première cellule et copiez-collez le code Python suivant dans la nouvelle cellule :
Pythonprint("hello {}".format(greeting))
Créer un Job
- Dans votre Workspace, cliquez
sur **Tâches et pipelines** dans la barre latérale.
- Cliquez sur Créer , puis sur Job .
- Cliquez sur la vignette Notebook pour configurer la première tâche. Si la vignette Notebook n'est pas disponible, cliquez sur Ajouter un autre type de tâche et recherchez Notebook .
- En option, remplacez le nom du Job, qui est par default
New Job <date-time>, par le nom de votre Job. - Dans Nom de la tâche , entrez un nom pour la tâche, par exemple
greeting-task. - Dans le menu déroulant Source , sélectionnez Workspace .
- Cliquez sur la zone de texte Chemin et utilisez l'explorateur de fichiers pour trouver le Notebook que vous avez créé, cliquez sur le nom du Notebook, et cliquez sur Confirmer .
- Click Add under parameter . Dans le champ Clé , saisissez
greeting. Dans le champ Valeur , saisissezAirflow user. - Cliquez sur Enregistrer la tâche .
Dans le panneau **Détails du Job**, copiez la valeur **ID du Job**. Cette valeur est requise pour trigger le Job depuis Airflow.
Exécuter le job
Pour tester votre nouveau job dans l'interface utilisateur de Lakeflow Jobs, cliquez sur dans le coin supérieur droit. Une fois l'exécution terminée, vous pouvez vérifier le résultat en consultant les détails d'exécution du Job.
Créer un nouveau DAG Airflow
Vous définissez un DAG Airflow dans un fichier Python. Pour créer un DAG pour Trigger le Job du Notebook d'exemple :
-
Dans un éditeur de texte ou un IDE, créez un nouveau fichier nommé
databricks_dag.pyavec le contenu suivant :Pythonfrom airflow import DAG
from airflow.providers.databricks.operators.databricks import DatabricksRunNowOperator
from airflow.utils.dates import days_ago
default_args = {
'owner': 'airflow'
}
with DAG('databricks_dag',
start_date = days_ago(2),
schedule_interval = None,
default_args = default_args
) as dag:
opr_run_now = DatabricksRunNowOperator(
task_id = 'run_now',
databricks_conn_id = 'databricks_default',
job_id = JOB_ID
)Remplacez
JOB_IDpar la valeur de l’ID du Job enregistré précédemment. -
Enregistrez le fichier dans le répertoire
airflow/dags. Airflow lit et installe automatiquement les fichiers DAG stockés dansairflow/dags/.
Installer et vérifier le DAG dans Airflow
Pour Trigger et vérifier le DAG dans l'interface utilisateur d'Airflow :
- Dans une fenêtre de navigateur, ouvrez
http://localhost:8080/home. L’écran **DAGs** Airflow apparaît. - Localisez
databricks_daget cliquez sur le bouton bascule Suspendre/Reprendre le DAG pour reprendre le DAG. - Déclenchez le DAG en cliquant sur le bouton Trigger DAG .
- Cliquez sur une exécution dans la colonne Exécutions pour afficher le statut et les détails de l'exécution.