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.cfgAIRFLOW_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