Cum să creezi un trigger DAG în Airflow folosind Experimental API

Î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 un articol pe Habr și am participat la Moscow Spark Meetup.

Povestea

De data aceasta, vom vorbi despre un alt program – Inginer de date. Î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 Experimental API. 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.0

Dacă 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 airflow

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

Baza 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.0

Dacă comanda a funcționat, atunci Airflow a creat fișierul său de configurare airflow.cfg în AIRFLOW_HOME:

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

Airflow are o interfață web. Aceasta poate fi pornită prin executarea comenzii:

(venv) $ airflow webserver --port 8081

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

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

Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS 🔥 Cumpără un hosting fiabil pentru site-uri cu protecție DDoS, servere VPS VDS | ProHoster