Kuidas luua DAG’i käivitajat Airflow’s, kasutades Experimental API-d

Meie haridusprogrammide ettevalmistamisel seisame aeg-ajalt silmitsi keerukustega teatud tööriistade kasutamisel. Ja sel hetkel, kui nendega kokku puutume, ei ole alati piisavalt dokumentatsiooni ja artikleid, mis aitaksid selle probleemiga tegeleda.

Nii oli näiteks 2015. aastal, mil meie programm "Suured andmed" kasutas Hadoopi klusterit Sparkiga 35 samaaegse kasutajaga. Kuidas seda sellise kasutusjuhtumi jaoks YARN-i abil ette valmistada, ei olnud selge. Lõpuks, pärast iseseisvat lahenduse leidmist, tegime postituse Habré ja esinesime samuti Moscow Spark Meetup.

Eelalugu

Seekord räägime teisest programmist – Andmeinsener. Selle raames ehitavad meie osalejad kahte tüüpi arhitektuuri: lambda ja kappa. Lambda-arhitektuuris kasutatakse batch-töötlemise kontekstis Airflowd logide edastamiseks HDFS-ist ClickHouse'i.

Üldiselt on kõik hästi. Las nad ehitavad oma vooge. Üks asi on aga: kõik meie programmid on tehnilised ka õppimise protsessi mõttes. Laborite kontrollimiseks kasutame automaatseid kontrollijaid: osaleja peab sisenema isiklikku kabinetti, vajutama nuppu "Kontrolli" ja mõne aja pärast näeb ta tagasisidet selle kohta, mida tegi. Just sel hetkel hakkame oma probleemile lähenema.

Selle labori kontrollimine on korraldatud nii: me saatame kontrollpaketi osaleja Kafka'sse, seejärel Gobblin edastab selle andmepaketi HDFS-i, siis Airflow võtab selle andmepaketi ja asetab selle ClickHouse'i. Probleem on selles, et Airflow ei pea seda reaalajas tegema, vaid plaani järgi: iga 15 minuti tagant võtab ta hunniku faile ja laadib need üles.

See tähendab, et me peame kuidagi nende DAG-i iseseisvalt aktiveerima vastavalt meie nõudmisele seadistaja töö ajal. Otsides leidsime, et hilisematel Airflow versioonidel on nii nimetatud Eksperimentaalne API. Sõna eksperimentaalne, muidugi, kõlab hirmutavalt, aga mis teha… Äkki õnnestub.

Edasi kirjeldame kogu teed: alates Airflow paigaldamisest kuni POST-päringu loomisega, mis aktiveerib DAG-i, kasutades Experimental API-d. Töötame Ubuntu 16.04 süsteemiga.

1. Airflow paigaldamine

Kontrollime, kas meil on Python 3 ja virtualenv.

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

Kui midagi on puudu, siis installige.

Nüüd loome katalooge, kus hakkame Airflow-ga edasi töötama.

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

Installime Airflow:

(venv) $ pip install airflow

Versioon, millega töötasime: 1.10.

Nüüd peame looma katalooge airflow_home, kuhu paigutatakse DAG-failid ja Airflow pluginad. Pärast katalooge loomist seadistame keskkonnamuutuja AIRFLOW_HOME.

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

Järgmine samm on käivitada käsk, mis loob ja algatab andmevoo andmebaasi SQLite-s:

(venv) $ airflow initdb

Andmebaas luuakse airflow.db vaikimisi.

Kontrollime, kas Airflow on edukalt paigaldatud:

$ 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

Kui käsk toimis, siis Airflow loo oma konfiguratsioonifaili airflow.cfg ühes AIRFLOW_HOME:

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

Airflow'il on veebiliides. Selle saab käivitada käsu abil:

(venv) $ airflow webserver --port 8081

Nüüd pääsete veebiliidesele brauseris, kasutades porti 8081 hostis, kus Airflow käivitatud, näiteks: <hostname:8081>.

2. Töö Experimental API-ga

Nüüd on Airflow seadistatud ja valmis tööks. Siiski peame käivitama ka Experimental API. Meie kontrollid on kirjutatud Pythonis, seega on kõik järgmised päringud selles kasutades teeki requests.

Tegelikult töötab API juba lihtsate päringute jaoks. Näiteks järgmine päring võimaldab selle toimimist testida:

>>> import requests
>>> host = 
>>> airflow_port = 8081 # meie puhul selline, default on 8080
>>> requests.get('http://{}:{} /{}'.format(host, airflow_port, 'api/experimental/test')).text
'OK'

Kui olete saanud sellise vastuse, tähendab see, et kõik töötab.

Kuid kui me soovime DAG-i aktiveerida, seisame silmitsi sellega, et taolist päringut ei saa teha ilma autentimiseta.

Selleks tuleb teha veel mõned toimingud.

Esiteks peab konfiguratsiooni lisama järgmise:

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

Seejärel tuleb luua oma kasutaja admin õigustega:

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

Seejärel tuleb luua kasutaja tavapäraste õigustega, kellele on lubatud DAG-i aktiveerida.

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

Nüüd on kõik valmis.

3. POST-päringu käivitamine

Ise POST-päring näeb välja selline:

>>> 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'

Päring töötles edukalt.

Seega anname DAG-ile töötlemiseks aega ja teeme päringu ClickHouse tabelisse, püüdes tabada kontrollpaketti.

Kontrollimine on lõppenud.

Allikas: habr.com

Osta usaldusväärne veebihosting DDoS kaitsega, VPS VDS serverid 🔥 Osta usaldusväärne veebihosting DDoS kaitsega, VPS VDS serverid | ProHoster