Comment créer un déclencheur DAG dans Airflow en utilisant l'API expérimentale

Lors de la prĂ©paration de nos programmes Ă©ducatifs, nous rencontrons pĂ©riodiquement des difficultĂ©s liĂ©es Ă  certains outils. Et au moment oĂč nous sommes confrontĂ©s Ă  ces problĂšmes, il n'y a pas toujours suffisamment de documentation et d'articles pour nous aider Ă  les rĂ©soudre.

C'Ă©tait le cas, par exemple, en 2015, lorsque dans le programme "SpĂ©cialiste des Big Data", nous utilisions un cluster Hadoop avec Spark pour 35 utilisateurs simultanĂ©s. Comment le prĂ©parer pour ce cas d'utilisation avec YARN n'Ă©tait pas clair. Finalement, aprĂšs avoir compris et fait le chemin par nous-mĂȘme, nous avons créé un post sur Habr et nous avons Ă©galement prĂ©sentĂ© lors du Moscow Spark Meetup.

Contexte

Cette fois, nous parlerons d'un autre programme – Data Engineer. Nos participants y construisent deux types d'architectures : lambda et kappa. Dans l'architecture lambda, dans le cadre du traitement batch, nous utilisons Airflow pour transfĂ©rer les logs d'HDFS vers ClickHouse.

Tout va plutĂŽt bien. Qu'ils construisent leurs pipelines. Cependant, il y a un mais : tous nos programmes sont technologiques en ce qui concerne le processus d'apprentissage lui-mĂȘme. Pour vĂ©rifier les labs, nous utilisons des vĂ©rificateurs automatiques : le participant doit se connecter Ă  son espace personnel, cliquer sur le bouton "VĂ©rifier", et aprĂšs un certain temps, il reçoit des retours dĂ©taillĂ©s sur ce qu'il a fait. Et c'est prĂ©cisĂ©ment Ă  ce moment-lĂ  que nous commençons Ă  nous approcher de notre problĂšme.

La vérification de ce lab est organisée comme suit : nous envoyons un paquet de données de contrÎle dans Kafka du participant, ensuite Gobblin dépose ce paquet de données dans HDFS, puis Airflow prend ce paquet de données et le place dans ClickHouse. L'astuce est qu'Airflow ne doit pas le faire en temps réel, il le fait selon un calendrier : toutes les 15 minutes, il prend un lot de fichiers et les traite.

Il s'avĂšre qu'il nous faut d'une maniĂšre ou d'une autre dĂ©clencher leur DAG nous-mĂȘmes selon notre demande Ă  ce moment-lĂ  lors de l'exĂ©cution du vĂ©rificateur. En cherchant sur Google, nous avons dĂ©couvert que pour les versions rĂ©centes d'Airflow, il existe ce qu'on appelle un API ExpĂ©rimental. Le mot expĂ©rimental, bien sĂ»r, fait peur, mais que faire... Peut-ĂȘtre que ça marchera.

Nous allons maintenant dĂ©crire tout le processus : de l'installation d'Airflow Ă  la formation de la requĂȘte POST qui dĂ©clenche le DAG en utilisant l'API ExpĂ©rimentale. Nous travaillerons avec Ubuntu 16.04.

1. Installation d'Airflow

Vérifions que nous avons Python 3 et virtualenv installés.

$ python3 --version
Python 3.6.6
$ virtualenv --version
15.2.0

Si l'un d'eux n'est pas présent, veillez à les installer.

Nous allons maintenant créer un répertoire dans lequel nous allons travailler avec Airflow.

$ mkdir 
$ cd /chemin/vers/votre/nouveau/répertoire
$ virtualenv -p which python3 venv
$ source venv/bin/activate
(venv) $

Installons Airflow :

(venv) $ pip install airflow

Version que nous avons utilisée : 1.10.

Nous devons maintenant crĂ©er un rĂ©pertoire airflow_home, oĂč seront situĂ©s les fichiers DAG et les plugins Airflow. AprĂšs avoir créé le rĂ©pertoire, nous dĂ©finirons la variable d'environnement AIRFLOW_HOME.

(venv) $ cd /chemin/vers/mon/airflow/espace_de_travail
(venv) $ mkdir airflow_home
(venv) $ export AIRFLOW_HOME=

L'étape suivante consiste à exécuter la commande qui créera et initialisera la base de données du flux dans SQLite :

(venv) $ airflow initdb

La base de données sera créée dans airflow.db par défaut.

Vérifions si Airflow est bien installé :

$ airflow version
[2018-11-26 19:38:19,607] {__init__.py:57} INFO - Utilisation de l'exécuteur SequentialExecutor
[2018-11-26 19:38:19,745] {driver.py:123} INFO - Génération des tables de grammaire à partir de /usr/lib/python3.6/lib2to3/Grammar.txt
[2018-11-26 19:38:19,771] {driver.py:123} INFO - Génération des tables de grammaire à partir de /usr/lib/python3.6/lib2to3/PatternGrammar.txt
  ____________       _____________
 ____    |__( )_________  __/__  /________      __
____  /| |_  /__  ___/_  /_ __  /_  __ _ | /| / \\
___  ___ |  / _  /   _  __/ _  / / /_/_/ \_|/|/ \
 _/_/_/  |_/_/_/  _/_/    _/_/    _/_/  ____/____/|__/
   v1.10.0

Si la commande a réussi, Airflow a créé son fichier de configuration airflow.cfg dans AIRFLOW_HOME:

$ tree
.
├── airflow.cfg
└── unittests.cfg

Airflow dispose d'une interface web. Vous pouvez la lancer en exécutant la commande :

(venv) $ airflow webserver --port 8081

Vous pouvez maintenant accéder à l'interface web dans votre navigateur à l'adresse http://:8081, par exemple : <hostname:8081>.

2. Travailler avec l'API Expérimentale

À ce stade, Airflow est configurĂ© et prĂȘt Ă  l'emploi. Cependant, nous devons Ă©galement dĂ©marrer l'API ExpĂ©rimentale. Nos contrĂŽleurs sont Ă©crits en Python, donc toutes les requĂȘtes suivantes seront rĂ©alisĂ©es en utilisant la bibliothĂšque requests.

En rĂ©alitĂ©, l'API fonctionne dĂ©jĂ  pour des requĂȘtes simples. Par exemple, cette requĂȘte permet de tester son fonctionnement :

>>> import requests
>>> host = 
>>> airflow_port = 8081 # dans notre cas, c'est 8081, par défaut c'est 8080
>>> requests.get('http://{}:{} /{}'.format(host, airflow_port, 'api/experimental/test').text
'OK'

Si vous avez reçu ce message en réponse, cela signifie que tout fonctionne.

Cependant, lorsque nous souhaitons dĂ©clencher un DAG, nous constaterons que ce type de requĂȘte ne peut ĂȘtre effectuĂ© sans authentification.

À cet effet, il sera nĂ©cessaire d'effectuer quelques autres actions.

Tout d'abord, ajoutez ceci dans la configuration :

[api]
auth_backend = airflow.contrib.auth.backends.password_auth

Ensuite, vous devez créer un utilisateur avec des droits d'administration :

>>> import airflow
>>> from airflow import models, settings
>>> from airflow.contrib.auth.backends.password_auth import PasswordUser
>>> user = PasswordUser(models.Admin())
>>> user.username = 'new_user_name'
>>> user.password = 'set_the_password'
>>> session = settings.Session()
>>> session.add(user)
>>> session.commit()
>>> session.close()
>>> exit()

Ensuite, il faut créer un utilisateur avec des droits standard, qui sera autorisé à déclencher le DAG.

>>> import airflow
>>> from airflow import models, settings
>>> from airflow.contrib.auth.backends.password_auth import PasswordUser
>>> user = PasswordUser(models.User())
>>> user.username = 'newprolab'
>>> user.password = 'Newprolab2019!'
>>> session = settings.Session()
>>> session.add(user)
>>> session.commit()
>>> session.close()
>>> exit()

Maintenant tout est prĂȘt.

3. Lancement de la requĂȘte POST

La requĂȘte POST ressemblera Ă  ceci :

>>> dag_id = newprolab
>>> url = 'http://{}:{}//{}/{}/{}'.format(host, airflow_port, 'api/experimental/dags', dag_id, 'dag_runs')
>>> data = {"conf":"{"key":"value"}"}
>>> headers = {'Content-type': 'application/json'}
>>> auth = ('newprolab', 'Newprolab2019!')
>>> uri = requests.post(url, data=json.dumps(data), headers=headers, auth=auth)
>>> uri.text
'{n  "message": "Créé "n}n'

La requĂȘte a Ă©tĂ© traitĂ©e avec succĂšs.

Nous donnons ensuite un certain temps au DAG pour traiter et faisons une requĂȘte Ă  la table ClickHouse, en essayant d'attraper le paquet de contrĂŽle de donnĂ©es.

La vérification est terminée.

Source : habr.com

Acheter un hĂ©bergement fiable pour les sites avec protection DDoS, serveurs VPS VDS đŸ”„ Acheter un hĂ©bergement fiable pour les sites avec protection DDoS, serveurs VPS VDS | ProHoster