
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 отвечает отдельный процесс — Scheduler. Собственно, Scheduler занимается всей механикой постановки задачек на исполнение. Задача, прежде чем попасть на исполнение, проходит несколько этапов:
- 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
