
Përshendetje, Habr! Në këtë artikull dëshiroj të flas për një mjet të shkëlqyer për zhvillimin e proceseve batch të përpunimit të të dhënave, si në infrastrukturën e DWH të korporatave ashtu edhe në DataLake tuaj. Bëhet fjalë për Apache Airflow (më poshtë Airflow). Ai është padrejtësisht i injoruar në Habr, dhe në pjesën kryesore do të përpiqem t'ju bind për atë që të paktën duhet ta merrni parasysh Airflow si një planifikues për proceset tuaja ETL/ELT.
Më parë kam shkruar një seri artikujsh mbi temën DWH, kur punoja në Tinkoff Bank. Tani jam bërë pjesë e ekipit të Mail.Ru Group dhe merrem me zhvillimin e platformës për analizën e të dhënave në drejtimin e lojërave. Në thelb, me kalimin e kohës dhe zhvillimin e lajmeve dhe zgjidhjeve interesante, ekipi ynë do të flasë këtu për platformën tonë të analizës së 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 njĂ« shumĂ«llojshmĂ«ri pĂ«rfitimesh pĂ«r organizimin e projektit tuaj dhe zhvillimin: nĂ« thelb, projekti juaj ETLâĂ«shtĂ« thjesht njĂ« projekt Python, dhe ju mund ta organizoni siç dĂ«shironi, duke marrĂ« parasysh veçoritĂ« e infrastrukturĂ«s, madhĂ«sinĂ« e ekipit dhe kĂ«rkesa tĂ« tjera. Mjetet janĂ« shumĂ« tĂ« thjeshta. 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. Pasi të kuptoni thelbin dhe qëllimin e tyre, do ta organizoni në mënyrë të optimizuar arkitekturën e proceseve. Mbase entiteti kryesor është Directed Acyclic Graph (më poshtë DAG).
DAG
DAG Ă«shtĂ« njĂ« bashkim kuptimor i detyrave tuaja, tĂ« cilat dĂ«shironi t'i ekzekutoni nĂ« njĂ« rend tĂ« caktuar sipas njĂ« orari tĂ« caktuar. Airflow ofron njĂ« ndĂ«rfaqe web tĂ« pĂ«rshtatshme pĂ«r punĂ«n me DAGâĂ«t dhe entitete tĂ« tjera:

DAG mund të duket kështu:

Zhvilluesi, duke projektuar DAG, vendos njĂ« grup operatorĂ«sh, mbi tĂ« cilat do tĂ« ndĂ«rtohen detyrat brenda DAGâit. KĂ«tu arrijmĂ« edhe 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 instancia të detyrave, që përshkruan se çfarë do të ndodhë gjatë ekzekutimit të instancës së detyrës. tashmë përmbajnë një grup operatorësh, të gatshëm për t'u përdorur. Shembuj:
- BashOperator â operator pĂ«r ekzekutimin e komandave bash.
- PythonOperator â operator pĂ«r thirrjen e kodit Python.
- EmailOperator â operator pĂ«r dĂ«rgimin e email-eve.
- HTTPOperator â operator pĂ«r punĂ«n me kĂ«rkesat http.
- SqlOperator â operator pĂ«r ekzekutimin e kodit SQL.
- Sensor â operator qĂ« pret njĂ« ngjarje (siç Ă«shtĂ« arritja e kohĂ«s sĂ« nevojshme, shfaqja e njĂ« skedari tĂ« kĂ«rkuar, njĂ« string nĂ« bazĂ«n e tĂ« dhĂ«nave, njĂ« pĂ«rgjigje nga API - etj.).
Ka operatorë më specifikë: DockerOperator, HiveOperator, S3FileTransferOperator, PrestoToMysqlOperator, SlackOperator.
Ju gjithashtu mund të zhvilloni operatorë, duke u bazuar në veçoritë 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, dhe disa operatorë për të punuar me : CHLoadFromHiveOperator dhe CHTableLoaderOperator. Në thelb, sa herë që në projekt lind një kod që përdoret shpesh, i ndërtuar mbi operatorët bazë, mund të mendoni të grumbulloni atë në një operator të ri. Kjo do ta thjeshtojë zhvillimin e mëtejshëm, dhe do t'i shtoni bibliotekës suaj të operatorëve në projekt.
Më pas, të gjithë këta instanca të tareas duhet të ekzekutohen, dhe tani është fjala për planifikuesin.
Dy ndryshime të dukshme në planifikim (të dyja në versionin alfa):
Planifikuesi i Đ·Đ°ĐŽĐ°Ń nĂ« Airflow Ă«shtĂ« i ndĂ«rtuar mbi . Celery Ă«shtĂ« njĂ« bibliotekĂ« Python qĂ« lejon organizimin e njĂ« rruge plus ekzekutimin asinkron dhe tĂ« shpĂ«rndarĂ« tĂ« detyrave. Nga ana e Airflow, tĂ« gjitha detyrat ndahen nĂ« pula. Pulat krijohen manualisht. Zakonisht, qĂ«llimi i tyre Ă«shtĂ« tĂ« kufizojnĂ« ngarkesĂ«n e punĂ«s me burimin ose tĂ« tipizojnĂ« detyrat brenda DWH. Pulat mund tĂ« menaxhohen pĂ«rmes ndĂ«rfaqes nĂ« web:

Ădo pul ka njĂ« kufizim nĂ« numrin e vendeve. Kur krijohet DAG, i caktohet njĂ« pul:
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-së mund të anulohet në nivelin e detyrës.
PĂ«r planifikimin e tĂ« gjitha detyrave nĂ« Airflow Ă«shtĂ« pĂ«rgjegjĂ«s njĂ« proces i veçantĂ« â Planifikuesi. NĂ« thelb, Planifikuesi merret me tĂ« gjithĂ« mekanikĂ«n e vendosjes sĂ« detyrave pĂ«r ekzekutim. NjĂ« detyrĂ«, pĂ«rpara se tĂ« shkojĂ« nĂ« ekzekutim, kalon disa etapa:
- Në DAG janë kryer detyrat e mëparshme, një e re mund të vendoset në radhë.
- Radhitja organizohet në bazë të prioritetit të detyrave (prioritetet gjithashtu mund të menaxhohen), dhe nëse ka një slot të lirë në pool, detyra mund të merret për punë.
- Nëse ka një worker celery të lirë, detyra dërgohet atje; fillon puna që keni programuar në detyrë, duke përdorur një operator të ndryshëm.
Mjafton.
Scheduler punon në shumicën e të gjitha DAG-eve dhe të gjitha detyrave brenda DAG-eve.
Për të filluar punën me DAG-un, DAG-u duhet të ketë një orar:
dag = DAG(DAG_NAME, default_args=default_args, schedule_interval='@hourly')Ka një grup presetësh të gatshëm: @once, @hourly, @daily, @weekly, @monthly, @yearly.
Gjithashtu mund të përdoren shprehjet 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 çfarë është Data e Ekzekutimit për DAG-un. Në Airflow, DAG ka një dimension të Datës së Ekzekutimit, pra, 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 mund të punojë njëkohësisht në disa Data Ekzekutimi. Kjo është e ilustruar qartë këtu:

Për fat të keq (ndoshta edhe për fat të mirë: varet nga situata), nëse ndryshohet implementimi i detyrës në DAG, atëherë ekzekutimi në Data Ekzekutimi të mëparshme do të bëhet tashmë duke pasur parasysh korrigjimet. Kjo është mirë nëse është e nevojshme të ripërllogaritetin të dhënat në periudhat e kaluara me një algoritëm të ri, por është keq, sepse humbet riprodhueshmëria e rezultatit (sigurisht, askush nuk e ndalon të rikthejë versionin e nevojshëm nga Git dhe të llogarisë një herë atë që nevojitet, ashtu si duhet).
Generimi i detyrave
Implementimi i DAG-ut është kod në Python, prandaj kemi një mënyrë shumë të përshtatshme për të reduktuar volumin e kodit kur punojmë, për shembull, me burime të shard-uara. Le të themi se keni si burim tri sharda MySQL, duhet të shkoni në çdo një dhe të merrni disa të dhëna. Dhe kjo të bëhet pa varësi dhe paralelisht. 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 i tillë:

Në të njëjtën kohë, mund të shtoni ose hiqni një shard, duke korrigjuar thjesht konfigurimin dhe duke përditësuar DAG-un. Shumë e përshtatshme!
Mund të përdorni gjithashtu një gjenerim më të komplikuar të kodit, për shembull, të punoni me burime si Baza të Dhënash ose të përshkruani strukturën e tabelave, algoritmin e funksionimit me tabelat dhe, duke marrë parasysh veçoritë e infrastrukturës DWH, të gjeneroni procesin e ngarkesës së N tabelave në ruajtjen tuaj. Ose, për shembull, të punoni me një API që nuk mbështet përdorimin e parametrave në formë liste, mund të gjeneroni N detyra në DAG nga kjo listë, të kufizoni paralelizmin e kërkesave në API me një pool dhe të nxirrni të dhënat e nevojshme nga API. Shumë fleksibël!
Repo
Në Airflow ka repo backend të vet, Baza e Dhënash (mund të jetë MySQL ose Postgres, ne kemi Postgres), ku ruhen gjendjet e detyrave, DAG-eve, konfigurimet e lidhjes, variablat globale dhe të tjera. Këtu do të doja të them se repozitori në Airflow është shumë i thjeshtë (rreth 20 tabela) dhe i përshtatshëm, nëse dëshironi të ndërtoni ndonjë proces tuajin mbi të. Më kujtohet 100500 tabela në repo Informatica, të cilat duhej të merren më kohë për t'u kuptuar se si të ndërtohej një kërkesë.
Monitorimi
Duke marrë parasysh thjeshtësinë e repo-it, mund të ndërtoni vetë një proces të përshtatshëm për monitorimin e detyrave. Ne përdorim një shënim në Zeppelin, ku shohim gjendjen e detyrave:

Kjo mund të jetë gjithashtu një ndërfaqe web e vetë Airflow-it:

Kodi i Airflow Ă«shtĂ« i hapur, prandaj ne kemi shtuar alarmin nĂ« Telegram. Ădo instancĂ« e detyrĂ«s qĂ« punon, nĂ«se ndodh njĂ« gabim, spamon nĂ« grupin nĂ« Telegram, ku ndodhet e gjithĂ« ekipi i zhvillimit dhe mbĂ«shtetjes.
Marrim reagime tĂ« shpejta nĂ«pĂ«rmjet Telegram-it (nĂ«se Ă«shtĂ« e nevojshme), pĂ«rmes Zeppelin-it â njĂ« pamje tĂ« pĂ«rgjithshme mbi detyrat nĂ« Airflow.
Përveç kësaj
Airflow nĂ« radhĂ« tĂ« parĂ« Ă«shtĂ« open source, dhe nuk duhet tĂ« prisni mrekulli nga ai. BĂ«huni tĂ« gatshĂ«m tĂ« kaloni kohĂ« dhe pĂ«rpjekje pĂ«r tĂ« ndĂ«rtuar njĂ« zgjidhje funksionale. 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, duhet tĂ« kushtoni shumĂ« vĂ«mendje organizatĂ«s sĂ« projektit, stabilitetit tĂ« funksionimit tĂ« vetĂ« Airflow-it: nuk ka mrekulli.
Tani kemi Airflow qĂ« punon çdo ditĂ« rreth 6,5 mijĂ« detyra. Karakteri i tyre Ă«shtĂ« mjaft i ndryshĂ«m. Ka detyra pĂ«r ngarkimin e tĂ« dhĂ«nave nĂ« DWH-in kryesor nga shumĂ« burime tĂ« ndryshme dhe shumĂ« specifike, ka detyra pĂ«r llogaritjen e tregtarĂ«ve brenda DWH-it kryesor, ka detyra pĂ«r publikimin e tĂ« dhĂ«nave nĂ« njĂ« DWH tĂ« shpejtĂ«, ka shumĂ« shumĂ« detyra tĂ« ndryshme â dhe Airflow i pĂ«rpunon ato çdo ditĂ«. NĂ«se flasim me shifra, atĂ«herĂ« kjo Ă«shtĂ« 2.3 mijĂ« detyra ELT tĂ« ndryshme nĂ« DWH (Hadoop), rreth 250 bazash tĂ« dhĂ«nash burimesh, kjo Ă«shtĂ« njĂ« ekip prej 4 ETL zhvilluesish, tĂ« cilĂ«t ndahen nĂ« procesimin ETL tĂ« tĂ« dhĂ«nave nĂ« DWH dhe nĂ« procesimin ELT tĂ« tĂ« dhĂ«nave brenda DWH dhe natyrisht edhe njĂ« administrator, i cili merret me infrastrukturĂ«n e shĂ«rbimit.
Planet për të ardhmen
Numri i proceseve po rritet dhe e vetmja gjë me të cilën do merren në aspektin e infrastrukturës Airflow, është shkallëzimi. Ne duam të ndërtrojmë një klaster Airflow, të ndajmë disa resurse për punëtorët Celery dhe të krijojmë një kokë të dyfishtë me proceset e planifikimit të detyrave dhe depozitën.
Epilogu
Kjo, sigurisht, nuk Ă«shtĂ« gjithçka pĂ«r tĂ« folur pĂ«r Airflow, por pĂ«rpjekjet kryesore i kam paraqitur. Atdheu vjen gjatĂ« ngrĂ«nies, provoje â dhe do tĂ« tĂ« pĂ«lqejĂ« đ
Burimi: habr.com
