Jak utworzyć wyzwalacz DAG w Airflow, używając Experimental API

Podczas przygotowywania naszych programów edukacyjnych regularnie napotykamy trudności związane z korzystaniem z niektórych narzędzi. W chwilach, kiedy się z nimi stykamy, nie zawsze mamy wystarczającą dokumentację i artykuły, które pomogłyby rozwiązać ten problem.

Tak było na przykład w 2015 roku, kiedy w programie „Specjalista ds. dużych danych” korzystaliśmy z klastra Hadoop z frameworkiem Spark dla 35 jednoczesnych użytkowników. Jak go przygotować dla tego przypadku użycia z wykorzystaniem YARN, nie było jasne. W końcu, po zrozumieniu i przejściu ścieżki samodzielnie, stworzyliśmy post na Habrahabr a także wystąpiliśmy na Moscow Spark Meetup.

Tło

W tym razem mowa będzie o innym programie – Data Engineer. W trakcie niego nasi uczestnicy budują dwa typy architektury: lambda i kappa. W architekturze lambda w ramach przetwarzania wsadowego wykorzystuje się Airflow do przenoszenia logów z HDFS do ClickHouse.

Ogólnie rzecz biorąc, wszystko jest w porządku. Mogą budować swoje pipeline'y. Jednak jest jedno „ale”: wszystkie nasze programy są technologiczne pod względem samego procesu nauczania. Do weryfikacji laboratoriów używamy automatycznych sprawdzaczy: uczestnik musi zalogować się do swojego panelu, kliknąć przycisk „Sprawdź”, a po chwili widzi szczegółowy feedback na temat tego, co zrobił. I właśnie w tym momencie zaczynamy zbliżać się do naszego problemu.

Weryfikacja tego laboratorium jest zorganizowana w ten sposób: wysyłamy kontrolną paczkę danych do Kafka uczestnika, następnie Gobblin przenosi tę paczkę danych na HDFS, a potem Airflow bierze tę paczkę danych i umieszcza ją w ClickHouse. Kluczowe jest to, że Airflow nie powinien tego robić w czasie rzeczywistym, lecz działa zgodnie z harmonogramem: co 15 minut pobiera paczkę plików i wrzuca je.

Okazuje się, że musimy jakoś wyzwolić ich DAG samodzielnie na nasze żądanie w trakcie pracy sprawdzacza tutaj i teraz. Po poszukiwaniu w Google stwierdziliśmy, że dla późniejszych wersji Airflow istnieje tak zwane Experimental API. Słowo experimental, oczywiście, brzmi przerażająco, ale co zrobić… Może się uda.

Dalej opiszemy całą ścieżkę: od instalacji Airflow do tworzenia zapytania POST, które wyzwala DAG, używając Experimental API. Będziemy pracować na Ubuntu 16.04.

1. Instalacja Airflow

Sprawdzimy, czy mamy zainstalowanego Pythona 3 i virtualenv.

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

Jeśli czegoś z tego brakuje, zainstaluj to.

Teraz utworzymy katalog, w którym będziemy dalej pracować z Airflow.

$ mkdir 
$ cd /ścieżka/do/twojego/nowego/katalogu
$ virtualenv -p which python3 venv
$ source venv/bin/activate
(venv) $

Zainstalujemy Airflow:

(venv) $ pip install airflow

Wersja, na której pracowaliśmy: 1.10.

Teraz musimy utworzyć katalog airflow_home, w którym będą przechowywane pliki DAG i wtyczki Airflow. Po utworzeniu katalogu ustalimy zmienną środowiskową AIRFLOW_HOME.

(venv) $ cd /ścieżka/do/mojego/airflow/roboczego
(venv) $ mkdir airflow_home
(venv) $ export AIRFLOW_HOME=

Następnym krokiem jest wykonanie polecenia, które utworzy i zainicjalizuje bazę danych strumienia danych w SQLite:

(venv) $ airflow initdb

Baza danych zostanie utworzona w airflow.db domyślnie.

Sprawdźmy, czy Airflow został zainstalowany:

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

Jeśli polecenie zostało wykonane pomyślnie, Airflow stworzył swój plik konfiguracyjny airflow.cfg do AIRFLOW_HOME:

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

Airflow ma interfejs webowy. Można go uruchomić, wykonując polecenie:

(venv) $ airflow webserver --port 8081

Teraz możesz uzyskać dostęp do interfejsu webowego w przeglądarce na porcie 8081 na hoście, na którym uruchomiono Airflow, na przykład: <hostname:8081>.

2. Praca z Experimental API

Teraz Airflow jest skonfigurowany i gotowy do pracy. Niemniej jednak musimy uruchomić również Experimental API. Nasze kontrolery są napisane w Pythonie, więc poniższe wszystkie zapytania będą wykorzystywały bibliotekę requests.

Tak naprawdę API już działa dla prostych zapytań. Na przykład, takie zapytanie pozwala przetestować jego działanie:

>>> import requests
>>> host = 
>>> airflow_port = 8081 # w naszym przypadku taki, a domyślnie 8080
>>> requests.get('http://{}:{}//{}'.format(host, airflow_port, 'api/experimental/test').text
'OK'

Jeśli otrzymasz takie powiadomienie w odpowiedzi, oznacza to, że wszystko działa.

Jednak gdy będziemy chcieli wywołać DAG, napotkamy na to, że tego rodzaju zapytanie nie może być przeprowadzone bez uwierzytelnienia.

W tym celu należy wykonać jeszcze kilka kroków.

Po pierwsze, należy dodać to do konfiguracji:

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

Następnie trzeba stworzyć własnego użytkownika z prawami administratora:

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

Następnie należy utworzyć użytkownika z normalnymi uprawnieniami, któremu będzie dozwolone uruchamianie DAG-a.

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

Teraz wszystko jest gotowe.

3. Uruchomienie żądania POST

Żądanie POST będzie wyglądać następująco:

>>> 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": "Utworzono "n}n'

Żądanie zostało pomyślnie przetworzone.

W związku z tym, dajemy jakąś chwilę DAG-owi na przetworzenie i wykonujemy zapytanie do tabeli ClickHouse, starając się uchwycić kontrolny pakiet danych.

Weryfikacja zakończona.

Źródło: habr.com

Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS 🔥 Kup solidny hosting stron z ochroną przed DDoS, serwery VPS VDS | ProHoster