Bij het ontwikkelen van onze opleidingsprogramma's stuiten we af en toe op moeilijkheden met betrekking tot sommige tools. En op het moment dat we deze tegenkomen, is er niet altijd voldoende documentatie en artikelen die ons helpen met deze problemen.
Dit was bijvoorbeeld het geval in 2015, toen we in het programma 'Specialist in Big Data' een Hadoop-cluster met Spark gebruikten voor 35 gelijktijdige gebruikers. Hoe je dit moest voorbereiden voor zo'n use case met YARN was ons onduidelijk. Uiteindelijk, door alles zelf uit te zoeken, hebben we en ook gepresenteerd op .
Achtergrond
Deze keer gaat het over een ander programma - . Onze deelnemers bouwen twee soorten architecturen: lambda en kappa. In de lambda-architectuur wordt Airflow gebruikt voor batchverwerking om logs uit HDFS naar ClickHouse te verplaatsen.
Over het algemeen gaat het goed. Laat ze hun pipelines bouwen. Maar, er is een maar: al onze programma's zijn technologisch goed in termen van het leerproces zelf. Voor het controleren van labwerken maken we gebruik van automatische checkers: de deelnemer moet inloggen op zijn persoonlijke account, op de knop 'Controleren' klikken en na een tijdje ziet hij een uitgebreide feedback over wat hij heeft gedaan. En precies op dat moment beginnen we ons probleem te benaderen.
De controle van dit lab is als volgt ingericht: we sturen een controledataset naar de Kafka van de deelnemer, vervolgens verplaatst Gobblin deze dataset naar HDFS, daarna haalt Airflow deze dataset op en plaatst deze in ClickHouse. Het ding is dat Airflow dit niet in realtime moet doen; het doet dit volgens een schema: elke 15 minuten haalt het een batch bestanden op en stopt deze.
Het blijkt dat we op de een of andere manier hun DAG zelf moeten triggeren op ons verzoek tijdens het werken met de checker hier en nu. Na wat googelen ontdekten we dat er voor latere versies van Airflow een zogenaamde bestaat. Het woord experimenteel, klinkt natuurlijk ontmoedigend, maar wat kunnen we doen... Misschien werkt het.
Hieronder beschrijven we het hele proces: van de installatie van Airflow tot het vormen van een POST-verzoek dat de DAG trigger, gebruikmakend van de Experimentale API. We zullen werken met Ubuntu 16.04.
1. Installatie van Airflow
Laten we controleren of we Python 3 en virtualenv hebben geïnstalleerd.
$ python3 --version
Python 3.6.6
$ virtualenv --version
15.2.0Als een van deze ontbreekt, installeer het dan.
Laten we nu een map aanmaken waar we verder zullen werken met Airflow.
$ mkdir
$ cd /pad/naar/je/nieuwe/directory
$ virtualenv -p which python3 venv
$ source venv/bin/activate
(venv) $Laten we Airflow installeren:
(venv) $ pip install airflowDe versie waarmee we werkten: 1.10.
Nu moeten we een map aanmaken airflow_home, waar de DAG-bestanden en Airflow-plugins zullen worden opgeslagen. Na het aanmaken van de map, zullen we de omgevingsvariabele instellen AIRFLOW_HOME.
(venv) $ cd /pad/naar/mijn/airflow/werkomgeving
(venv) $ mkdir airflow_home
(venv) $ export AIRFLOW_HOME=De volgende stap is het uitvoeren van het commando dat de gegevensstroomdatabase in SQLite aanmaakt en initialiseert:
(venv) $ airflow initdbDe database zal worden aangemaakt in airflow.db by default.
Laten we controleren of Airflow is geïnstalleerd:
$ 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.0Als het commando succesvol was, heeft Airflow zijn configuratiebestand aangemaakt airflow.cfg in AIRFLOW_HOME:
$ tree
.
├── airflow.cfg
└── unittests.cfgAirflow heeft een webinterface. Je kunt deze starten door het commando uit te voeren:
(venv) $ airflow webserver --port 8081Nu kun je de webinterface in de browser openen op poort 8081 op de host waar Airflow is gestart, bijvoorbeeld: <hostname:8081>.
2. Werken met de Experimentele API
Nu is Airflow ingesteld en klaar voor gebruik. Echter, we moeten ook de Experimentele API starten. Onze checkers zijn geschreven in Python, dus alle verzoeken hierna zullen hiermee gebeuren met de bibliotheek requests.
De API is overigens al actief voor eenvoudige verzoeken. Bijvoorbeeld, dit verzoek test of het werkt:
>>> import requests
>>> host =
>>> airflow_port = 8081 # in ons geval zo, standaard is 8080
>>> requests.get('http://{}:{}{}'.format(host, airflow_port, 'api/experimental/test')).text
'OK'Als je dit bericht als antwoord hebt gekregen, betekent dit dat alles werkt.
Echter, wanneer we de DAG willen triggeren, zullen we tegen het probleem aanlopen dat dit type verzoek niet zonder authenticatie kan gebeuren.
Daarvoor moeten we nog een aantal stappen ondernemen.
Allereerst moet je dit aan de configuratie toevoegen:
[api]
auth_backend = airflow.contrib.auth.backends.password_authVervolgens moet je een gebruiker met beheerdersrechten aanmaken:
>>> 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()Vervolgens moet je een gebruiker aanmaken met standaardrechten, die bevoegd is om de DAG-trigger te activeren.
>>> 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()Alles is nu klaar.
3. Start van de POST-aanroep
De POST-aanroep zal er als volgt uitzien:
>>> 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}'De aanvraag is succesvol verwerkt.
Daarom geven we de DAG enige tijd om te verwerken en doen we een verzoek aan de ClickHouse-tabel om een controlepakket gegevens te vangen.
Controle is voltooid.
Bron: habr.com
