
Tere, Habr! Selles artiklis soovin rÀÀkida ĂŒhest suurepĂ€rasest tööriistast andmete töötlemise batch-protsesside arendamiseks, nĂ€iteks ettevĂ”tte DWH vĂ”i teie DataLake infrastruktuuris. RÀÀgime Apache Airflow'st (edaspidi Airflow). Ta on Habr'is ebaĂ”iglaselt tĂ€helepanuta jĂ€etud ja peamises osas pĂŒĂŒan ma veenda teid, et vĂ€hemalt Airflow'le tasub tĂ€helepanu pöörata, kui valite ETL/ELT-protsesside ajakava.
Varem olen kirjutanud seeria artikleid DWH teemal, kui töötasin Tinkoff Pangas. NĂŒĂŒd olen osa Mail.Ru Group'ist ja tegelen andmete analĂŒĂŒsiplatvormi arendamisega mĂ€nguvaldkonnas. Tegelikult, uudiste ja huvitavate lahenduste ilmumise korral rÀÀgime me oma meeskonnaga siin meie andmete analĂŒĂŒtika platvormist.
Proloog
Nii, alustame. Mis on Airflow? See on teek (vÔi ) töövoogude arendamiseks, planeerimiseks ja jÀlgimiseks. Airflow peamine omadus: protsesside kirjeldamiseks (arendamiseks) kasutatakse Python'i keelt. Sealt tulenevad palju eeliseid oma projekti ja arenduse korraldamiseks: sisuliselt on teie (nÀiteks) ETL-projekt lihtsalt Python-projekt, mida saate korraldada vastavalt oma vajadustele, arvestades infrastruktuuri eripÀra, meeskonna suurust ja muid nÔudeid. Tööriistad on lihtsad. Kasutage nÀiteks PyCharm'i + Git'i. See on suurepÀrane ja vÀga mugav!
NĂŒĂŒd vaatame Airflow'i pĂ”hielemente. MĂ”istes nende olemust ja eesmĂ€rki, korraldate protsesside arhitektuuri optimaalselt. Peamine element on suunatud tsĂŒkliline graaf (edaspidi DAG).
DAG
DAG on teie ĂŒlesannete mĂ”tteline kogum, mida soovite teatavas jĂ€rjekorras teatud ajakava jĂ€rgi tĂ€ita. Airflow pakub mugavat veebiliidest DAG'ide ja teiste elementide haldamiseks:

DAG vÔib vÀlja nÀha selline:

Arendaja, kavandades DAG-i, sisestab operaatorite komplekti, mille pĂ”hjal ĂŒlesanded DAG-is ĂŒles ehitatakse. Siit jĂ”uame veel ĂŒhe olulise kontseptsioonini: Airflow Operaator.
Operaatorid
Operaator on ĂŒksus, mille alusel luuakse tĂ¶Ă¶ĂŒlesannete eksemplarid, kus kirjeldatakse, mis toimub tĂ¶Ă¶ĂŒlesande eksemplari tĂ€itmise ajal. kannavad juba endas komplekti valmis kasutamiseks mĂ”eldud operaatoritest. NĂ€ited:
- BashOperator â operaator bash-kĂ€skude tĂ€itmiseks.
- PythonOperator â operaator Python-koodi kutsumiseks.
- EmailOperator â operaator e-kirjade saatmiseks.
- HTTPOperator â operaator HTTP-pĂ€ringute tegemiseks.
- SqlOperator â operaator SQL-koodi tĂ€itmiseks.
- Sensor â operaator, mis ootab sĂŒndmust (nagu vajaliku aja saabumine, nĂ”utava faili ilmumine, rida andmebaasis, API-st saadud vastus jne).
On ka spetsiifilisemaid operaatoreid: DockerOperator, HiveOperator, S3FileTransferOperator, PrestoToMysqlOperator, SlackOperator.
Samuti saate arendada operaatorite, lĂ€htudes oma eripĂ€radest, ja kasutada neid projektis. NĂ€iteks lĂ”ime MongoDBToHiveViaHdfsTransfer, operaatori dokumentide eksportimiseks MongoDB-st Hive'i, ja mitu operaatorit tööks koos. : CHLoadFromHiveOperator ja CHTableLoaderOperator. Tegelikult, kui projektis tekib sageli kasutatav kood, mis on ĂŒles ehitatud pĂ”hioperatoritele, siis tasub kaaluda selle kokkupanekut uueks operaatoriks. See lihtsustab edasist arendust ning tĂ€idate oma projekti operaatorite raamatukogu.
SeejĂ€rel tuleb need ĂŒlesande eksemplarid tĂ€ita ning nĂŒĂŒd rÀÀgime ajastamisest.
Ajakava
Airflow ĂŒlesannete ajastaja pĂ”hineb . Celery on Python'i teek, mis vĂ”imaldab korraldada jĂ€rjekorvi ning asĂŒnkroonset ja jaotatud ĂŒlesannete tĂ€itmist. Airflow'i poolelt jagunevad kĂ”ik ĂŒlesanded basseinideks. Basseinid luuakse kĂ€sitsi. Reeglina on nende eesmĂ€rgiks piirata koormust allikaga töötamisel vĂ”i tĂŒĂŒpeerida ĂŒlesandeid DWH sees. Basseine saab hallata veebiliidese kaudu:

Igal basseinil on piirang salveslotide arvu osas. DAG'i loomisel mÀÀratakse sellele bassein:
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__DAG tasemel mÀÀratud bassein saab ĂŒlesande tasemel ĂŒle kirjutada.
Airflowis vastutab kĂ”igi ĂŒlesannete ajastamise eest eraldi protsess â Scheduler. Scheduler tegeleb ĂŒlesannete tĂ€itmise mehaanika kĂ”igega, mis on seotud ĂŒlesannete tegemisega. Ălesanne peab enne tĂ€itmist lĂ€bima mitu etappi:
- DAG-is on eelnevad ĂŒlesanded lĂ”petatud, uue saab jĂ€rjekorda panna.
- JĂ€rjekord sorteeritakse vastavalt ĂŒlesannete prioriteetidele (prioriteetidega saab samuti hallata) ning kui basseinil on vaba koht, saab ĂŒlesande tööle vĂ”tta.
- Kui vaba celery worker on olemas, suunatakse ĂŒlesanne sellesse; algab töö, mille te olete ĂŒlesandes mÀÀratlenud, kasutades seda vĂ”i seda operaatorit.
Lihtne.
Scheduler töötab mitmete DAG-ide ja DAG-ides asuvate ĂŒlesannete peal.
Kuna Scheduler peab hakkama töötama DAG-iga, peab DAG-il olema ajakava:
dag = DAG(DAG_NAME, default_args=default_args, schedule_interval='@hourly')On olemas hulk valmispreset'e: @once, @hourly, teeb pĂ”hiosa tööst, aga mitte nĂŒĂŒd. Hetkel lihtsalt vĂ€ljastame meie konteksti logisse., @weekly, @monthly, @yearly.
Samuti on vÔimalik kasutada cron-vÀljendeid:
dag = DAG(DAG_NAME, default_args=default_args, schedule_interval='*/10 * * * *')TÀideviimise kuupÀev
Airflow'i toimimise mĂ”istmiseks on oluline mĂ”ista, mis on DAG-i jaoks TĂ€ideviimise kuupĂ€ev. Airflow's on DAG-il TĂ€ideviimise kuupĂ€eva mÔÔde, st sĂ”ltuvalt DAG-i töö ajakavast luuakse iga TĂ€ideviimise kuupĂ€eva jaoks ĂŒlesande eksemplarid. Ja iga TĂ€ideviimise kuupĂ€eva jaoks saab ĂŒlesandeid uuesti tĂ€ita â vĂ”i nĂ€iteks vĂ”ib DAG töötada samaaegselt mitmel TĂ€ideviimise kuupĂ€eval. See on selgelt nĂ€idatud siin:

Kahjuks (vĂ”i vĂ”ib-olla ka Ă”nneks, sĂ”ltub olukorrast) on, kui DAG'i ĂŒlesande elluviimist muudetakse, siis eelnevate tĂ€itmise kuupĂ€evade tulemused arvestavad juba parandusi. See on hea, kui on vaja andmeid varasemate perioodide jooksul uuesti arvutada uue algoritmiga, kuid halb, kuna tulemuse korduvus kaob (loomulikult ei takista miski tagasi minna Git'i ja saada vajalik lĂ€htekood ning ĂŒhekordselt arvutada, kuidas vaja).
Ălesannete genereerimine
DAG'i elluviimine â kood Pythonis, seega on meil vĂ€ga mugav vĂ”imalus koodi mahtu vĂ€hendada, töötades nĂ€iteks sharditud allikatega. Oletame, et teil on kolm MySQL shard'i allikana, peate minema igasse ja vĂ”tma mingid andmed. Samas sĂ”ltumatult ja paralleelselt. Pythonis DAG'is vĂ”ib kood vĂ€lja nĂ€ha selline:
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 on selline:

Selle kÀigus saab shardi lisada vÔi eemaldada, lihtsalt seadistust kohandades ja DAG-i vÀrskendades. Mugav!
Saate kasutada ka keerukamat koodigeneratsiooni, nĂ€iteks töötades andmebaasi allikatest vĂ”i kirjeldades tabelistruktuuri, tabelite töötlemise algoritmi ning laenguprotsessi genereerimist, arvestades DWH infrastruktuuri eripĂ€rasid. VĂ”i nĂ€iteks, kui API ei toeta loendina parametri töötlemist, saate selle loendi pĂ”hjal DAG-is genereerida N ĂŒlesannet, piirata API-pĂ€ringute paralleelsust ja koguda vajalikud andmed API-st. Paindlik!
Repo
Airflow'il on oma tagavarala, andmebaas (see vĂ”ib olla MySQL vĂ”i Postgres, meil on Postgres), kus hoitakse ĂŒlesannete, DAG'ide, ĂŒhenduste seadete, globaalsete muutujate jne olekuid. Siinkohal tahaksin öelda, et Airflow'i tagavara on vĂ€ga lihtne (umbes 20 tabelit) ja mugav, kui soovite selle ĂŒle oma protsessi luua. Tuleb meelde 100500 tabelit Informatica tagavaras, mis nĂ”udsid pikaajalist tutvumist, enne kui sain aru, kuidas pĂ€ringut koostada.
JĂ€lgimine
Arvestades tagavara lihtsust, saate ise luua endale sobiva ĂŒlesannete jĂ€lgimise protsessi. Me kasutame Zeppelin'is mĂ€rkmikku, kus vaatame ĂŒlesannete olekut:

See vÔib olla ka Airflow'i veebiliides:

Airflow'i kood on avatud, seega lisasime me teavituse Telegram'i. Iga ĂŒlesande aktiivne instants, kui toimub viga, spammib Telegram'i gruppi, kus on kogu arendus- ja tugimeeskond.
Saame Telegram'i kaudu kiire reageerimise (kui see on vajalik), Zeppelin'i kaudu â ĂŒldise ĂŒlevaate Airflow'i ĂŒlesannetest.
Kokku
Airflow on avatud lĂ€htekoodiga ja ei tasu oodata, et see imepĂ€raselt töötaks. Olge valmis investeerima aega ja vaeva, et luua toimiv lahendus. EesmĂ€rk on saavutatav ja uskuge, see on seda vÀÀrt. Arenduse kiirus, paindlikkus, uute protsesside lihtne lisamine â teile meeldib see. Loomulikult tuleb palju tĂ€helepanu pöörata projekti korraldamisele ja Airflow enda tööstabiilsusele: imesid ei juhtu.
Praegu töötleb meie Airflow igapĂ€evaselt umbes 6,5 tuhat ĂŒlesannet. Need ĂŒlesanded on iseloomult ĂŒsna erinevad. On ĂŒlesandeid andmete laadimiseks peamisse DWH-st paljusid erinevaid ja vĂ€ga spetsiifilisi allikaid, on ĂŒlesandeid vitriinide arvutamiseks peamise DWH sees, on ĂŒlesandeid andmete avaldamiseks kiiremas DWH-s, ja palju, palju erinevaid ĂŒlesandeid â ja Airflow töötleb kĂ”iki neid pĂ€evast pĂ€eva. Kui rÀÀkida numbritest, siis see on 2,3 tuhat ELT erineva keerukusega ĂŒlesannet DWH-s (Hadoop), umbes 250 erinevat andmebaasi allikat, see on meeskond neljast ETL arendajast, kes jagunevad ETL andmete töötlemise ja ELT andmete töötlemise vahel DWH-s ja muidugi veel ĂŒks administraator, kes tegeleb teenuse infrastruktuuriga.
Tulevikuplaanid
Protsesside arv kasvab ĂŒha enam ning meie peamine tegevusala Airflow infrastruktuuris on skaleerimine. Tahame luua Airflow klastrit, eraldada mĂ”ned jalad Celery töötajatele ja teha dubleeriva pea, mis haldab ĂŒlesannete planeerimise protsesse ja hoidlat.
Epilog
See ei ole muidugi kĂ”ik, mida Airflow kohta rÀÀkida tahaks, kuid olen pĂŒĂŒdnud peamised punktid esile tuua. Isu tuleb söögi ajal, proovi â ja sulle meeldib đ
Allikas: habr.com
