
Përshëndetje, Habr! Në këtë artikull do të flas për një mjet të shkëlqyer për zhvillimin e proceseve batch për përpunimin e të dhënave, për shembull, në infrastrukturën e DWH të korporatave ose të DataLake tuaj. Tema do të jetë Apache Airflow (në vazhdim Airflow). Ai është ndoshta tërhequr më pak vëmendje në Habr dhe në pjesën kryesore do të përpiqem t'ju bind se së paku duhet të shikoni Airflow kur zgjidhni një planifikues për proceset tuaja ETL/ELT.
Më parë kam shkruar një seri artikujsh mbi temën e DWH, kur punoja në Bankën Tinkoff. Tani kam bërë pjesë të ekipit Mail.Ru Group dhe jam duke u marrë me zhvillimin e platformës për analizën e të dhënave në drejtimin e lojërave. Ashtu siç do të shfaqen lajme dhe zgjidhje interesante, ne me ekipin do të tregojmë këtu për platformën tonë për analizën e të dhënave.
Prologu
Pra, le tĂ« fillojmĂ«. ĂfarĂ« Ă«shtĂ« Airflow? Ky Ă«shtĂ« njĂ« bibliotekĂ« (ose ) pĂ«r zhvillimin, planifikimin dhe monitorimin e proceseve tĂ« punĂ«s. Karakteristika kryesore e Airflow: pĂ«r pĂ«rshkrimin (zhvillimin) e proceseve pĂ«rdoret kodi nĂ« gjuhĂ«n Python. Kjo sjell shumĂ« pĂ«rparĂ«si pĂ«r organizimin e projektit tuaj dhe zhvillimin: nĂ« thelb, projekti juaj (p.sh.) ETL Ă«shtĂ« thjesht njĂ« projekt nĂ« Python, dhe ju mund ta organizoni si tĂ« doni, duke marrĂ« parasysh karakteristikat e infrastrukturĂ«s, madhĂ«sinĂ« e ekipit dhe kĂ«rkesa tĂ« tjera. gjithçka Ă«shtĂ« e thjeshtĂ« nga ana instrumentale. PĂ«rdorni, pĂ«r shembull, PyCharm + Git. Kjo Ă«shtĂ« e shkĂ«lqyer dhe shumĂ« e pĂ«rshtatshme!
Tani le të shqyrtojmë entitetet kryesore të Airflow. Duke kuptuar natyrën dhe qëllimin e tyre, do të organizoni optimalisht arkitekturën e proceseve. Ndoshta, entiteti kryesor është Directed Acyclic Graph (në vazhdim DAG).
DAG
DAG është një bashkim kuptimor i detyrave tuaja që dëshironi të kryeni në një rend të caktuar sipas një orari të caktuar. Airflow ofron një ndërfaqe të përshtatshme web për punën me DAG'ët dhe entitete të tjera:

DAG mund të duket kështu:

Një zhvillues, duke dizajnuar DAG, përcakton një grup operatorësh, mbi të cilët do të ndërtohen detyrat brenda DAG-ut. Këtu arrijmë në një entitet tjetër të rëndësishëm: Operatorin e Airflow.
Operatorët
Operatori është një entitet mbi bazën e të cilit krijohen instanca të detyrave, ku përshkruhet se çfarë do të ndodhë gjatë ekzekutimit të një instance të detyrës. mbledhin tashmë një grup operatorësh, të gatshëm për t'u përdorur. Shembuj:
- BashOperator â operator pĂ«r ekzekutimin e njĂ« komande bash.
- PythonOperator â operator pĂ«r thirrjen e kodit Python.
- EmailOperator â operator pĂ«r dĂ«rgimin e njĂ« emaili.
- HTTPOperator â operator pĂ«r punĂ«n me kĂ«rkesat http.
- SqlOperator â operator pĂ«r ekzekutimin e kodit SQL.
- Sensor â operator pĂ«r pritjen e njĂ« ngjarjeje (pĂ«rmbyllje tĂ« njĂ« kohĂ« tĂ« caktuar, shfaqje tĂ« njĂ« skedari, njĂ« rresht nĂ« bazĂ«n e tĂ« dhĂ«nave, pĂ«rgjigje nga API â etj.).
Ka operatorë më specifikë: DockerOperator, HiveOperator, S3FileTransferOperator, PrestoToMysqlOperator, SlackOperator.
Ju gjithashtu mund të zhvilloni operatorë duke u bazuar në karakteristikat tuaja dhe t'i përdorni ato në projekt. Për shembull, ne krijuam MongoDBToHiveViaHdfsTransfer, operatori për eksportimin e dokumenteve nga MongoDB në Hive, si dhe disa operatorë për të punuar me : CHLoadFromHiveOperator dhe CHTableLoaderOperator. Në thelb, sa më shpejt që në projekt të shfaqet një kod i përdorur shpesh, i ndërtuar mbi operatorët bazë, duhet të mendoni për ta grumbulluar atë në një operator të ri. Kjo do ta thjeshtojë zhvillimin e mëtejshëm dhe do të pasuroni bibliotekën tuaj të operatorëve në projekt.
Më pas, të gjitha këto instanca të detyrave duhet të ekzekutohen, dhe tani do të flasim për planifikuesin.
Planifikuesi
Planifikuesi i detyrave në Airflow është ndërtuar mbi . Celery është një bibliotekë Python që mundëson organizimin e një rreshti plus ekzekutimin asinkron dhe të shpërndarë të detyrave. Nga ana e Airflow, të gjitha detyrat ndahen në grupe. Grupet krijohen manualisht. Përgjithësisht, qëllimi i tyre është të kufizojnë ngarkesën në punën me burimin ose të klasifikojnë detyrat brenda DWH. Grupeve mund t'u jepet menaxhimi përmes ndërfaqes web:

Ădo grup ka njĂ« kufizim tĂ« numrit tĂ« vendeve. Kur krijohet njĂ« DAG, i caktohet njĂ« grup:
ALERT_MAILS = Variable.get("gv_mail_admin_dwh")
DAG_NAME = 'dma_load'
OWNER = 'Vasya Pupkin'
DEPENDS_ON_PAST = True
EMAIL_ON_FAILURE = True
EMAIL_ON_RETRY = True
RETRIES = int(Variable.get('gv_dag_retries'))
POOL = 'dma_pool'
PRIORITY_WEIGHT = 10
start_dt = datetime.today() - timedelta(1)
start_dt = datetime(start_dt.year, start_dt.month, start_dt.day)
default_args = {
'owner': OWNER,
'depends_on_past': DEPENDS_ON_PAST,
'start_date': start_dt,
'email': ALERT_MAILS,
'email_on_failure': EMAIL_ON_FAILURE,
'email_on_retry': EMAIL_ON_RETRY,
'retries': RETRIES,
'pool': POOL,
'priority_weight': PRIORITY_WEIGHT
}
dag = DAG(DAG_NAME, default_args=default_args)
dag.doc_md = __doc__Puli i caktuar në nivelin e DAG-it mund të tejkalohet në nivelin e detyrës.
PĂ«r planifikimin e tĂ« gjitha detyrave nĂ« Airflow pĂ«rgjigjet njĂ« proces i veçantĂ« â Planifikuesi. Praktikisht, Planifikuesi merret me tĂ« gjithĂ« mekanikĂ«n e vendosjes sĂ« detyrave pĂ«r t'u ekzekutuar. NjĂ« detyrĂ«, para se tĂ« arrijĂ« pĂ«r ekzekutim, kalon disa etapa:
- Në DAG janë kryer detyrat e mëparshme, një e re mund të vendoset në radhë.
- Radhët klasifikohen në varësi të përparësisë së detyrave (përparësitë gjithashtu mund të menaxhohen), dhe, nëse ka një vend të lirë në pul, detyra mund të merret për punë.
- Nëse ka një punëtor të lirë celery, detyra i dërgohet atij; fillon puna që ju keni programuar në detyrë, duke përdorur operatorin përkatës.
Mjafton e thjeshtë.
Scheduler funksionon në shumë DAG-e dhe të gjitha detyrat brenda DAG-eve.
Për të filluar punën me një DAG, duhet t'i caktohet një orar:
dag = DAG(DAG_NAME, default_args=default_args, schedule_interval='@hourly')Ka njĂ« grup presetâesh tĂ« gatshme: @once, @hourly, @daily, @weekly, @monthly, @yearly.
Po ashtu mund të përdoren shprehje cron:
dag = DAG(DAG_NAME, default_args=default_args, schedule_interval='*/10 * * * *')Data e Ekzekutimit
PĂ«r tĂ« kuptuar se si funksionon Airflow, Ă«shtĂ« e rĂ«ndĂ«sishme tĂ« kuptoni se çfarĂ« Ă«shtĂ« Data e Ekzekutimit pĂ«r DAG-un. NĂ« Airflow, njĂ« DAG ka dimensionin e DatĂ«s sĂ« Ekzekutimit, dmth. sipas orarit tĂ« punĂ«s sĂ« DAG-ut, krijohen instanca tĂ« detyrave pĂ«r çdo DatĂ« Ekzekutimi. Dhe pĂ«r çdo DatĂ« Ekzekutimi, detyrat mund tĂ« ekzekutohen pĂ«rsĂ«ri â ose, pĂ«r shembull, DAG-u mund tĂ« punojĂ« nĂ« tĂ« njĂ«jtĂ«n kohĂ« nĂ« disa Data Ekzekutimi. Kjo ilustrohet kĂ«tu:

Fatkeqësisht (ndoshta edhe për mirë: varet nga situata), nëse realizohet një ndryshim në implementimin e problemit në DAG, ekzekutimi në Datat e Ekzekutimit të mëparshëm do të marrë parasysh ato rregullime. Kjo është mirë nëse nevojitet për të ripërllogaritur të dhënat në periudha të kaluara me algoritmin e ri, por është keq sepse humbet riprodhueshmëria e rezultatit (sigurisht, askush nuk e pengon të kthesh versionin e duhur nga Git dhe një herë ta llogaritësh atë që nevojitet, ashtu siç nevojitet).
Generimi i problemeve
Implementimi i DAG-ës është kod në Python, prandaj ne kemi një mënyrë shumë të përshtatshme për të reduktuar volumin e kodit kur punojmë, për shembull, me burime të ndara. Le të ketë tre sharda MySQL si burim, ju duhen të shkoni në secilin dhe të merrni disa të dhëna. Dhe kjo ndodhi në mënyrë të pavarur dhe paralele. Kodi në Python në DAG mund të duket kështu:
connection_list = lv.get('connection_list')
export_profiles_sql = '''
SELECT
id,
user_id,
nickname,
gender,
{{params.shard_id}} as shard_id
FROM profiles
'''
for conn_id in connection_list:
export_profiles = SqlToHiveViaHdfsTransfer(
task_id='export_profiles_from_' + conn_id,
sql=export_profiles_sql,
hive_table='stg.profiles',
overwrite=False,
tmpdir='/data/tmp',
conn_id=conn_id,
params={'shard_id': conn_id[-1:], },
compress=None,
dag=dag
)
export_profiles.set_upstream(exec_truncate_stg)
export_profiles.set_downstream(load_profiles)DAG rezulton si vijon:

Në këtë mënyrë, mund të shtoni apo hiqni një shard, thjesht duke rregulluar konfigurimin dhe duke përditësuar DAG. Shumë e përshtatshme!
Mund të përdorni edhe një gjenerim më të komplikuar të kodit, për shembull, të punoni me burimet si bazat e të dhënave ose të përshkruani strukturën e tabelave, algoritmin për punën me tabelën dhe duke marrë parasysh veçoritë e infrastrukturës DWH, të gjeneroni procesin e ngarkimit të N tabelave në magazinën tuaj. Ose, për shembull, të punoni me një API, e cila nuk mbështet funksionimin me parametrin si një listë, mund të gjeneroni N detyrat në DAG nga ky listë, të kufizoni paralelitetin e kërkesave në API me një pool dhe të nxirrni të dhënat e nevojshme nga API. Mjaft fleksibël!
Repozitori
Në Airflow ka një depo backend, DB (mund të jetë MySQL ose Postgres, ne kemi Postgres), ku ruhen gjendjet e detyrave, DAG-ve, konfigurimet e lidhjeve, variablat globale, etj., etj. Këtu do të doja të theksoja se depoja në Airflow është shumë e thjeshtë (rreth 20 tabela) dhe e lehtë për t'u ndërtuar nëse dëshironi të krijoni ndonjë proces mbi të. Më kujtohen 100500 tabela në depozitat Informatica që duhej të kuptoheshin për një kohë të gjatë para se të dija si të ndërtoja një kërkesë.
Monitorimi
Duke marrë parasysh thjeshtësinë e depozitës, mund të krijoni vetë një proces të përshtatshëm për monitorimin e detyrave. Ne përdorim një blloknot në Zeppelin, ku shohim gjendjen e detyrave:

Kjo mund të jetë gjithashtu ndërfaqja web e vetë Airflow:

Kodi i Airflow Ă«shtĂ« i hapur, kĂ«shtu qĂ« ne shtuam alarmin nĂ« Telegram. Ădo instancĂ« e punĂ«s sĂ« detyrĂ«s, nĂ«se ndodh njĂ« gabim, dĂ«rgon njĂ« mesazh nĂ« grupin nĂ« Telegram, ku Ă«shtĂ« e gjithĂ« ekipi i zhvillimit dhe mbĂ«shtetjes.
Marrim pĂ«rgjigje tĂ« shpejta pĂ«rmes Telegramit (nĂ«se kĂ«rkohet), pĂ«rmes Zeppelin â pamjen e pĂ«rgjithshme mbi detyrat nĂ« Airflow.
Në përfundim
Airflow Ă«shtĂ« nĂ« radhĂ« tĂ« parĂ« open source, dhe nuk duhet tĂ« prisni mrekulli prej tij. Jini tĂ« gatshĂ«m tĂ« investoni kohĂ« dhe energji pĂ«r tĂ« ndĂ«rtuar njĂ« zgjidhje funksionuese. QĂ«llimi Ă«shtĂ« i arritshĂ«m, besoni, ia vlen. ShpejtĂ«sia e zhvillimit, fleksibiliteti, thjeshtĂ«sia e shtimit tĂ« proceseve tĂ« reja â do t'ju pĂ«lqejĂ«. Sigurisht, Ă«shtĂ« e nevojshme tĂ« kushtoni shumĂ« vĂ«mendje organizimit tĂ« projektit, stabilitetit tĂ« funksionimit tĂ« Airflow-it: mrekullitĂ« nuk ndodhin.
Tani, Airflow pĂ«rpunon çdo ditĂ« rreth 6,5 mijĂ« detyrash. Sipas natyrĂ«s, ato janĂ« mjaft tĂ« ndryshme. Ka detyra qĂ« lidhen me ngarkimin e tĂ« dhĂ«nave nĂ« DWH kryesor nga shumĂ« burime tĂ« ndryshme dhe shumĂ« specifike, ka detyra pĂ«r llogaritjen e vitrinave brenda DWH kryesor, ka detyra pĂ«r publikimin e tĂ« dhĂ«nave nĂ« DWH tĂ« shpejtĂ«, ka shumĂ« e shumĂ« detyra tĂ« ndryshme â dhe Airflow i pĂ«rpunon ato çdo ditĂ«. NĂ«se flasim pĂ«r numra, atĂ«herĂ« kjo Ă«shtĂ« 2,3 mijĂ« detyra ELT tĂ« ndryshme nĂ« DWH (Hadoop), rreth 250 burime tĂ« tĂ« dhĂ«nave , kjo Ă«shtĂ« njĂ« ekip prej 4 zhvilluesish ETL, tĂ« cilĂ«t ndahen nĂ« pĂ«rpunimin ETL tĂ« tĂ« dhĂ«nave nĂ« DWH dhe pĂ«rpunimin ELT tĂ« tĂ« dhĂ«nave brenda DWH, dhe sigurisht edhe njĂ« admin, i cili merret me infrastrukturĂ«n e shĂ«rbimit.
Planet për të ardhmen
Numri i proceseve padyshim po rritet, dhe ajo që do të bëjmë në lidhje me infrastrukturën e Airflow do të jetë zgjerimi. Ne duam të ndërtuam një klaster Airflow, të bënim disa këmbë për worker'ët Celery dhe të krijonim një krye që dyfishohet me procese të planifikimit të punëve dhe një depo.
Epilogu
Sigurisht, kjo nuk Ă«shtĂ« e gjitha qĂ« do doja tĂ« flisja pĂ«r Airflow, por pĂ«rpavaçëm pĂ«rpjekja ime ishte tĂ« theksoja pikat kryesore. Apetiti vjen gjatĂ« ngrĂ«nies, provoni â dhe do t'ju pĂ«lqejĂ« đ
Burimi: habr.com
