Kuidas luua Airflow'is DAG'i kÀivitust, kasutades Eksperimentaalset API-d

Koolituste valmistamisel puutume aeg-ajalt kokku probleemidega, mis on seotud mÔnede tööriistadega. Ja sel hetkel, kui nendega silmitsi seisame, ei ole alati piisavalt dokumentatsiooni ja artikleid, mis aitaksid selle probleemiga toime tulla.

Nii juhtus nÀiteks 2015. aastal, kui kasutasime suure andmehulga spetsialisti programmis Hadoopi klastrit koos Sparkiga 35 samaaegse kasutaja jaoks. Kuidas sellise kasutusjuhtumi jaoks YARN-iga töötamiseks ette valmistada, oli arusaamatu. LÔpuks, pÀrast iseenda uurimist, tegime postituse Habras ja esinesime ka Moscow Spark Meetupil.

Eellugu

Sel korral rÀÀgime teisest programmist – Data Engineer. Sellel kursusel ehitavad meie osalejad kahte tĂŒĂŒpi arhitektuuri: lambda ja kappa. Ja lambda-arhitektuuris kasutatakse partii töötlemise raames logide edastamiseks HDFS-ist ClickHouse'i Airflow't.

KĂ”ik on ĂŒldiselt hĂ€sti. Las nad ehitavad oma torujuhtmeid. Siiski on ĂŒks aga: kĂ”ik meie programmid on tehnoloogilised, arvestades Ă”ppimisprotsessi. Laborite kontrollimiseks kasutame automaatseid kontrollijaid: osaleja peab minema isiklikku kabinetti, vajutama nuppu "Kontrolli" ja teatud aja pĂ€rast nĂ€eb ta mingit laiendatud tagasisidet selle kohta, mida ta tegi. Ja just sel hetkel hakkame lĂ€henema meie probleemile.

Selle labori kontroll toimub nii: saadame osaleja Kafka kaudu kontrollpaketi andmeid, seejÀrel Gobblin edastab selle andmepaketi HDFS-ile, seejÀrel Airflow vÔtab selle andmepaketi ja paneb selle ClickHouse'i. Asi on selles, et Airflow ei pea seda tegema reaalajas, ta teeb seda graafiku jÀrgi: kord 15 minuti jooksul vÔtab ta hulga faile ja edastab need.

Tuleb vĂ€lja, et peame kuidagi kĂ€ivitama nende DAGi iseseisvalt meie nĂ”udmisel kontrollija töö ajal siin ja praegu. Otsides selgitasime vĂ€lja, et hilisemate versioonide korral on Airflow'l nii nimetatud Experimental API. SĂ”na experimental, muidugi, kĂ”lab hirmutatavalt, aga mis seal ikka... Äkki lĂ€heb korda.

Edasi kirjeldame kogu teekonda: alates Airflow'i installimisest kuni POST-pÀringu koostamiseni, mis kÀivitab DAGi, kasutades Experimental API-d. Töötame Ubuntu 16.04-ga.

1. Airflow'i installimine

Kontrollime, kas meil on installitud Python 3 ja virtualenv.

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

Kui ĂŒhtegi neist pole, siis Installige need.

NĂŒĂŒd loome katalooge, kus töötame edaspidi Airflow'ga.

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

Paigaldame Airflow:

(venv) $ pip install airflow

Versioon, millega me töötasime: 1.10.

NĂŒĂŒd peame looma katalooge airflow_home, kuhu paigutatakse DAG-failid ja Airflowi pistikprogrammid. PĂ€rast katalooge loomist seadistame keskkonnamuutuja AIRFLOW_HOME.

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

JÀrgmiseks sammuks kÀivitame kÀsu, mis loob ja initsialiseerib andmebaasi SQLite-s:

(venv) $ airflow initdb

Andmebaas luuakse airflow.db vaikimisi.

Kontrollime, kas Airflow on installitud:

$ 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 töötas, siis on Airflow loonud oma konfiguratsioonifaili airflow.cfg ja AIRFLOW_HOME:

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

Airflowil on veebiliides. Selle kÀivitamiseks kasutage kÀsku:

(venv) $ airflow webserver --port 8081

NĂŒĂŒd pÀÀsete veebiliidesele brauseris port 8081 kaudu masinas, kus Airflow kĂ€ivitati, nĂ€iteks: <hostname:8081>.

2. Töö eksperimenteerimise API-ga

Sellega on Airflow seadistatud ja valmis töötama. Siiski peame kÀivitama ka eksperimenteerimise API. Meie kontrollijad on kirjutatud Pythonis, seega on kÔik edasised pÀringud kirjutatud selles, kasutades teeki taotlused.

Tegelikult töötab API juba lihtsate pÀringute jaoks. NÀiteks selline pÀring vÔimaldab testida selle tööd:

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

Kui saite sellise vastuse, tÀhendab see, et kÔik töötab.

Kuid kui me tahame DAG-i vĂ€lja kutsuda, siis seisame silmitsi faktiga, et seda tĂŒĂŒpi pĂ€ringut pole vĂ”imalik teha ilma autentimiseta.

Selle jaoks tuleb teha veel paar sammu.

Esiteks tuleb konfi lisada see:

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

SeejÀrel tuleb luua oma kasutaja administraatori Ô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 tavakasutaja, kellel on lubatud DAG-i kÀivitamine.

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

POST-pÀring nÀeb vÀlja jÀrgmine:

>>> 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 on edukalt töödeldud.

SeetĂ”ttu anname DAG-ile veidi aega töötlemiseks ja teeme pĂ€ringu ClickHouse tabelisse, pĂŒĂŒdes jÀÀda kontrollpakettide juurde.

Kontrollimine on lÔpetatud.

Allikas: habr.com

Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid | ProHoster