
Tere, Habr! Selles artiklis tahan rÀÀkida ĂŒhte suurepĂ€rasest tööriistast andmete töötlemise batch-protsesside arendamiseks, nĂ€iteks ettevĂ”tte DWH vĂ”i teie DataLake infrastruktuuris. Jutt kĂ€ib Apache Airflow'ist (edaspidi Airflow). See on ebaausalt tĂ€helepanuta jĂ€etud Habril ja peamises osas pĂŒĂŒan teid veenda, et Airflow'i tasub vaatama hakata, kui valite ETL/ELT-protsesside planeerijat.
Varem kirjutasin seeriat artikleid DWH teemal, kui töötasin Tinkoff Pangas. NĂŒĂŒd olen osa Mail.Ru Grupi meeskonnast ja tegelen andmeanalĂŒĂŒsi platvormi arendamisega mĂ€ngusuunal. Nagu alati, anname meeskonnaga siinkohal teada uudistest ja huvitavatest lahendustest meie andmeanalĂŒĂŒtika platvormil.
Proloog
Nii et alustame. Mis on Airflow? See on teek (vÔi ) tööprotsesside arendamiseks, planeerimiseks ja jÀlgimiseks. Airflow'i peamine omadus on see, et protsesside kirjeldamiseks (arendamiseks) kasutatakse Python'i keelt. Sealt tulenevad mitmed eelised teie projekti korraldamiseks ja arendamiseks: teie (nÀiteks) ETL-projekt on lihtsalt Python'i projekt ja saate seda korraldada nii, nagu soovite, arvestades infrastruktuuri eripÀra, meeskonna suurust ja muid nÔudeid. Tööriistade poolelt on kÔik lihtne. Kasutage nÀiteks PyCharm'i + Git'i. See on suurepÀrane ja vÀga mugav!
NĂŒĂŒd vaatame Airflow'i peamisi komponente. Nende olemuse ja otstarbe mĂ”istmine aitab teil protsesside arhitektuuri optimaalselt korraldada. TĂ”enĂ€oliselt on peamine element Directed Acyclic Graph (edaspidi DAG).
DAG
DAG on mĂ”ningane tĂ€henduslik kogum teie ĂŒlesannetest, mida soovite teatud jĂ€rjekorras kindla ajakava jĂ€rgi tĂ€ita. Airflow pakub mugavat veebiliidest DAG'ide ja muude komponentide jaoks:

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

Arendaja, kes projekteerib DAG'i, mÀÀrab komplekti operaatoritest, millele DAG'i sisesed ĂŒlesanded pĂ”hinevad. Siit jĂ”uame veel ĂŒhe olulise elemendini: Airflow'i operaator.
Operaatorid
Operaator on element, mille alusel luuakse ĂŒlesande eksemplare, kus kirjeldatakse, mis toimub ĂŒlesande eksemplari tĂ€itmise ajal. sisaldavad juba komplekti kasutamiseks valmis operaatoritest. NĂ€idised:
- BashOperator â kĂ€sk, et tĂ€ita bash-kĂ€sku.
- PythonOperator â kĂ€sk, et kutsuda esile Pythonâi kood.
- EmailOperator â kĂ€sk, et saata e-kirja.
- HTTPOperator â kĂ€sk, et töötada http-pĂ€ringutega.
- SqlOperator â kĂ€sk, et tĂ€ita SQL-koodi.
- Sensor â kĂ€sk, mis ootab sĂŒndmust (Ă”ige aja saabumist, vajaliku faili tekkimist, rida andmebaasis, vastust API-st jne.).
On olemas ka rohkem spetsiifilisi kÀske: DockerOperator, HiveOperator, S3FileTransferOperator, PrestoToMysqlOperator, SlackOperator.
Saate samuti arendada kĂ€skusid, pidades silmas oma omadusi, ja kasutada neid projektis. NĂ€iteks oleme loonud MongoDBToHiveViaHdfsTransfer, kĂ€sku dokumentide eksportimiseks MongoDB-st Hive'i, ning mitmeid kĂ€skusid tööks CHLoadFromHiveOperator ja CHTableLoaderOperator. Ăldiselt, kui projektis tekib tihti kasutatav kood, mis on ĂŒles ehitatud pĂ”hikĂ€sude peale, tasub mĂ”elda selle kogumisse uueks kĂ€skuks. See lihtsustab edasist arendust ja rikastab teie projektis kĂ€sude raamatukogu. Edasi tuleb kĂ”iki neid ĂŒlesanne eksemplaare tĂ€ita, ning nĂŒĂŒd rÀÀgime planeerijast.
Airflow'i ĂŒlesannete planeerija pĂ”hineb
Planeerija
Celery Iga basseini jaoks on mÀÀratud maksimaalne slotide arv. 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 = 10start_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__
Bassein, mis on mÀÀratud DAGâi tasemel, saab ĂŒle kirjutada ĂŒlesande tasemel.Kogu ĂŒlesannete planeerimise eest Airflow'is vastutab eraldi protsess â Planeerija. Tegelikult tegeleb Planeerija kogu ĂŒlesannete tĂ€itmise mehhanikaga. Enne kui ĂŒlesanne tĂ€itmiseks jĂ”uab, lĂ€bib see mitu etappi:
Airflow'i kĂ”igi ĂŒlesannete planeerimise eest vastutab eraldi protsess â Scheduler. Tegelikult tegeleb Scheduler kogu ĂŒlesannete tĂ€itmise mehhaanikaga. Enne, kui ĂŒlesanne tĂ€itmisele lĂ€heb, lĂ€bib see mitu etappi:
- DAG'is on tĂ€idetud eelnevad ĂŒlesanded, uue saab jĂ€rjekorda seadma.
- JĂ€rjekord sorteeritakse ĂŒlesannete prioriteedi alusel (prioriteediga saab samuti töötada) ja kui grupis on vaba koht, saab ĂŒlesande tööle vĂ”tta.
- Kui on vaba celery worker, suunatakse ĂŒlesanne sinna; algab töö, mille olete ĂŒlesandes kodeerinud, kasutades ĂŒhte vĂ”i teist operaatorit.
Piisavalt lihtne.
SĂŒsteemi planeerija töötab kĂ”igi DAG'ide ja kĂ”igi ĂŒlesannete ĂŒle DAG'ides.
Kuna sĂŒsteemi planeerija vĂ”iks DAG'iga töötama hakata, peab DAG'ile mÀÀrama ajakava:
dag = DAG(DAG_NAME, default_args=default_args, schedule_interval='@hourly')On olemas valmis preset'id: @once, @hourly, teeb peamise töö, aga mitte praegu. Praegu viskame lihtsalt meie konteksti logisse., @weekly, @monthly, @yearly.
Samuti vÔib kasutada cron-vÀljendeid:
dag = DAG(DAG_NAME, default_args=default_args, schedule_interval='*/10 * * * *')TÀideviimise kuupÀev
Et aru saada, kuidas Airflow töötab, on oluline mĂ”ista, mis on DAG'i jaoks tĂ€ideviimise kuupĂ€ev. Airflow's on DAG'il tĂ€ideviimise kuupĂ€ev, st vastavalt DAG'i töö ajakavale luuakse iga tĂ€ideviimise kuupĂ€eva jaoks ĂŒlesannete instantsid. Iga tĂ€ideviimise kuupĂ€eva korral saab ĂŒlesandeid uuesti tĂ€ita â vĂ”i nĂ€iteks DAG vĂ”ib töötada samaaegselt mitmel tĂ€ideviimise kuupĂ€eval. See on siin selgelt kujutatud:

Kahjuks (vĂ”i vĂ”ib-olla ka Ă”nneks: sĂ”ltub olukorrast), kui ĂŒlesande rakendust DAG'is muudetakse, siis tĂ€itmine eelnevates tĂ€ideviimise kuupĂ€evades toimub juba parandustega arvestades. See on hea, kui on vaja andmeid minevikus uue algoritmiga ĂŒmber arvutada, kuid halb, sest tulemuse reproduktsiooni vĂ”imekus kaob (loomulikult ei takista miski teil vajaliku lĂ€htekoodi versiooni Git'ist tagasi tuua ja ĂŒhekordselt arvutada vajalikku, nagu vaja).
Ălesannete genereerimine
DAG'i rakendamine on Python'i kood, seega on meil vĂ€ga mugav vĂ”imalus koodi mahtu vĂ€hendada, nĂ€iteks sharditud allikate töötlemisel. Oletame, et teil on kolm MySQL shardi allikana, peate kĂŒlastama igaĂŒht ja vĂ”tma mĂ”ned andmed. Ja seda iseseisvalt ja paralleelselt. Python'i kood DAG'is vĂ”ib vĂ€lja nĂ€ha nii:
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 abil saab ƥardi lisada vÔi eemaldada, lihtsalt seadistust kohandades ja DAG-i uuendades. Mugav!
VĂ”ib kasutada ka keerukamat koodigeneratsiooni, nĂ€iteks töötada andmebaasi tĂŒĂŒpi allikatega vĂ”i kirjeldada tabelistruktuuri, töötab tabeliga ja arvestades DWH infrastruktuuri eripĂ€rasid, genereerida N tabeli laadimisprotsessi teie andmesalvestisse. VĂ”i nĂ€iteks töö API-ga, mis ei toeta nimekirja parameetriga töötamist, saate selle nimekirja pĂ”hjal genereerida DAG-is N ĂŒlesannet, piirata API pĂ€ringute paralleelsust puhkuse kaudu ja saada API-st vajalikud andmed. Paindlik!
Repo
Airflow'l on oma tagatuba, andmebaas (see vĂ”ib olla MySQL vĂ”i Postgres, meil on Postgres), kuhu salvestatakse ĂŒlesannete, DAG-ide, ĂŒhenduste seadistuste, globaalsete muutujate jms seisundid. Sellega seoses tuleks öelda, et Airflow'i tagatuba on vĂ€ga lihtne (uin umbes 20 tabelit) ja mugav, kui soovite sellel pĂ”hjal oma protsessi luua. Tuleb meelde 100500 tabelit Informatica tagatubades, mille mĂ”istmiseks tuli palju vaeva nĂ€ha.
JĂ€lgimine
Kuna tagatuba on lihtne, saate ise luua endale mugava ĂŒlesannete jĂ€lgimise protsessi. Me kasutame Zeppelinis mĂ€rkmikku, kus vaatame ĂŒlesannete seisundit:

See vÔib olla ka Airflow'i enda veebiliides:

Airflow'i kood on avatud, seega oleme lisanud Telegraafi hĂ€ireteade. Iga töötav ĂŒlesande instants, kui tekib viga, saadab rĂ€mpspostina Telegraafis rĂŒhma, kus on kogu arendus- ja tugimeeskond.
Saame Telegraafi kaudu kiire reageerimise (kui see on vajalik), Zeppelinist â ĂŒldise ĂŒlevaate Airflow'i ĂŒlesannetest.
KokkuvÔttes
Airflow on esmajoones avatud lĂ€htekoodiga ja te ei tohi oodata imesid. Olge valmis kulutama aega ja vaeva, et vĂ€lja töötada toimiv lahendus. EesmĂ€rk on saavutatav, uskuge, see on seda vÀÀrt. Arenduse kiirus, paindlikkus, uute protsesside lisamise lihtsus â teile meeldib. Loomulikult tuleb projektiorganiseerimisele, Airflow'i enda töö stabiilsusele palju tĂ€helepanu pöörata: imesid ei juhtu.
Praegu töötab meie Airflow iga pĂ€ev umbes 6500 ĂŒlesannet. Iseloomult on need piisavalt erinevad. On andmete laadimise ĂŒlesandeid peamise DWH-sse paljusid erinevaid ja vĂ€ga spetsiifilisi allikaid, on andmete vitriinide arvutamise ĂŒlesandeid peamise DWH-siseselt, on ĂŒlesandeid andmete avaldamiseks kiire DWH-s, on palju erinevaid ĂŒlesandeid â ja Airflow töötleb neid kĂ”iki pĂ€evast pĂ€eva. Kui rÀÀkida numbritest, siis see 2,3 tuhat ELT ĂŒlesandeid erineva keerukusega DWH-s (Hadoop), umbes 2,5 sajandit andmebaasi allikaid, see on meeskond 4 ETL arendajat, kes jagunevad ETL andmete töötlemise ja ELT andmete töötlemise vahel DWH-s ning loomulikult on veel ĂŒks administraator, kes tegeleb teenuse infrastruktuuriga.
Tulevikuplaanid
Protsesside arv kasvab paratamatult ja peamine, millega me tegelema hakkame Airflow infrastruktuuri osas, on skaleerimine. Me tahame ĂŒles ehitada Airflow klusteri, eraldada paar jalgadele Celery töötajatele ja teha ennast dubleeriv pea ĂŒlesannete ajakava ja hoidla protsesside jaoks.
Epilog
See ei ole muidugi kĂ”ik, mida ma Airflow kohta rÀÀkida sooviksin, kuid peamised punktid olen pĂŒĂŒdnud vĂ€lja tuua. Isu tuleb söögi kĂ”rvale, proovige â ja see meeldib teile đ
Allikas: habr.com
