Wie man einen DAG-Trigger in Airflow mit der Experimental API erstellt

Bei der Entwicklung unserer Bildungsprogramme stoßen wir gelegentlich auf Herausforderungen im Umgang mit bestimmten Tools. Zu den Zeiten, in denen wir mit diesen Schwierigkeiten konfrontiert werden, ist nicht immer ausreichend Dokumentation und Artikel verfügbar, die uns bei der Lösung dieser Probleme helfen könnten.

So war es zum Beispiel im Jahr 2015, als wir im Programm "Spezialist für Big Data" einen Hadoop-Cluster mit Spark für 35 gleichzeitige Benutzer nutzten. Wie man ihn für diesen Anwendungsfall mit YARN vorbereitet, war unklar. Letztendlich haben wir, nachdem wir die Sache selbst durchdrungen haben, es geschafft, Post auf Habré. und haben sogar auf Moscow Spark Meetup.

Hintergrund

In diesem Fall wird es um ein anderes Programm gehen – Data Engineer. In diesem bauen unsere Teilnehmer zwei Arten von Architekturen auf: Lambda und Kappa. In der Lambda-Architektur wird im Rahmen der Batch-Verarbeitung Airflow verwendet, um Logs von HDFS nach ClickHouse zu übertragen.

Insgesamt läuft alles gut. Lassen Sie sie ihre Pipelines aufbauen. Es gibt jedoch einen Punkt: alle unsere Programme sind technologisch hinsichtlich des gesamten Lernprozesses fortgeschritten. Für die Überprüfung der Labore verwenden wir automatische Checker: Der Teilnehmer muss sich in sein persönliches Konto einloggen, die Schaltfläche „Überprüfen“ drücken, und nach einer gewissen Zeit erhält er ein umfassendes Feedback zu seinen Ergebnissen. Genau in diesem Moment nähern wir uns unserem Problem.

Die Überprüfung dieses Labors ist so aufgebaut: Wir senden ein Kontrollpaket an die Kafka des Teilnehmers, danach überträgt Gobblin dieses Datenpaket auf HDFS, und schließlich nimmt Airflow dieses Datenpaket und speichert es in ClickHouse. Der Clou ist, dass Airflow dies nicht in Echtzeit tun muss, sondern nach einem Zeitplan: Alle 15 Minuten wird ein Satz von Dateien abgeholt und hochgeladen.

Das bedeutet, dass wir irgendwie die DAGs selbstständig auf Anforderung während der Arbeit des Checkers hier und jetzt auslösen müssen. Nach einer kurzen Internetrecherche haben wir herausgefunden, dass es für neuere Versionen von Airflow so etwas wie eine Experimental APIgibt. Das Wort experimentalklingt zwar beängstigend, aber was soll man machen... Vielleicht klappt es ja.

Nun beschreiben wir den gesamten Weg: von der Installation von Airflow bis zur Erstellung einer POST-Anfrage, die den DAG über die Experimental API auslöst. Wir arbeiten mit Ubuntu 16.04.

1. Installation von Airflow

Überprüfen wir, ob Python 3 und virtualenv installiert sind.

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

Wenn eines davon fehlt, installieren Sie es bitte.

Jetzt erstellen wir ein Verzeichnis, in dem wir weiter mit Airflow arbeiten werden.

$ mkdir 
$ cd /path/to/your/new/directory
$ virtualenv -p $(which python3) venv
$ source venv/bin/activate
(venv) $

Installieren wir Airflow:

(venv) $ pip install airflow

Die Version, die wir verwendet haben: 1.10.

Jetzt müssen wir ein Verzeichnis erstellen airflow_home, in dem die DAG-Dateien und Plugins von Airflow abgelegt werden. Nach der Erstellung des Verzeichnisses setzen wir die Umgebungsvariable AIRFLOW_HOME.

(venv) $ cd /path/to/my/airflow/workspace
(venv) $ mkdir airflow_home
(venv) $ export AIRFLOW_HOME=

Der nächste Schritt besteht darin, den Befehl auszuführen, der die Datenbank für den Datenstrom in SQLite erstellen und initialisieren wird:

(venv) $ airflow initdb

Die Datenbank wird unter airflow.db als Standard.

Überprüfen wir, ob Airflow erfolgreich installiert wurde:

$ airflow version
[2018-11-26 19:38:19,607] {__init__.py:57} INFO - Verwende Executor SequentialExecutor
[2018-11-26 19:38:19,745] {driver.py:123} INFO - Generieren von Grammatiktabellen aus /usr/lib/python3.6/lib2to3/Grammar.txt
[2018-11-26 19:38:19,771] {driver.py:123} INFO - Generieren von Grammatiktabellen aus /usr/lib/python3.6/lib2to3/PatternGrammar.txt
  ____________       _____________
 ____    |__( )_________  __/__  /________      __
____  /| |_  /__  ___/_  /_ __  /_  __ _ | /| / /
___  ___ |  / _  /   _  __/ _  / / /_/ /_ |/ |/ /
 _/_/  |_/_/  /_/    /_/    /_/  ____/____/|__/
   v1.10.0

Wenn der Befehl erfolgreich ausgeführt wurde, hat Airflow seine Konfigurationsdatei erstellt. airflow.cfg in AIRFLOW_HOME:

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

Airflow verfügt über eine Web-Oberfläche, die Sie mit dem folgenden Befehl starten können:

(venv) $ airflow webserver --port 8081

Jetzt können Sie über den Webbrowser auf die Web-Oberfläche unter Port 8081 auf dem Host zugreifen, auf dem Airflow gestartet wurde, zum Beispiel: <hostname:8081>.

2. Arbeiten mit der Experimental-API

Jetzt ist Airflow eingerichtet und bereit zur Verwendung. Allerdings müssen wir auch die Experimental-API starten. Unsere Checker sind in Python geschrieben, daher werden alle nachfolgenden Anfragen damit unter Verwendung der Bibliothek erfolgen. requests.

Tatsächlich funktioniert die API bereits für einfache Anfragen. Zum Beispiel ermöglicht dieser Befehl, ihre Funktionsweise zu testen:

>>> import requests
>>> host = 
>>> airflow_port = 8081 # in unserem Fall so, standardmäßig 8080
>>> requests.get('http://{}:{}/{}'.format(host, airflow_port, 'api/experimental/test')).text
'OK'

Wenn Sie diese Nachricht als Antwort erhalten haben, bedeutet das, dass alles funktioniert.

Wenn wir jedoch den DAG auslösen möchten, stellen wir fest, dass diese Art von Anfrage nicht ohne Authentifizierung durchgeführt werden kann.

Dafür müssen noch einige Schritte unternommen werden.

Zunächst müssen Sie Folgendes in die Konfiguration hinzufügen:

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

Dann müssen Sie einen Benutzer mit Administratorrechten erstellen:

>>> 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()

Anschließend müssen Sie einen Benutzer mit normalen Rechten erstellen, der Berechtigungen zum Auslösen des DAG hat.

>>> 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()

Jetzt ist alles bereit.

3. Starten der POST-Anfrage

Die POST-Anfrage wird folgendermaßen aussehen:

>>> 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": "Erstellt <DagRun newprolab @ 2019-03-27 10:24:25+00:00: manual__2019-03-27T10:24:25+00:00, extern ausgelöst: True>"n}n'

Die Anfrage wurde erfolgreich bearbeitet.

Wir geben dem DAG nun etwas Zeit zur Verarbeitung und machen eine Anfrage an die ClickHouse-Datenbank, um den Kontrollpaket abzurufen.

Überprüfung abgeschlossen.

Quelle: habr.com

Kaufen Sie zuverlässiges Hosting für Websites mit DDoS-Schutz, VPS VDS-Servern 🔥 Kaufen Sie zuverlässiges Hosting für Websites mit DDoS-Schutz, VPS VDS-Servern | ProHoster