Kā izveidot DAG aktivizētāju programmā Airflow, izmantojot eksperimentālo API

Sagatavojot savas izglītības programmas, mēs periodiski saskaramies ar grūtībām darbā ar dažiem rīkiem. Un brīdī, kad ar tiem sastopamies, ne vienmēr ir pietiekami daudz dokumentācijas un rakstu, kas palīdzētu tikt galā ar šo problēmu.

Tā tas bija, piemēram, 2015. gadā, un mēs izmantojām Hadoop kopu ar Spark 35 vienlaicīgiem lietotājiem programmā Big Data Specialist. Nebija skaidrs, kā to sagatavot šādam lietotāja gadījumam, izmantojot DZIJU. Rezultātā, paši izdomājuši un izstaigājuši ceļu, viņi to arī izdarīja ieraksts vietnē Habré un arī uzstājās Maskavas dzirksteles tikšanās.

Aizvēsture

Šoreiz mēs runāsim par citu programmu - Datu inženieris. Uz tā mūsu dalībnieki veido divu veidu arhitektūru: lambda un kappa. Un lamdba arhitektūrā Airflow tiek izmantota kā daļa no pakešapstrādes, lai pārsūtītu žurnālus no HDFS uz ClickHouse.

Kopumā viss ir kārtībā. Ļaujiet viņiem pašiem veidot savus mācību procesus. Tomēr ir viens "bet": visas mūsu programmas ir tehnoloģiski attīstītas pašā mācību procesā. Lai pārbaudītu laboratorijas rezultātus, mēs izmantojam automātiskos pārbaudītājus: dalībniekam jāpiesakās savā personīgajā kontā, jānoklikšķina uz pogas "Pārbaudīt", un pēc brīža viņš redz paplašinātu atsauksmi par savu darbu. Un tieši šajā brīdī mēs sākam risināt savu problēmu.

Šīs laboratorijas pārbaude tiek sakārtota šādi: mēs nosūtām kontroles datu paketi dalībnieka Kafka, tad Gobblin pārsūta šo datu paketi uz HDFS, pēc tam Airflow paņem šo datu paketi un ievieto to ClickHouse. Viltība ir tāda, ka Airflow tas nav jādara reāllaikā, tas dara to pēc grafika: reizi 15 minūtēs tas aizņem vairākus failus un augšupielādē tos.

Izrādās, ka mums pašiem pēc mūsu pieprasījuma kaut kā jāiedarbina viņu DAG, kamēr pārbaudītājs darbojas šeit un tagad. Googlējot noskaidrojām, ka jaunākām Airflow versijām ir t.s Eksperimentālā API. Vārds experimental, protams, izklausās biedējoši, bet ko lai dara... Pēkšņi paceļas.

Tālāk mēs aprakstīsim visu procesu: no Airflow instalēšanas līdz POST pieprasījuma izveidei, kas aktivizē DAG, izmantojot eksperimentālo API. Mēs strādāsim ar Ubuntu 16.04.

1. Gaisa plūsmas uzstādīšana

Pārbaudīsim, vai mums ir Python 3 un virtualenv.

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

Ja kāda no tām trūkst, instalējiet to.

Tagad izveidosim direktoriju, kurā turpināsim strādāt ar Airflow.

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

Instalējiet gaisa plūsmu:

(venv) $ pip install airflow

Versija, pie kuras strādājām: 1.10.

Tagad mums ir jāizveido direktorijs airflow_home, kur atradīsies DAG faili un Airflow spraudņi. Pēc direktorija izveides iestatiet vides mainīgo AIRFLOW_HOME.

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

Nākamais solis ir palaist komandu, kas izveidos un inicializēs datu plūsmas datu bāzi programmā SQLite:

(venv) $ airflow initdb

Datubāze tiks izveidota airflow.db noklusējuma.

Pārbaudiet, vai ir uzstādīta gaisa plūsma:

$ 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

Ja komanda darbojās, Airflow izveidoja savu konfigurācijas failu airflow.cfg в AIRFLOW_HOME:

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

Airflow ir tīmekļa saskarne. To var palaist, izpildot komandu:

(venv) $ airflow webserver --port 8081

Tagad varat piekļūt tīmekļa saskarnei pārlūkprogrammā, kas atrodas resursdatorā, kurā darbojās Airflow, portā 8081, piemēram: <hostname:8081>.

2. Darbs ar eksperimentālo API

Šajā gaisa plūsma ir konfigurēta un gatava darbam. Tomēr mums ir jāpalaiž arī eksperimentālā API. Mūsu dambrete ir rakstīta Python valodā, tāpēc turpmāk visi pieprasījumi tiks uz to, izmantojot bibliotēku requests.

Faktiski API jau darbojas vienkāršiem pieprasījumiem. Piemēram, šāds pieprasījums ļauj pārbaudīt tā darbu:

>>> import requests
>>> host = <your hostname>
>>> airflow_port = 8081 #в нашем случае такой, а по дефолту 8080
>>> requests.get('http://{}:{}/{}'.format(host, airflow_port, 'api/experimental/test').text
'OK'

Ja atbildē saņēmāt šādu ziņojumu, tas nozīmē, ka viss darbojas.

Tomēr, kad vēlamies aktivizēt DAG, mēs saskaramies ar faktu, ka šāda veida pieprasījumu nevar veikt bez autentifikācijas.

Lai to izdarītu, jums būs jāveic vairākas darbības.

Pirmkārt, jums tas jāpievieno konfigurācijai:

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

Pēc tam jums ir jāizveido savs lietotājs ar administratora tiesībām:

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

Tālāk jums ir jāizveido lietotājs ar normālām tiesībām, kam būs atļauts veikt DAG trigeri.

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

Tagad viss ir gatavs.

3. POST pieprasījuma palaišana

Pats POST pieprasījums izskatīsies šādi:

>>> 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 <DagRun newprolab @ 2019-03-27 10:24:25+00:00: manual__2019-03-27T10:24:25+00:00, externally triggered: True>"n}n'

Pieprasījums veiksmīgi apstrādāts.

Attiecīgi mēs dodam DAG kādu laiku apstrādāt un iesniegt pieprasījumu ClickHouse tabulai, mēģinot noķert kontroles datu paketi.

Verifikācija pabeigta.

Avots: www.habr.com

Iegādājieties uzticamu mitināšanu vietnēm ar DDoS aizsardzību, VPS VDS serveriem 🔥 Iegādājieties uzticamu tīmekļa vietņu mitināšanu ar DDoS aizsardzību, VPS VDS serveriem | ProHoster