При подготовката на нашите образователни програми периодично се сблъскваме със сложности от гледна точка на работата с някои инструменти. А в момента, когато се натъкваме на тях, не винаги има достатъчно документация и статии, които да помогнат да се справим с този проблем.
Така беше, например, през 2015 година, когато в програмата “Специалист по големи данни” използвахме Hadoop клъстър със Spark на 35 едновременни потребители. Как да го подготвим за такъв случай с YARN не беше ясно. В крайна сметка, след като се ориентирахме и преминахме през процеса сами, направихме и освен това участвахме на .
Предистория
Този път ще говорим за друга програма – . На нея нашите участници строят два типа архитектура: lambda и kappa. И в lambda архитектурата при батч обработка се използва Airflow за прехвърляне на логовете от HDFS в ClickHouse.
Всичко е общо взето добре. Нека да строят своите пайплайни. Но тук идва основния проблем: всички наши програми са технологични от гледна точка на самия процес на обучение. За проверка на лаби използваме автоматични чекери: участникът трябва да влезе в личния си кабинет, да натисне бутона “Проверка”, и след известно време вижда разширена обратна връзка за това, което е направил. И точно в този момент започваме да се сблъскваме с нашия проблем.
Проверка на тази лаборатория е организирана така: изпращаме контролен пакет данни в Kafka на участника, след което Gobblin прехвърля този пакет данни на HDFS, след това Airflow взема този пакет данни и го поставя в ClickHouse. Идеята е, че Airflow не трябва да го прави в реално време, а го прави на график: на всеки 15 минути взема пакет файлове и ги поставя.
Получава се, че трябва да тригерираме техния DAG сами по нашето изискване по време на работа на чекера тук и сега. След като потърсихме, установихме, че за по-късните версии на Airflow съществува т.н. . Думата 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 - Използва се изпълнителят SequentialExecutor
[2018-11-26 19:38:19,745] {driver.py:123} INFO - Генериране на граматични таблици от /usr/lib/python3.6/lib2to3/Grammar.txt
[2018-11-26 19:38:19,771] {driver.py:123} INFO - Генериране на граматични таблици от /usr/lib/python3.6/lib2to3/PatternGrammar.txt
____________ _____________
____ |__( )_________ __/__/________ __
____ /| |_ /__ ___/_ /_ __ /_ __ _ | /| /
___ ___ | / _ / _ __/ _ / / \/ _/_ |/ |/
_/_/ |_/_/_/ /_/_/ /_/_/ /_/_/ ____/____/|__/
v1.10.0Ако командата е работила, то Airflow е създал своя конфигурационен файл airflow.cfg в AIRFLOW_HOME:
$ tree
.
├── airflow.cfg
└── unittests.cfgAirflow разполага с уеб интерфейс. Може да бъде стартиран, като изпълните командата:
(venv) $ airflow webserver --port 8081Сега можете да влезете в уеб интерфейса в браузъра на порт 8081 на хоста, където Airflow е стартиран, например: <hostname:8081>.
2. Работа с Experimental API
С това Airflow е конфигуриран и готов за работа. Въпреки това, трябва да стартираме и Experimental API. Нашите проверяващи програми са написани на Python, така че следващите запитвания ще бъдат на него с използване на библиотеката requests.
Всъщност API вече работи за прости запитвания. Например, такава заявка позволява да се тества работата му:
>>> 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
