Как да създадете тригер за DAG в Airflow, използвайки Experimental API

При подготовката на нашите образователни програми периодично се сблъскваме със затруднения при работата с някои инструменти. И в момента, в който се сблъскваме с тях, не винаги има достатъчно документация и статии, които да помогнат за разрешаването на проблема.

Така беше, например, през 2015 година, когато в програмата „Специалист по големи данни“ ползвахме Hadoop клъстер със Spark с 35 едновременно ползващи. Как да го подготвим за такъв use case с използване на YARN беше неясно. В крайна сметка, след като се запознахме и преминахме сами през процеса, направихме пост в Хабра и също така се представихме на Moscow Spark Meetup.

Предистория

Този път ще говорим за друга програма – Data Engineer. В нея нашите участници изграждат два типа архитектура: lambda и kappa. И в lambda-архитектурата в рамките на batch обработката се използва Airflow за преместване на логовете от HDFS в ClickHouse.

Всичко е в общи линии добре. Нека изграждат своите пайплайни. Обаче, има нещо: всички наши програми са технологични от гледна точка на самия процес на обучение. За проверка на лабовете използваме автоматични чекери: участникът трябва да влезе в личния си кабинет, да натисне бутон „Проверка“, и след известно време вижда разширена обратна връзка за това, което е направил. И именно в този момент започваме да се добираме до нашия проблем.

Проверката на тази лабораторна работа е организирана така: ние изпращаме контролен пакет данни в Kafka на участника, след това Gobblin прехвърля този пакет данни на HDFS, след това Airflow взима този пакет данни и го поставя в ClickHouse. Интересното е, че Airflow не трябва да го прави в реално време, а го прави по график: на всеки 15 минути взима пачка файлове и ги прехвърля.

Въз основа на това, трябва по някакъв начин да тригерираме нашия DAG сами по нашето искане, докато работи чекера тук и сега. След като потърсихме в Google, разбрахме, че за по-късни версии на Airflow съществува т.н. Experimental API. Думата experimental, разбира се, звучи плашещо, но какво да се прави… Възможно е да проработи.

По-долу ще опишем целия процес: от инсталиране на Airflow до формулиране на POST заявка, която тригерира DAG, използвайки Experimental API. Ще работим с Ubuntu 16.04.

1. Инсталиране на Airflow

Нека проверим дали имаме инсталирани Python 3 и virtualenv.

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

Ако нещо от това липсва, моля инсталирайте.

Сега ще създадем каталог, в който ще работим с Airflow.

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

Нека инсталираме Airflow:

(venv) $ pip install airflow

Версията, с която работим: 1.10.

Сега трябва да създадем каталог airflow_home, където ще се намират DAG файловете и плъгините на Airflow. След създаване на каталога, ще установим променливата на средата AIRFLOW_HOME.

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

Следващата стъпка е да изпълним команда, която ще създаде и инициализира базата данни на потока данни в SQLite:

(venv) $ airflow initdb

Базата данни ще бъде създадена в airflow.db по подразбиране.

Нека проверим дали Airflow е инсталиран:

$ 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

Ако командата е изпълнена успешно, Airflow е създал своя конфигурационен файл airflow.cfg в AIRFLOW_HOME:

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

Airflow разполага с уеб интерфейс. Можете да го стартирате, като изпълните командата:

(venv) $ airflow webserver --port 8081

Сега можете да отидете в уеб интерфейса в браузъра на порт 8081 на хоста, където е стартиран Airflow, например: <hostname:8081>.

2. Работа с Experimental API

Sега Airflow е настроен и готов за работа. Въпреки това, трябва да стартираме и Experimental API. Нашите чекери са написани на Python, така че следващите заявки ще бъдат на него с използването на библиотеката requests.

Наистина, API вече работи за прости заявки. Например, тази заявка позволява да teствате работата му:

>>> import requests
>>> host = 
>>> airflow_port = 8081 #в нашия случай такъв, а по подразбиране 8080
>>> requests.get('http://{}:{}//{}'.format(host, airflow_port, 'api/experimental/test').text
'OK'

Ако получите такова съобщение в отговор, това означава, че всичко работи.

Въпреки това, когато искаме да задействаме DAG, ще се сблъскаме с факта, че този вид заявка не може да се направи без удостоверяване.

За това ще трябва да извършим още няколко стъпки.

На първо място, в конфигурацията трябва да добавим следното:

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

След това, трябва да създадем свой потребител с администраторски права:

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

След това трябва да създадем потребител с обикновени права, на който ще бъде разрешено да задейства 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()

Сега всичко е готово.

3. Изпълнение на POST заявката

Самата POST заявка ще изглежда така:

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

Заявката е обработена успешно.

Съответно, след това даваме малко време на DAG да се обработи и правим заявка към таблицата ClickHouse, опитвайки се да уловим контролните пакети данни.

Проверката е завършена.

Източник: habr.com

Купете надежден хостинг за сайтове със защита от DDoS, VPS и VDS сървъри 🔥 Купете надежден хостинг за сайтове със защита от DDoS, VPS и VDS сървъри | ProHoster