În pregătirea programelor noastre educaționale, ne confruntăm periodic cu dificultăți în utilizarea unor instrumente. Și în momentul în care ne întâlnim cu ele, nu întotdeauna există suficiente documentații și articole care să ne ajute să rezolvăm această problemă.
Așa a fost, de exemplu, în 2015, când în programul „Specialist în date mari” am utilizat un cluster Hadoop cu Spark pentru 35 de utilizatori simultan. Cum să-l configurăm pentru un astfel de caz de utilizare cu YARN, era neclar. În final, după ce ne-am descurcat și am parcurs singuri drumul, am realizat și am participat la .
Povestea
De data aceasta, vom vorbi despre un alt program – . În cadrul acestuia, participanții noștri construiesc două tipuri de arhitectură: lambda și kappa. Iar în arhitectura lambda, în cadrul procesării pe loturi, se folosește Airflow pentru a transfera jurnalele din HDFS în ClickHouse.
În general, totul este bine. Să-și construiască propriile canalizări. Însă, există o problemă: toate programele noastre sunt tehnologice în ceea ce privește procesul de învățare. Pentru verificarea laboratoarelor, folosim verificatori automați: participanții trebuie să se conecteze în contul personal, să apese butonul „Verifică” și, după un timp, să primească un feedback extins asupra a ceea ce au făcut. Și exact în acel moment începem să ne apropiem de problema noastră.
Verificarea acestui laborator este organizată astfel: trimitem un pachet de date de control în Kafka participantului, apoi Gobblin transferă acest pachet de date în HDFS, iar Airflow preia acest pachet de date și îl plasează în ClickHouse. Ideea este că Airflow nu trebuie să facă asta în timp real, ci o face conform unui program: odată la 15 minute ia un set de fișiere și le încarcă.
Așadar, trebuie să găsim o modalitate de a declanșa DAG-ul lor pe cont propriu, la cererea noastră, în timpul lucrului verificatorului aici și acum. Căutând pe Google, am aflat că pentru versiunile mai recente de Airflow există așa-numitul . Cuvântul experimental, desigur, sună înfricoșător, dar ce să facem... Poate va funcționa.
Mai departe, vom descrie întregul proces: de la instalarea Airflow până la formarea cererii POST care declanșează DAG-ul, folosind Experimental API. Vom lucra cu Ubuntu 16.04.
1. Instalarea Airflow
Vom verifica dacă avem instalat Python 3 și virtualenv.
$ python3 --version
Python 3.6.6
$ virtualenv --version
15.2.0Dacă ceva din acestea nu este instalat, vă rugăm să le instalați.
Acum vom crea un director în care vom continua să lucrăm cu Airflow.
$ mkdir
$ cd /path/to/your/new/directory
$ virtualenv -p which python3 venv
$ source venv/bin/activate
(venv) $Să instalăm Airflow:
(venv) $ pip install airflowVersiunea cu care am lucrat: 1.10.
Acum trebuie să creăm un director airflow_home, unde vor fi plasate fișierele DAG și pluginurile Airflow. După ce am creat directorul, vom seta variabila de mediu AIRFLOW_HOME.
(venv) $ cd /path/to/my/airflow/workspace
(venv) $ mkdir airflow_home
(venv) $ export AIRFLOW_HOME=Următorul pas este să executăm comanda care va crea și inițializa baza de date a fluxului de date în SQLite:
(venv) $ airflow initdbBaza de date va fi creată în airflow.db implicit.
Să verificăm dacă Airflow s-a instalat:
$ airflow version
[2018-11-26 19:38:19,607] {__init__.py:57} INFO - Folosind executorul SequentialExecutor
[2018-11-26 19:38:19,745] {driver.py:123} INFO - Generând tabele de gramatică din /usr/lib/python3.6/lib2to3/Grammar.txt
[2018-11-26 19:38:19,771] {driver.py:123} INFO - Generând tabele de gramatică din /usr/lib/python3.6/lib2to3/PatternGrammar.txt
____________ _____________
____ |__( )_________ __/__ /________ __
____ /| |_ /__ ___/_ /_ __ /_ __ _ | /| /
___ ___ | / _ / _ __/ _ / / \/_\/ _ |/ |/
_/_/ |_/_/ /_\/ /_/ /_\/ ____/____/|__/
v1.10.0Dacă comanda a funcționat, atunci Airflow a creat fișierul său de configurare airflow.cfg în AIRFLOW_HOME:
$ tree
.
├── airflow.cfg
└── unittests.cfgAirflow are o interfață web. Aceasta poate fi pornită prin executarea comenzii:
(venv) $ airflow webserver --port 8081Acum poți accesa interfața web în browser pe portul 8081 pe gazda unde a fost pornit Airflow, de exemplu: <hostname:8081>.
2. Lucrând cu API Experimental
Aici Airflow este configurat și pregătit pentru utilizare. Cu toate acestea, trebuie să pornim și API-ul Experimental. Checker-ele noastre sunt scrise în Python, așadar, toate cererile vor fi în acest limbaj folosind biblioteca requests..
De fapt, API-ul este deja funcțional pentru cereri simple. De exemplu, această cerere permite testarea funcționalității sale:
>>> import requests
>>> host =
>>> airflow_port = 8081 # în cazul nostru acesta, iar în mod implicit 8080
>>> requests.get('http://{}:{}//{}'.format(host, airflow_port, 'api/experimental/test').text
'OK'Dacă ai primit acest mesaj în răspuns, înseamnă că totul funcționează.
Cu toate acestea, când dorim să declanșăm DAG-ul, ne vom confrunta cu faptul că acest tip de cerere nu poate fi efectuat fără autentificare.
Pentru aceasta va fi necesar să parcurgem câțiva pași suplimentari.
În primul rând, în configurație trebuie să adăugăm acest lucru:
[api]
auth_backend = airflow.contrib.auth.backends.password_authApoi, este necesar să creăm un utilizator propriu cu drepturi de administrator:
>>> 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()Apoi, trebuie creat un utilizator cu drepturi normale, căruia îi va fi permis să declanșeze DAG-ul.
>>> 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()Acum totul este pregătit.
3. Inițierea cererii POST
Cererile POST vor arăta astfel:
>>> 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'Cererea a fost procesată cu succes.
Prin urmare, mai departe îi oferim un timp DAG-ului pentru a procesa și facem o cerere în tabelul ClickHouse, încercând să prindem un pachet de date de control.
Verificarea a fost finalizată.
Sursa: habr.com
