Exécuter les pipelines dans un workflow
Vous pouvez exécuter un pipeline dans le cadre d'un workflow de traitement de données avec Lakeflow Jobs, Apache Airflow ou Azure Data Factory.
Un pipeline résout automatiquement les dépendances entre ses datasets, il gère donc l'orchestration simple et interne au pipeline de manière autonome. Pour l'orchestration pour laquelle un pipeline n'est pas conçu, comme l'exécution conditionnelle, la création de branches en fonction des résultats des tâches, les nouvelles tentatives ou la coordination d'un pipeline avec d'autres types de travail, utilisez un orchestrateur de workflow dédié au lieu d'intégrer la logique dans le pipeline.
Préparez votre pipeline pour l'orchestration
L'orchestration fonctionne mieux lorsque chaque pipeline couvre une unité de travail distincte que vous souhaitez planifier, valider ou exécuter indépendamment. Concevez vos pipelines autour de ces limites afin qu'un workflow puisse les coordonner comme des tâches distinctes, y compris un flux de contrôle approprié entre les tâches en amont et en aval.
Si vous disposez déjà d'un grand pipeline qui combine un travail que vous souhaitez orchestrer séparément, divisez-le en pipelines plus petits en déplaçant les tables vers un nouveau pipeline. Voir Déplacer des tables entre les pipelines.
Lakeflow Jobs
Vous pouvez orchestrer plusieurs tâches dans les Lakeflow Jobs afin d'implémenter un workflow de traitement de données. Pour inclure un pipeline dans un Job, utilisez la tâche Pipeline lorsque vous créez un Job. Voir la tâche de pipeline pour les Jobs.
Apache Airflow
Apache Airflow est une solution open source pour la gestion et la planification des workflows de données. Airflow représente les workflows sous forme de graphes orientés acycliques (DAG) d'opérations. Vous définissez un workflow dans un fichier Python et Airflow gère la planification et l'exécution. Pour des information sur l'installation et l'utilisation d'Airflow avec Databricks, consultez Orchestrer les Lakeflow Jobs avec Apache Airflow.
Pour exécuter un pipeline dans le cadre d'un workflow Airflow, utilisez le DatabricksSubmitRunOperator.
Exigences
Les éléments suivants sont requis pour utiliser la prise en charge d'Airflow pour les LakeFlow Pipelines :
- Version Airflow 2.1.0 ou version ultérieure.
- La version 2.1.0 du package du fournisseur Databricks ou version ultérieure.
Exemple
L'exemple suivant crée un DAG Airflow qui déclenche une mise à jour pour le pipeline avec l'identifiant 8279d543-063c-4d63-9926-dae38e35ce8b:
from airflow import DAG
from airflow.providers.databricks.operators.databricks import DatabricksSubmitRunOperator
from airflow.utils.dates import days_ago
default_args = {
'owner': 'airflow'
}
with DAG('ldp',
start_date=days_ago(2),
schedule_interval="@once",
default_args=default_args
) as dag:
opr_run_now=DatabricksSubmitRunOperator(
task_id='run_now',
databricks_conn_id='CONNECTION_ID',
pipeline_task={"pipeline_id": "8279d543-063c-4d63-9926-dae38e35ce8b"}
)
Remplacez CONNECTION_ID par l’identifiant d’une connexion Airflow à votre Workspace.
Enregistrez cet exemple dans le répertoire airflow/dags et utilisez l'interface utilisateur Airflow pour afficher et Trigger le DAG. Utilisez l'interface utilisateur du pipeline pour afficher les détails de la mise à jour du pipeline.
Azure Data Factory
Lakeflow pipelines et Azure Data Factory incluent chacun des options pour configurer le nombre de tentatives en cas d'échec. Si les valeurs de nouvelle tentative sont configurées sur votre pipeline et sur l'activité Azure Data Factory qui appelle le pipeline, le nombre de nouvelles tentatives est la valeur de nouvelle tentative d'Azure Data Factory multipliée par la valeur de nouvelle tentative du pipeline.
Par exemple, si une mise à jour de pipeline échoue, le pipeline réessaie la mise à jour jusqu'à cinq fois par default. Si la nouvelle tentative d’Azure Data Factory est définie sur trois, et que votre pipeline utilise le default de cinq nouvelles tentatives, votre pipeline défaillant pourrait être relancé jusqu’à quinze fois. Pour éviter les tentatives de nouvelle tentative excessives lorsque les mises à jour de pipeline échouent, Databricks recommande de limiter le nombre de nouvelles tentatives lors de la configuration du pipeline ou de l'activité Azure Data Factory qui appelle le pipeline.
Pour modifier la configuration de relance de votre pipeline, utilisez le paramètre pipelines.numUpdateRetryAttempts lors de la configuration du pipeline.
Azure Data Factory est un service ETL basé sur le cloud qui vous permet d'orchestrer l'intégration de données et les flux de travail de transformation. Azure Data Factory prend directement en charge l'exécution de tâches Databricks dans un workflow, y compris les Notebooks, les tâches JAR et les scripts Python. Vous pouvez également inclure un pipeline dans un workflow en appelant l'API REST du pipeline à partir d'une activité web Azure Data Factory. Par exemple, pour Trigger une mise à jour de pipeline à partir d'Azure Data Factory :
-
Créez une data factory ou ouvrez une data factory existante.
-
Lorsque la création est terminée, ouvrez la page de votre data factory et cliquez sur la vignette Ouvrir Azure Data Factory Studio . L'interface utilisateur d'Azure Data Factory apparaît.
-
Créez un nouveau pipeline Azure Data Factory en sélectionnant **Pipeline** dans le menu déroulant **Nouveau** de l'interface utilisateur d'Azure Data Factory Studio.
-
Dans la boîte à outils Activités , développez Général et faites glisser l'activité Web vers le canevas de pipeline. Cliquez sur l'onglet Paramètres et entrez les valeurs suivantes :
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.
-
URL :
https://<databricks-instance>/api/2.0/pipelines/<pipeline-id>/updates.Remplacez
<get-workspace-instance>.Remplacez
<pipeline-id>par l'identifiant du pipeline. -
Méthode : sélectionnez POST dans le menu déroulant.
-
En-têtes : cliquez sur + Nouveau . Dans la zone de texte Nom , saisissez
Authorization. Dans la zone de texte Valeur , saisissezBearer <personal-access-token>.Remplacez
<personal-access-token>par un jeton d’accès personnel Databricks. -
Corps : Pour transmettre des paramètres de requête supplémentaires, saisissez un document JSON contenant les paramètres. Par exemple, pour start une mise à jour et retraiter toutes les données pour le pipeline :
{"full_refresh": "true"}. S'il n'y a pas de paramètres de requête supplémentaires, saisissez des accolades vides ({}).
Pour tester l'activité Web, cliquez sur Déboguer dans la barre d'outils du pipeline de l'interface utilisateur de Data Factory. Le résultat et le statut de l'exécution, y compris les erreurs, s'affichent dans l'onglet **Output** de l'Azure Data Factory pipeline. Utilisez l'interface utilisateur des pipelines pour afficher les détails de la mise à jour du pipeline.
Une exigence de workflow courante est de start une tâche après l'achèvement d'une tâche précédente. Étant donné que la requête du pipeline updates est asynchrone, renvoyant après le début de la mise à jour mais avant qu'elle ne soit terminée, les tâches de votre pipeline Azure Data Factory qui dépendent de la mise à jour du pipeline doivent attendre que la mise à jour soit terminée. Une option pour attendre la fin de la mise à jour consiste à ajouter une activité Until après l'activité Web qui Trigger la mise à jour du pipeline. Dans l'activité Until :
- Ajoutez une activité d’attente pour attendre un nombre configuré de secondes pour l’achèvement de la mise à jour.
- Ajoutez une activité Web après l'activité d'attente qui utilise la requête de détails de mise à jour du pipeline pour obtenir le statut de la mise à jour. Le champ
statede la réponse renvoie l'état actuel de la mise à jour, y compris si elle est terminée. - Utilisez la valeur du champ
statepour définir la condition d'arrêt de l'activité Jusqu'à. Vous pouvez également utiliser une activité Définir la variable pour ajouter une variable de pipeline basée sur la valeurstateet utiliser cette variable pour la condition de fin.