Gjatë përgatitjes së programeve tona arsimore, herë pas here na ndodhin vështirësi në lidhje me përdorimin e disa mjeteve. Dhe në momentin kur përballemi me ta, nuk ka gjithmonë dokumentacion dhe artikuj të mjaftueshëm që mund të na ndihmojnë të zgjidhim këtë problem.
Kështu ndodhi, për shembull, në vitin 2015 dhe ne në programin "Specialist për të Dhënat e Mëdha" përdornim një kluster Hadoop me Spark për 35 përdorues njëkohësisht. Si ta përgatitnim atë për këtë rast përdorimi me YARN, nuk ishte e qartë. Në fund, duke u marrë me këtë dhe duke kaluar nëpër të gjithë procesin vetë, arritëm dhe gjithashtu të flasim në .
Historia e mëparshme
Këtë herë do të flasim për një program tjetër – . Në të, pjesëmarrësit tanë ndërtuan dy lloje arkitekturash: lambda dhe kappa. Në arkitekturën lambda, në kuadër të përpunimit batch, përdorim Airflow për të transferuar log-et nga HDFS në ClickHouse.
Gjithçka është mirë në përgjithësi. Le të ndjejnë ndërtimin e pipeline-ve të tyre. Sidoqoftë, ka një por: të gjitha programet tona janë teknologjike në lidhje me procesin e mësimit. Për verifikimin e lab-eve, ne përdorim kontrollet automatike: pjesëmarrësi duhet të hyjë në kabinetin e tij personal, të klikojë butonin "Verifiko", dhe pas një kohe të caktuar ai sheh ndonjë feedback të zgjeruar për atë që bëri. Dhe pikërisht në këtë moment fillojmë të përballemi me problemin tonë.
Kontrolli i këtij lab-i është i organizuar kështu: ne dërgojmë një paketë kontrolli të dhënash në Kafka për pjesëmarrësin, pastaj Gobblin e transferon këtë paketë të dhënash në HDFS, më pas Airflow merr këtë paketë të dhënash dhe e vendos në ClickHouse. Gjalpi është se Airflow nuk duhet ta bëjë këtë në real-time, ai e bën atë sipas një grafiku: çdo 15 minuta merr një sërë skedash dhe i dërgon.
Kështu që na nevojitet ndonjë mënyrë për të aktivizuar DAG-un e tyre vetvetiu sipas kërkesës sonë gjatë punës së kontroluesit këtu e tani. Pas kërkimit në Google, zbuluam se për versionet e vonshme të Airflow ekziston një . Fjala eksperimental, natyrisht, tingëllon disi frikësuese, por çfarë të bëjë... Ndoshta do të funksionojë.
Më vonë do të përshkruajmë të gjithë rrugën: nga instalimi i Airflow deri në formimin e një kërkese POST, e cila aktivizon DAG-un duke përdorur Experimental API. Do të punojmë me Ubuntu 16.04.
1. Instalimi i Airflow
Do të kontrollojmë se kemi Python 3 dhe virtualenv të instaluar.
$ python3 --version
Python 3.6.6
$ virtualenv --version
15.2.0Nëse ndonjëra prej tyre nuk është aty, instaloni.
Tani do të krijojmë një katalog ku do të vazhdojmë të punojmë me Airflow.
$ mkdir
$ cd /path/to/your/new/directory
$ virtualenv -p which python3 venv
$ source venv/bin/activate
(venv) $Le të instalojmë Airflow:
(venv) $ pip install airflowVersioni me të cilin punuam: 1.10.
Tani na nevojitet të krijojmë një katalog airflow_home, ku do të ruhen skedaret DAG dhe pllakos të Airflow. Pas krijimit të katalogut, do të vendosim variablën e mjedisit AIRFLOW_HOME.
(venv) $ cd /path/to/my/airflow/workspace
(venv) $ mkdir airflow_home
(venv) $ export AIRFLOW_HOME=Hapi tjetër është të ekzekutojmë komandën që do të krijojë dhe inicializojë bazën e të dhënave të rrjedhës në SQLite:
(venv) $ airflow initdbBaza e të dhënave do të krijohet në airflow.db si të paracaktuar.
Të kontrollojmë nëse Airflow u instalua:
$ airflow version
[2018-11-26 19:38:19,607] {__init__.py:57} INFO - Duke përdorur ekzekutorin SequentialExecutor
[2018-11-26 19:38:19,745] {driver.py:123} INFO - Duke gjeneruar tabela gramatikore nga /usr/lib/python3.6/lib2to3/Grammar.txt
[2018-11-26 19:38:19,771] {driver.py:123} INFO - Duke gjeneruar tabela gramatikore nga /usr/lib/python3.6/lib2to3/PatternGrammar.txt
____________ _____________
____ |__( )_________ __/__ /________ __
____ /| |_ /__ ___/_ /_ __ /_ __ _ | /| /
___ ___ | / _ / _ __/ _ / / \/_ |\/ |/
_/_/_ |_/_/_/ /_/_/ /_/_/ /_/_/ ____/____/__/
v1.10.0Nëse komanda është ekzekutuar, atëherë Airflow krijoi skedarin e tij të konfigurimit airflow.cfg në AIRFLOW_HOME:
$ tree
.
├── airflow.cfg
└── unittests.cfgAirflow ka një ndërfaqe në internet. Ajo mund të aktivizohet duke ekzekutuar komandën:
(venv) $ airflow webserver --port 8081Tani mund të hyni në ndërfaqen në internet në shfletuesin tuaj në portin 8081 në host-in ku Airflow është aktivizuar, për shembull: <hostname:8081>.
2. Puna me API-në Eksperimentale
Me këtë, Airflow është konfiguruar dhe i gatshëm për punë. Megjithatë, na nevojitet të aktivizojmë gjithashtu API-në Eksperimentale. Kontrolluesit tanë janë shkruar në Python, prandaj të gjitha kërkesat më poshtë do të jenë në të duke përdorur bibliotekën requests.
Në të vërtetë, API-ja është tashmë në funksion për kërkesa të thjeshta. Për shembull, një kërkesë e tillë lejon të testoni funksionimin e saj:
>>> import requests
>>> host =
>>> airflow_port = 8081 #në rastin tonë është kështu, ndërsa për paracaktim 8080
>>> requests.get('http://{}:{}{}'.format(host, airflow_port, '/api/experimental/test').text
'OK'Nëse keni marrë një mesazh të tillë në përgjigje, atëherë do të thotë se gjithçka funksionon.
Megjithatë, kur do të duam të aktivizojmë DAG, do të ngjitemi me faktin se ky lloj kërkese nuk mund të kryhet pa autentifikim.
Për këtë, do të nevojitet të përdorim disa hapa të tjerë.
Së pari, në konfigurim duhet të shtojmë këtë:
[api]
auth_backend = airflow.contrib.auth.backends.password_authPastaj, duhet të krijojmë përdoruesin tonë me të drejtat e administratorit:
>>> 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()Pastaj, duhet të krijoni një përdorues me të drejta të zakonshme, të cilit do t'i lejohet të bëjë trigger të DAG.
>>> 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()Tani gjithçka është gati.
3. Nisja e kërkesës POST
Kërkesa POST do të duket kështu:
>>> dag_id = newprolab
>>> url = 'http://{}:{}//api/experimental/dags/{}/dag_runs'.format(host, airflow_port, dag_id)
>>> 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'Kërkesa është përpunuar me sukses.
Për rrjedhojë, më pas i japim DAG-it disa kohë për të përpunuar dhe bëjmë një kërkesë në tabelën ClickHouse, duke u përpjekur të kapim një paketë kontrolli të të dhënave.
Kontrolli është përfunduar.
Burimi: habr.com
