Si të krijosh një trigger DAG në Airflow, duke përdorur API-në Eksperimentale

Në përgatitjen e programeve tona edukative, herë pas here hasim vështirësi në lidhje me përdorimin e disa mjeteve. Dhe në momentin kur përballemi me to, nuk ka gjithmonë dokumentacion të mjaftueshëm dhe artikuj që do të ndihmonin për të zgjidhur këtë problem.

KĂ«shtu ndodhi, pĂ«r shembull, nĂ« vitin 2015 dhe ne nĂ« programin “Specialist pĂ«r tĂ« dhĂ«na tĂ« mĂ«dha” pĂ«rdornim njĂ« klaster Hadoop me Spark pĂ«r 35 pĂ«rdorues tĂ« njĂ«kohshĂ«m. Si ta pĂ«rgatitim atĂ« pĂ«r njĂ« rast pĂ«rdorimi tĂ« tillĂ« me pĂ«rdorimin e YARN, nuk ishte e qartĂ«. Duke iu kthyer mbrapa dhe duke kaluar rrugĂ«n vetĂ«, arritĂ«m post tĂ« gĂ«zueshĂ«m nĂ« Habr. dhe madje u paraqitĂ«m nĂ« Moscow Spark Meetup.

Pas historia

Tani do tĂ« flasim pĂ«r njĂ« program tjetĂ«r – Data Engineer. NĂ« tĂ«, pjesĂ«marrĂ«sit tanĂ« ndĂ«rtuan dy lloje arkitekture: lambda dhe kappa. Dhe nĂ« arkitekturĂ«n lambda, pĂ«r procesimin e grumbujve, pĂ«rdoret Airflow pĂ«r tĂ« transferuar regjistrat nga HDFS nĂ« ClickHouse.

Gjithçka Ă«shtĂ« mirĂ«. Le tĂ« ndĂ«rtojnĂ« tubacionet e tyre. Por ka njĂ« por: tĂ« gjithĂ« programet tona janĂ« teknologjikĂ« nĂ« lidhje me vetĂ« procesin e edukimit. PĂ«r tĂ« verifikuar laboratorĂ«t pĂ«rdorim kontrollues automatikĂ«: pjesĂ«marrĂ«si duhet tĂ« hyjĂ« nĂ« llogarinĂ« e tij, tĂ« klikohet nĂ« butonin “Verifiko”, dhe pas njĂ« kohe ai sheh njĂ« reagim tĂ« zgjeruar mbi atĂ« qĂ« bĂ«ri. Dhe nĂ« atĂ« moment ne fillojmĂ« tĂ« afrojmĂ« nĂ« problemin tonĂ«.

Verifikimi i këtij laboratori është i strukturuar kështu: ne dërgojmë një paketë kontrolli të dhënash në Kafka të pjesëmarrësit, pastaj Gobblin transferon këtë paketë të dhënash në HDFS, më pas Airflow merr këtë paketë të dhënash dhe e vendos në ClickHouse. Meraku është se Airflow nuk duhet ta bëjë këtë në kohë reale, e bën këtë sipas një programi: çdo 15 minuta merr një grup skedarësh dhe i dërgon.

Pra, na nevojitet tĂ« aktivizojmĂ« DAG-un e tyre vetĂ« sipas kĂ«rkesĂ«s sonĂ« nĂ« kohĂ«n kur punon kontrolluesi kĂ«tu dhe tani. Pasi e kĂ«rkuam nĂ« Google, zbuluam se pĂ«r versionet e vonshme tĂ« Airflow ekziston njĂ« e ashtuquajtur Experimental API. Fjala experimental, natyrisht, tingĂ«llon frikshĂ«m, por çfarĂ« tĂ« bĂ«jmë  Ndoshta do tĂ« ketĂ« sukses.

Më pas do të përshkruajmë të gjithë rrugën: nga instalimi i Airflow deri te formimi i kërkesës POST, e cila aktivizon DAG-un, duke përdorur Experimental API. Do të punojmë me Ubuntu 16.04.

1. Instalimi i Airflow

Le të kontrollojmë nëse kemi Python 3 dhe virtualenv.

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

Nëse ndonjëra nga këto nuk është e instaluar, atëherë instaloni.

Tani le të krijojmë një katalog, ku do të punojmë më tej 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 airflow

Versioni, me të cilin punuam, ishte: 1.10.

Tani na nevojitet të krijojmë një katalog airflow_home, ku do të vendosen skedarët e DAG-ve dhe pluginet e Airflow. Pasi ta krijojmë katalogun, le të vendosim variablin mjedisor AIRFLOW_HOME.

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

Hapi tjetĂ«r – tĂ« ekzekutojmĂ« komandĂ«n qĂ« do tĂ« krijojĂ« dhe inicializojĂ« bazĂ«n e tĂ« dhĂ«nave tĂ« rrjedhĂ«s nĂ« SQLite:

(venv) $ airflow initdb

Baza e të dhënave do të krijohet në airflow.db si dyshim.

Le të kontrollojmë nëse është instaluar Airflow:

$ 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

Nëse komanda u ekzekutua, atëherë Airflow krijoi skedarin e tij të konfigurimit airflow.cfg në AIRFLOW_HOME:

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

Airflow ka një ndërfaqe të uebit. Ajo mund të nisë, duke ekzekutuar komandën:

(venv) $ airflow webserver --port 8081

Tani mund të hyni në ndërfaqen e uebit në shfletues në portin 8081 në hostin ku Airflow u nisi, për shembull: <hostname:8081>.

2. Puna me Experimental API

Kështu, Airflow është i konfiguruar dhe i gatshëm për punë. Megjithatë, na nevojitet të aktivizojmë gjithashtu Experimental API. Kontrolluesit tanë janë shkruar në Python, kështu që më poshtë do të jenë të gjitha kërkesat në të, duke përdorur bibliotekën requests.

Në të vërtetë, API tashmë funksionon për kërkesa të thjeshta. Për shembull, kjo kërkesë lejon të testojmë funksionimin e tij:

>>> import requests
>>> host = 
>>> airflow_port = 8081 #në rastin tonë të tillë, ndryshe 8080
>>> requests.get('http://{}:{} /{}'.format(host, airflow_port, 'api/experimental/test')).text
'OK'

Nëse kemi marrë një mesazh të tillë në përgjigje, atëherë kjo do të thotë se gjithçka funksionon.

Megjithatë, kur dëshirojmë të aktivizojmë DAG-un, do të përballemi me faktin se ky lloj kërkese nuk mund të bëhet pa autentifikim.

Për këtë, do të duhet të bëjmë disa hapa shtesë.

Së pari, në konfigurim duhet të shtojmë këtë:

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

Pastaj, na nevojitet të krijojmë një përdorues me drejtë admin:

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

Më pas, duhet të krijoni një përdorues me të drejta normale, i cili do të ketë leje për të aktivizuar DAGun.

>>> 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. Aktivizimi i kërkesës POST

Kërkesa POST do të duket kështu:

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

Kërkesa është procesuar me sukses.

Në përputhje me këtë, ne i japim pak kohë DAG-ut për të procesuar dhe bëjmë një kërkesë në tabelën ClickHouse, duke u përpjekur të kapim paketën e kontrolit të të dhënave.

Kontrollimi përfundoi.

Burimi: habr.com

Bleni hostim tĂ« besueshĂ«m pĂ«r faqe me mbrojtje nga DDoS, serverĂ« VPS VDS đŸ”„ Bleni hostim tĂ« besueshĂ«m pĂ«r faqe me mbrojtje nga DDoS, serverĂ« VPS VDS | ProHoster