Managed Airflow (3e génération) | Managed Airflow (2e génération) | Managed Airflow (1re génération héritée)
Ce tutoriel explique comment utiliser Managed Service pour Apache Airflow afin de créer un DAG Apache Airflow (graphe orienté acyclique) qui exécute une tâche Apache Hadoop de décompte de mots sur un cluster Managed Service pour Apache Spark.
Objectifs
- Accéder à votre environnement Managed Airflow et utiliser l' interface utilisateur Airflow.
- Créer et afficher des variables d'environnement Airflow
- Créer et exécuter un DAG comprenant les tâches suivantes :
- Créer un cluster Managed Service pour Apache Spark
- Exécuter une tâche Apache Hadoop de décompte de mots sur le cluster
- Envoyer les résultats du décompte de mots dans un bucket Cloud Storage
- Supprimer le cluster
Coûts
Dans ce document, vous utilisez les composants facturables suivants de Google Cloud:
- Managed Airflow
- Managed Service for Apache Spark
- Cloud Storage
Obtenez une estimation des coûts en fonction de votre utilisation prévue,
utilisez le simulateur de coût.
Avant de commencer
Assurez-vous que les API suivantes sont activées dans votre projet :
Console
Activez les API Managed Service pour Apache Spark et Cloud Storage.
Rôles requis pour activer les API
Pour activer les API, vous devez disposer de l'autorisation
serviceusage.services.enable. Si vous avez créé le projet, vous disposez probablement déjà de cette autorisation via le rôle Propriétaire (roles/owner). Sinon, vous pouvez l'obtenir via le rôle Administrateur d'utilisation du service (roles/serviceusage.serviceUsageAdmin). Découvrez comment attribuer des rôles.gcloud
Activez les API Managed Service pour Apache Spark et Cloud Storage :
Rôles requis pour activer les API
Pour activer les API, vous devez disposer de l'autorisation
serviceusage.services.enable. Si vous avez créé le projet, vous disposez probablement déjà de cette autorisation via le rôle Propriétaire (roles/owner). Sinon, vous pouvez l'obtenir via le rôle Administrateur d'utilisation du service (roles/serviceusage.serviceUsageAdmin). Découvrez comment attribuer des rôles.gcloud services enable dataproc.googleapis.com
storage-component.googleapis.com Dans votre projet, créez un bucket Cloud Storage de n'importe quelle région et classe de stockage pour stocker les résultats de la tâche Hadoop de décompte de mots.
Notez le chemin du bucket que vous avez créé, par exemple
gs://example-bucket. Vous allez définir une variable Airflow pour ce chemin et l'utiliser dans l'exemple de DAG plus loin dans ce tutoriel.Créez un environnement Managed Airflow avec les paramètres par défaut. Attendez la fin de la création de l'environnement. Lorsque vous avez terminé, une coche verte s'affiche à gauche du nom de l'environnement.
Notez la région dans laquelle vous avez créé votre environnement, par exemple
us-central. Vous allez définir une variable Airflow pour cette région et l'utiliser dans l'exemple de DAG afin d'exécuter un cluster Managed Service pour Apache Spark dans la même région.
Définir les variables Airflow
Définissez les variables Airflow à utiliser ultérieurement dans l'exemple de DAG. Par exemple, vous pouvez définir des variables Airflow dans l'interface utilisateur Airflow.
| Variable Airflow | Valeur |
|---|---|
gcp_project
|
L'ID du projet du projet
que vous utilisez pour ce tutoriel, tel que example-project. |
gcs_bucket
|
URI du bucket Cloud Storage que vous avez créé pour ce tutoriel,
tel que gs://example-bucket |
gce_region
|
Région dans laquelle vous avez créé votre environnement, telle que us-central1.
Il s'agit de la région dans laquelle votre cluster Managed Service pour Apache Spark sera créé. |
Afficher l'exemple de workflow
Un DAG Airflow est un ensemble de tâches organisées que vous souhaitez programmer et exécuter. Les DAG sont définis dans des fichiers Python standards. Le code affiché dans hadoop_tutorial.py correspond au code du workflow.
Opérateurs
Pour orchestrer les trois tâches dans l'exemple de workflow, le DAG importe les trois opérateurs Airflow suivants :
DataprocClusterCreateOperator: crée un cluster Managed Service pour Apache Spark.DataProcHadoopOperator: envoie une tâche Hadoop de décompte de mots et écrit les résultats dans un bucket Cloud Storage.DataprocClusterDeleteOperator: supprime le cluster pour éviter que des frais Compute Engine ne continuent d'être facturés.
Dépendances
Vous organisez les tâches que vous souhaitez exécuter de façon à refléter leurs relations et leurs dépendances. Les tâches de ce DAG sont exécutées de manière séquentielle.
Planification
Le nom du DAG est composer_hadoop_tutorial, et le DAG s'exécute une fois par jour. Comme la valeur start_date transmise à default_dag_args est définie sur yesterday, Managed Airflow programme le workflow de façon qu'il commence immédiatement après l'importation du DAG dans le bucket de l'environnement.
Importer le DAG dans le bucket de l'environnement
Managed Airflow stocke les DAG dans le dossier /dags du bucket de votre environnement.
Pour importer le DAG :
Sur votre machine locale, enregistrez
hadoop_tutorial.py.Dans la Google Cloud console, accédez à la page Environnements.
Dans la liste des environnements, dans la colonne Dossier des DAG de votre environnement, cliquez sur le lien DAG.
Cliquez sur Importer des fichiers.
Sélectionnez
hadoop_tutorial.pysur votre machine locale, puis cliquez sur Ouvrir.
Managed Airflow ajoute le DAG à Airflow et le planifie automatiquement. Les modifications sont appliquées au DAG après trois à cinq minutes.
Explorer les exécutions du DAG
Afficher l'état des tâches
Lorsque vous importez votre fichier de DAG dans le dossier dags/ de Cloud Storage, Managed Airflow l'analyse. Une fois la procédure terminée, le nom du workflow apparaît dans la liste des DAG et le workflow est mis en file d'attente pour être exécuté immédiatement.
Pour connaître l'état des tâches, accédez à l'interface Web Airflow, puis cliquez sur DAGs (DAG) dans la barre d'outils.
Pour ouvrir la page d'informations des DAG, cliquez sur
composer_hadoop_tutorial. Cette page comprend une représentation graphique des tâches et des dépendances du workflow.
Pour connaître l'état de chaque tâche, cliquez sur Graph View (Vue graphique), puis sur le graphique de chaque tâche.
Remettre le workflow en file d'attente
Pour exécuter de nouveau le workflow à partir de la vue graphique :
- Dans la vue graphique de l'interface utilisateur Airflow, cliquez sur le graphique
create_dataproc_cluster. - Pour réinitialiser les trois tâches, cliquez sur Clear (Effacer), puis sur OK pour confirmer.
- Cliquez de nouveau sur
create_dataproc_clusterdans la vue graphique. - Pour remettre le workflow en file d'attente, cliquez sur Run (Exécuter).
Afficher les résultats des tâches
Vous pouvez également vérifier l'état et les résultats du composer_hadoop_tutorial
workflow en accédant aux pages suivantes de la Google Cloud console :
Clusters Managed Service pour Apache Spark : pour surveiller la création et la suppression de clusters. Notez que le cluster créé par le workflow est éphémère : il n'existe que pour la durée du workflow et est supprimé dans le cadre de sa dernière tâche.
Tâches Managed Service pour Apache Spark : pour afficher ou surveiller la tâche Apache Hadoop de décompte de mots. Cliquez sur l'ID de job pour afficher la sortie du journal associée au job.
Navigateur Cloud Storage : pour afficher les résultats du décompte de mots dans le dossier
wordcountdu bucket Cloud Storage que vous avez créé pour ce tutoriel.
Nettoyage
Supprimez les ressources utilisées dans ce tutoriel :
Supprimez l'environnement Managed Airflow, y compris le bucket de l'environnement que vous avez supprimé manuellement.
Supprimez le bucket Cloud Storage qui stocke les résultats de la tâche Hadoop de décompte de mots.