PĂ«rshĂ«ndetje, unĂ« jam Dmitriy Logvinenko â Inxhinier i tĂ« DhĂ«nave nĂ« departamentin e analizĂ«s sĂ« grupit tĂ« kompanive «VezĂ«t».
Do t'ju tregoj pĂ«r njĂ« mjet tĂ« shkĂ«lqyer pĂ«r zhvillimin e proceseve ETL â Apache Airflow. Por Airflow Ă«shtĂ« aq universal dhe shumĂ«fishtĂ«, saqĂ« ia vlen ta shqyrtoni edhe nĂ«se nuk mereni me rrjedhat e tĂ« dhĂ«nave, por keni nevojĂ« tĂ« filloni nganjĂ«herĂ« ndonjĂ« proces dhe tĂ« ndihmoni nĂ« mbajtjen e tyre nĂ«n kontroll.
Po, nuk do të flas vetëm, por do të tregoj dhe: programi ka shumë kode, shkëmbime dhe rekomandime.

ĂfarĂ« zakonisht shihni kur kĂ«rkoni fjalĂ«n Airflow / Wikimedia Commons
Përmbajtja
Hyrje
Apache Airflow â Ă«shtĂ« ashtu si Django:
- shkruar në Python,
- ka një admin të shkëlqyer,
- është pa kufij për zgjerim,
â vetĂ«m mĂ« tĂ« mirat, dhe Ă«shtĂ« krijuar pĂ«r qĂ«llime krejtĂ«sisht tĂ« tjera, pĂ«rkatĂ«sisht (siç Ă«shtĂ« shkruar deri nĂ« kaptinĂ«):
- nisjen dhe monitorimin e detyrave në numër të pakufizuar makinash (sa do t'ju lejojë Celery/Kubernetes dhe ndërgjegjja juaj)
- me gjenerimin dinamik të flukseve të punës nga një kod shumë i lehtë për t'u shkruar dhe për t'u kuptuar në Python
- dhe me mundësinë për të lidhur çdo bazë të të dhënave dhe API me ndihmën e komponenteve të gatshme ose pluginëve të vetë-krijuar (gjë që bëhet jashtëzakonisht e lehtë).
Ne përdorim Apache Airflow në këtë mënyrë:
- mbledhim të dhëna nga burime të ndryshme (shumë instanca SQL Server dhe PostgreSQL, API të ndryshme me metrika të aplikacioneve, madje edhe 1C) në DWH dhe ODS (këtë e bëjmë me Vertica dhe Clickhouse).
- si një avanguardë
cron, i cili ekzekuton procese konsolidimi të të dhënave në ODS, si dhe monitoron shërbimin e tyre.
Derisa përpara pak kohësh ne e përmbushëm nevojën tonë me një server të vogël me 32 bërthama dhe 50 GB RAM. Në Airflow funksionon:
- më 200 DAG-ë (në thelb flukse pune, në të cilat kemi mbushur detyrat),
- në secilin në mesatarisht 70 detyra,
- kjo gjë ekzekutohet (po ashtu në mesatarisht) një herë në orë.
Dhe pĂ«r mĂ«nyrĂ«n se si ne u zgjeruam, do tĂ« shkruaj mĂ« poshtĂ«, por tani le tĂ« pĂ«rcaktojmĂ« ĂŒber-detyrĂ«n qĂ« do tĂ« zgjidhim:
Ka janĂ« tri servera SQL burimorĂ«, tĂ« cilĂ«t kanĂ« nga 50 baza tĂ« dhĂ«nash â instanca tĂ« njĂ« projekti, pĂ«rkatĂ«sisht struktura e tyre Ă«shtĂ« e njĂ«jtĂ« (nĂ« shumicĂ«n e rasteve, mwa-ha-ha), dhe pĂ«r pasojĂ«, çdo njĂ«ri ka njĂ« tabelĂ« Orders (fatmirĂ«sisht, njĂ« tabelĂ« me njĂ« emĂ«r tĂ« tillĂ« mund tĂ« futet nĂ« çdo biznes). Ne marrim tĂ« dhĂ«nat duke shtuar fushat ndihmĂ«se (serveri burim, baza burimore, identifikuesi i detyrĂ«s ETL) dhe naivisht do t'i hedhim ato nĂ«, tĂ« themi, Vertica.
Fillojmë!
Pjesa kryesore, praktike (edhe pak teorike)
Pse na duhet (dhe juve)
Kur pemët ishin të mëdha dhe unë isha një i thjeshtë SQL-çkuar në një shitje të madhe ruse, ne hidheshim në proceset ETL aka rrjedhat e të dhënave me dy mjete të disponueshme për ne:
- Informatica Power Center â njĂ« sistem ekstremisht kompleks, jashtĂ«zakonisht produktiv, me pajisjet e veta dhe versionimin e saj. Kam pĂ«rdorur ndoshta 1% tĂ« mundĂ«sive tĂ« saj. Pse? Epo, sĂ« pari, ky ndĂ«rfaqe ishte diku nga vitet '00 dhe na e shkaktonte presion psikologjik. SĂ« dyti, kjo gjĂ« Ă«shtĂ« e mbyllur nĂ« procese tejet komplekse, riblerje intensive tĂ« komponenteve dhe shumĂ« karakteristika tĂ« rĂ«ndĂ«sishme pĂ«r ndĂ«rmarrjet. PĂ«r atĂ« çmim qĂ« kushton, si fluturimi i njĂ« Airbus A380/njĂ« vit, do tĂ« heshtim.
Kujdes, screenshot mund t'i bëjë të rinjve nën 30 vjeç të ndihen disi keq

- SQL Server Integration Server â ne kemi pĂ«rdorur kĂ«tĂ« shok nĂ« proceset tona interne. Por nĂ« realitet: ne tashmĂ« pĂ«rdorim SQL Server, dhe do tĂ« ishte disi e paarsyeshme tĂ« mos i shfrytĂ«zonim mjetet ETL tĂ« tij. Gjithçka nĂ« tĂ« Ă«shtĂ« mirĂ«: dhe ndĂ«rfaqja Ă«shtĂ« e bukur, si dhe raportet e ekzekutimit... Por nuk pĂ«r kĂ«tĂ« e duam produktin software, oh jo pĂ«r kĂ«tĂ«. TĂ« versiononi atĂ«
dtsx(i cili është një XML me node të përzier gjatë ruajtjes) mund ta bëjmë, por çfarë do të thotë? Të bëjmë një paketë detyrash që do të transferonte njëqind tabela nga një server në tjetrin? Po çfarë njëqind, do t'ju bie gishti tregues pas njëzet copash, duke klikuar mbi butonin e mausit. Por duke parë, sigurisht duket më modern:
Kemi kërkuar pa dyshim zgjidhje. Punë madje gjysmë arritëm të krijojmë një gjenerator të vetë-ndërtuar për paketat SSIS...
... dhe pastaj më gjeti puna e re. Dhe aty më preku Apache Airflow.
Kur mĂ«sova se pĂ«rshkrimet e proceseve ETL janĂ« thjesht kod Python, pothuajse kĂ«rcen nga gĂ«zimi. KĂ«shtu qĂ« rrjedhat e tĂ« dhĂ«nave u nĂ«nshtruan versionimit dhe difes, dhe grumbullimi i tabelave me njĂ« strukturĂ« unike nga njĂ«qind baza tĂ« dhĂ«nash nĂ« njĂ« target u bĂ« njĂ« punĂ« e kodit Python nĂ« njĂ« ekran 13â me njĂ« tĂ« dhjetĂ«.
Po mblidhni një klaster
Le të mos organizojmë një kopësht të fëmijëve, dhe të mos flasim për gjëra të heltësisht të qarta, siç është instalimi i Airflow, baza e të dhënave që keni zgjedhur, Celery dhe punë të tjera të përshkruara në dokumente.
Që të mund të fillojmë menjëherë me eksperimentet, kam skicuar docker-compose.yml në të cilin:
- Le të ngrisim vetë Airflow: Scheduler, Webserver. Atje do të funksionojë gjithashtu Flower për monitorimin e detyrave të Celery (sepse është shtypur tashmë në
apache/airflow:1.10.10-python3.7, dhe ne nuk kemi asnjë problem me këtë); - PostgreSQL, në të cilin Airflow do të shkruajë informacionin e saj operativ (të dhënat e planifikuesit, statistikat e ekzekutimit etj.), ndërsa Celery do të shënojë detyrat e përfunduara;
- Redis, i cili do të veprojë si një broker detyrash për Celery;
- Celery worker, i cili do të merret me ekzekutimin e drejtpërdrejtë të detyrave.
- NĂ« dosjen
./dagsne do të vendosim skedarët tanë me përshkrimin e dags. Ato do të kapen në flakë, prandaj nuk është e nevojshme të tërhiqni të gjithë stekun pas çdo shkurti.
Disa herë kodi në shembuj është dhënë jo plotësisht (për të mos e mbingarkuar tekstin), dhe diku ai modifikohet gjatë procesit. Shembujt e plotë funksionues të kodit mund të shihen në repositorium. .
docker-compose.yml
version: '3.4'
x-airflow-config: &airflow-config
AIRFLOW__CORE__DAGS_FOLDER: /dags
AIRFLOW__CORE__EXECUTOR: CeleryExecutor
AIRFLOW__CORE__FERNET_KEY: MJNz36Q8222VOQhBOmBROFrmeSxNOgTCMaVp2_HOtE0=
AIRFLOW__CORE__HOSTNAME_CALLABLE: airflow.utils.net:get_host_ip_address
AIRFLOW__CORE__SQL_ALCHEMY_CONN: postgres+psycopg2://airflow:airflow@airflow-db:5432/airflow
AIRFLOW__CORE__PARALLELISM: 128
AIRFLOW__CORE__DAG_CONCURRENCY: 16
AIRFLOW__CORE__MAX_ACTIVE_RUNS_PER_DAG: 4
AIRFLOW__CORE__LOAD_EXAMPLES: 'False'
AIRFLOW__CORE__LOAD_DEFAULT_CONNECTIONS: 'False'
AIRFLOW__EMAIL__DEFAULT_EMAIL_ON_RETRY: 'False'
AIRFLOW__EMAIL__DEFAULT_EMAIL_ON_FAILURE: 'False'
AIRFLOW__CELERY__BROKER_URL: redis://broker:6379/0
AIRFLOW__CELERY__RESULT_BACKEND: db+postgresql://airflow:airflow@airflow-db/airflow
x-airflow-base: &airflow-base
image: apache/airflow:1.10.10-python3.7
entrypoint: /bin/bash
restart: always
volumes:
- ./dags:/dags
- ./requirements.txt:/requirements.txt
services:
# Redis si Celery broker
broker:
image: redis:6.0.5-alpine
# DB për metadata e Airflow
airflow-db:
image: postgres:10.13-alpine
environment:
- POSTGRES_USER=airflow
- POSTGRES_PASSWORD=airflow
- POSTGRES_DB=airflow
volumes:
- ./db:/var/lib/postgresql/data
# Kontejner kryesor me Webserver, Scheduler, Celery Flower të Airflow
airflow:
<<: *airflow-base
environment:
<<: *airflow-config
AIRFLOW__SCHEDULER__DAG_DIR_LIST_INTERVAL: 30
AIRFLOW__SCHEDULER__CATCHUP_BY_DEFAULT: 'False'
AIRFLOW__SCHEDULER__MAX_THREADS: 8
AIRFLOW__WEBSERVER__LOG_FETCH_TIMEOUT_SEC: 10
depends_on:
- airflow-db
- broker
command: >
-c " sleep 10 &&
pip install --user -r /requirements.txt &&
/entrypoint initdb &&
(/entrypoint webserver &) &&
(/entrypoint flower &) &&
/entrypoint scheduler"
ports:
# Celery Flower
- 5555:5555
# Webserver i Airflow
- 8080:8080
# Punëtori Celery, do të shkallëzohet duke përdorur `--scale=n`
worker:
<<: *airflow-base
environment:
<<: *airflow-config
command: >
-c " sleep 10 &&
pip install --user -r /requirements.txt &&
/entrypoint worker"
depends_on:
- airflow
- airflow-db
- brokerShënime:
- NĂ« ndĂ«rtimin e kompozitĂ«s kam mbĂ«shtetur shumĂ« nĂ« imazhin e njohur â patjetĂ«r shikoni. Ndoshta nĂ« jetĂ«n tuaj nuk keni nevojĂ« pĂ«r asgjĂ« tjetĂ«r.
- Të gjitha konfigurimet e Airflow janë të disponueshme jo vetëm përmes
airflow.cfg, por edhe përmes variablave të ambientit (falë zhvilluesve), të cilat i përdora keq. - Natyrisht, ai nuk është gati për prodhim: qëllimisht nuk kam vendosur heartbeats në kontenierë, nuk u shqetësova për sigurinë. Megjithatë, bëra një minimum të përshtatshëm për eksperimentet tona.
- VĂ« re se:
- Folderi me DAG-at duhet të jetë i aksesueshëm si për planifikuesin ashtu edhe për punëtorët.
- E njĂ«jta gjĂ« vlen edhe pĂ«r tĂ« gjitha bibliotekat e palĂ«ve tĂ« treta â ato duhet tĂ« jenĂ« tĂ« gjitha tĂ« instaluara nĂ« makinat me planifikuesin dhe punĂ«torĂ«t.
Tani thjesht:
$ docker-compose up --scale worker=3Pas ngritjes, mund të shikoni ndërfaqet e internetit:
- Airflow:
- Flower:
Koncepte themelore
Nëse nuk keni kuptuar asgjë nga këto «DAG», ja një fjalor i shkurtër:
- Planifikuesi â ŰŻŰłŰȘĆgĆ· (djali mĂ« i rĂ«ndĂ«sishĂ«m) nĂ« Airflow, qĂ« kontrollon qĂ« robotĂ«t tĂ« punojnĂ«, jo njerĂ«zit: monitoron orarin, pĂ«rditĂ«son DAG-at, nis detyrat.
Në versionet më të vjetra, ai kishte probleme me memorinë (jo, jo amnezia, por rrjedhje) dhe atje kishte mbetur një parametër legacy në konfigurime.
run_durationâ intervali i rivendosjes sĂ« tij. Por tani gjithçka Ă«shtĂ« nĂ« rregull. - DAG (edhe i njohur si «dag») â «grafik i orientuar jo-ciklik», por njĂ« pĂ«rcaktim i tillĂ« shumĂ« pak njerĂ«zve do t'u thotĂ« diçka, dhe nĂ« thelb, Ă«shtĂ« njĂ« kontejner pĂ«r detyrat qĂ« ndĂ«rveprojnĂ« me njĂ«ra-tjetrĂ«n (shih mĂ« poshtĂ«) ose njĂ« analog i PaketĂ«s nĂ« SSIS dhe Workflow nĂ« Informatica.
Përveç dag-ve, gjithashtu mund të ketë sub-dag, por ndoshta nuk do të arrijmë atje.
- DAG Run â njĂ« dag i inicializuar, tĂ« cilit i Ă«shtĂ« dhĂ«nĂ« e tij
execution_date. DAGRAT e njĂ« dag-u mund tĂ« punojnĂ« paralelisht (nĂ«se, sigurisht, keni bĂ«rĂ« detyrat tuaja idempotente). - Operator â kĂ«to janĂ« copa kode qĂ« janĂ« pĂ«rgjegjĂ«se pĂ«r ekzekutimin e njĂ« veprimi tĂ« caktuar. Ka tre lloje operatorĂ«sh:
- veprim, si për shembull operatori ynë i preferuar
PythonOperator, i cili mund të ekzekutojë çdo (të vlefshëm) kod Python; - transfer, të cilët transportojnë të dhëna nga një vend në një tjetër, për shembull,
MsSqlToHiveTransfer; - sensor do të lejojë të reagoni ose të ngadalësoni ekzekutimin e dag-ut deri në ndodhjen e një eventi të caktuar.
HttpSensormund të tërheqë endpoint-in e caktuar dhe kur të presë përgjigjen e duhur, të fillojë transferiminGoogleCloudStorageToS3Operator. Një mendje kurioze do të pyesë: "përse? Sepse mund të bëjmë ripërsëritje direkt në operator!" Dhe pastaj, që të mos bllokojmë rezervuarin e detyrave me operatorë që janë pezull. Sensorët aktivizohen, kontrollojnë dhe vdesin deri në përpjekjen e ardhshme.
- veprim, si për shembull operatori ynë i preferuar
- DetyrĂ« â operatorĂ«t e shpallur, pa marrĂ« parasysh llojin, dhe tĂ« lidhur me dagun rriten nĂ« nivelin e detyrave.
- Instanca e detyrĂ«s â kur gjenerali-planifikues vendos se detyrat janĂ« gati pĂ«r t'u dĂ«rguar nĂ« luftĂ« te punĂ«torĂ«t-ekzekutorĂ« (nĂ« vend, nĂ«se po pĂ«rdorim
LocalExecutorose në një nod të largët në rastin eCeleryExecutor), ai u cakton atyre një kontekst (dmth. një grup variablesh - parametrash ekzekutimi), zhvillon modelet e komandave ose kërkesave dhe i vendos ato në rezervuar.
Drejtojmë detyrat
Së pari do të shënojmë skemën e përgjithshme të dagut tonë dhe pastaj do të thellojmë gjithnjë e më shumë në detaje, sepse përdorim disa zgjidhje jo triviale.
Pra, në formën më të thjeshtë, një dag i tillë do të dukej kështu:
nga datetime import timedelta, datetime
nga airflow import DAG
nga airflow.operators.python_operator import PythonOperator
nga commons.datasources import sql_server_ds
dag = DAG('orders',
schedule_interval=timedelta(hours=6),
start_date=datetime(2020, 7, 8, 0))
def workflow(**context):
print(context)
për conn_id, schema në sql_server_ds:
PythonOperator(
task_id=schema,
python_callable=workflow,
provide_context=True,
dag=dag)Le të kuptojmë:
- Së pari importojmë libraritë e nevojshme dhe disa gjëra të tjera;
sql_server_dsâ kjo Ă«shtĂ«List[namedtuple[str, str]]me emrat e lidhjeve nga Airflow Connections dhe bazat e tĂ« dhĂ«nave nga tĂ« cilat do tĂ« marrim tabelĂ«n tonĂ«;dagâ deklarata e DAG-ut tonĂ«, e cila duhet tĂ« jetĂ« nĂ«globals(), pĂ«rndryshe Airflow nuk do ta gjente. DAG-ut i duhet gjithashtu tĂ« themi:- se quhet
ordersâ ky emĂ«r pastaj do tĂ« shfaqet nĂ« ndĂ«rfaqen e internetit, - se do tĂ« fillojĂ« punĂ«n nĂ« mesnatĂ« mĂ« 8 korrik,
- dhe duhet të ekzekutohet, afërsisht çdo 6 orë (për djemtë e mrekullueshëm këtu në vend të
timedelta()lejohetcron-linjĂ«0 0 0/6 ? * * *, pĂ«r ata mĂ« pak tĂ« mrekullueshĂ«m â njĂ« shprehje si@daily);
- se quhet
workflow()do të bëjë punën kryesore, por jo tani. Tani ne thjesht do të hedhim kontekstin tonë në log.- Dhe tani magjia e thjeshtë e krijimit të detyrave:
- ecim përmes burimeve tona;
- inicjalizojmë
PythonOperator, i cili do të realizojë shoshin tonëworkflow(). Mos u harroni të tregoni një emër unik (brenda DAG-ut) për detyrën dhe të lidheni me vetë DAG-un. Flagsiguroni_kontekstnga ana e tij do të hedhë në funksion argumente shtesë, të cilat ne do t'i mbledhim me kujdes përmes**konteksti.
Derisa në këtë pikë, çfarë kemi marrë:
- një DAG të ri në ndërfaqen web,
- njëqind e pesëdhjetë detyra që do të ekzekutohen paralelisht (nëse lejojnë konfigurimet e Airflow, Celery dhe fuqia e serverëve).
Të paktën pothuajse e morëm.

Kush do të instalohet varësitë?
Për ta thjeshtuar këtë, unë e kam ngjeshur në docker-compose.yml trajtim requirements.txt në të gjitha nodet.
Tani po fillon:

Kube gri â instanca tĂ« detyrave, tĂ« pĂ«rpunuara nga planifikuesi.
Disa minuta pritje, detyrat i kapin punëtorët:

TĂ« gjelbra, e qartĂ«, â qĂ« kanĂ« punuar me sukses. TĂ« kuqe â jo shumĂ« me sukses.
Meqenëse, në prodhim tonë nuk ka ndonjë dosje
./dags, e cila sinkronizohet midis makinave â tĂ« gjitha DAG-tĂ« ndodhen nĂ«gitnĂ« Gitlab tonĂ«, dhe Gitlab CI shpĂ«rndani pĂ«rditĂ«simet nĂ« makina kur behet bashkimi nĂ«master.
Pak për Flower
NdĂ«rsa punĂ«torĂ«t po punojnĂ« me detyrat tona boshe, kujtojmĂ« njĂ« mjet tjetĂ«r, i cili mund tĂ« na tregojĂ« diçka â Flower.
Faqja e parë me informacione përmbledhëse për nodet-punëtorë:

Faqja më e ngarkuar me detyra që kanë filluar punën:

Faqja më e mërzitshme me gjendjen e brokerit tonë:

Faqja mĂ« e ndritshme â me grafikat e gjendjes sĂ« detyrave dhe kohĂ«n e pĂ«rfundimit tĂ« tyre:

Po ngarkojmë të ngarkuarin
Pra, të gjithë detyrat kanë përfunduar, mund të çojmë të plagosurit.

Dhe numri i tĂ« plagosurve ishte i konsiderueshĂ«m â pĂ«r arsye tĂ« ndryshme. NĂ«se pĂ«rdoret siç duhet Airflow, kĂ«ta katrorĂ« tregojnĂ« se tĂ« dhĂ«nat me siguri nuk arritĂ«n.
Duhet të shohim logun dhe të rilançojmë instancat e detyrave që kanë dështuar.
Duke klikuar në çdo katror, do të shohim veprimet që na janë në dispozicion:

Mund të marrim dhe të bëjmë Clear për atë që ka dështuar. Pra, harrojmë se ka ndodhur diçka, dhe e njëjta instancë detyre do të shkojë te planifikuesi.

Natyrisht, nuk Ă«shtĂ« shumĂ« humane tĂ« bĂ«sh kĂ«shtu me tĂ« gjitha katrorĂ«t e kuq â nuk Ă«shtĂ« ky pritja jonĂ« nga Airflow. Sigurisht, kemi armĂ«n e masave shkatĂ«rruese: Browse/Task Instances

Do të zgjedhim gjithçka njëherësh dhe do të klikojmë opsionin e duhur për të kaluar në zero:

Pas pastrimit, taksitë tona duken kështu (ato tashmë po presin me padurim që planifikuesi t'i planifikojë):

Lidhje, likuj dhe variabla të tjerë
Koha më e mirë për të parë DAG-un tjetër, update_reports.py:
from collections import namedtuple
from datetime import datetime, timedelta
from textwrap import dedent
from airflow import DAG
from airflow.contrib.operators.vertica_operator import VerticaOperator
from airflow.operators.email_operator import EmailOperator
from airflow.utils.trigger_rule import TriggerRule
from commons.operators import TelegramBotSendMessage
dag = DAG('update_reports',
start_date=datetime(2020, 6, 7, 6),
schedule_interval=timedelta(days=1),
default_args={'retries': 3, 'retry_delay': timedelta(seconds=10)})
Report = namedtuple('Report', 'source target')
reports = [Report(f'{table}_view', table) for table in [
'reports.city_orders',
'reports.client_calls',
'reports.client_rates',
'reports.daily_orders',
'reports.order_duration']]
email = EmailOperator(
task_id='email_success', dag=dag,
to='{{ var.value.all_the_kings_men }}',
subject='Raportet DWH janë përditësuar',
html_content=dedent("""Të nderuar, raportet janë përditësuar"""),
trigger_rule=TriggerRule.ALL_SUCCESS)
tg = TelegramBotSendMessage(
task_id='telegram_fail', dag=dag,
tg_bot_conn_id='tg_main',
chat_id='{{ var.value.failures_chat }}',
message=dedent("""
Natasha, zgjohem, ne {{ dag.dag_id }} e hodhëm
"""),
trigger_rule=TriggerRule.ONE_FAILED)
for source, target in reports:
queries = [f"TRUNCATE TABLE {target}",
f"INSERT INTO {target} SELECT * FROM {source}"]
report_update = VerticaOperator(
task_id=target.replace('reports.', ''),
sql=queries, vertica_conn_id='dwh',
task_concurrency=1, dag=dag)
report_update >> [email, tg]Ădokush ka bĂ«rĂ« ndonjĂ«herĂ« njĂ« raport azhurnimi, apo jo? Ja pĂ«rsĂ«ri: ka njĂ« listĂ« burimesh nga ku tĂ« marrim tĂ« dhĂ«nat; ka njĂ« listĂ« ku t'i vendosim; mos harro tĂ« na lajmerosh kur ndodhi ndonjĂ« gjĂ« ose kur diçka dĂ«shtoi (mirĂ«, kjo nuk Ă«shtĂ« pĂ«r ne, jo).
Le të kalojmë përsëri në skedarin dhe të shikojmë gjërat e reja enigmatike:
from commons.operators import TelegramBotSendMessageâ asgjĂ« nuk na pengon tĂ« krijojmĂ« operatorĂ«t tanĂ«, siç e bĂ«mĂ« ne duke krijuar njĂ« mbĂ«shtjellĂ«s tĂ« vogĂ«l pĂ«r dĂ«rgimin e mesazheve nĂ« Razblokirovannyy. (PĂ«r kĂ«tĂ« operator do tĂ« flasim mĂ« poshtĂ«);default_args={}â dag-u mund tĂ« ndajĂ« argumentet e njĂ«jta me tĂ« gjithĂ« operatorĂ«t e tij;to='{{ var.value.all_the_kings_men }}'â fushatonuk do tĂ« jetĂ« e koduar, por do tĂ« formohet dinamikisht me anĂ« tĂ« Jinja dhe njĂ« variabli me listĂ«n e email-eve qĂ« e kam vendosur me kujdes nĂ«Admin/Variables;trigger_rule=TriggerRule.ALL_SUCCESSâ kushti pĂ«r aktivizimin e operatorit. NĂ« rastin tonĂ«, letra do tĂ« shkojĂ« te shefat vetĂ«m nĂ«se tĂ« gjitha varĂ«sitĂ« pĂ«rfunduan suksesshĂ«m;tg_bot_conn_id='tg_main'â argumentetconn_idpranojnĂ« identifikues tĂ« lidhjeve qĂ« krijojmĂ« nĂ«Admin/Connections;trigger_rule=TriggerRule.ONE_FAILEDâ mesazhet nĂ« Telegram do tĂ« dĂ«rgohen vetĂ«m nĂ« rast se ka detyra tĂ« dĂ«shtuar;task_concurrency=1â ndalojmĂ« ekzekutimin e shumĂ« instanceve tĂ« njĂ«jta tĂ« njĂ« detyre nĂ« tĂ« njĂ«jtĂ«n kohĂ«. NĂ« tĂ« kundĂ«rt, do tĂ« kemi njĂ« ekzekutim tĂ« shumĂ« instancaveVerticaOperator(tĂ« shikojnĂ« njĂ« tabelĂ«);report_update >> [email, tg]â tĂ« gjithaVerticaOperatordo tĂ« pĂ«rfundojnĂ« duke dĂ«rguar e-mail dhe mesazhe, kĂ«shtu:

Por, pasi operatorët e njoftimeve kanë kushte të ndryshme për t'u ekzekutuar, do të punojë vetëm njëri. Në Pamjen e Dendrës duket pak më pak e qartë:

Do tĂ« them disa fjalĂ« pĂ«r makrotĂ« dhe miqtĂ« e tyre â variablat.
Makrotë janë placeholder-e Jinja, të cilat mund të vendosin informacion të ndryshëm të dobishëm në argumentet e operatorëve. Për shembull, kështu:
SELECT
id,
payment_dtm,
payment_type,
client_id
FROM orders.payments
WHERE
payment_dtm::DATE = '{{ ds }}'::DATE{{ ds }} do të zhvillohet në përmbajtjen e variablit të kontekstit execution_date në formatin YYYY-MM-DD: 2020-07-14. E bukura është se variablat e kontekstit lidhen fort me një instancë të caktuar të detyrës (katrorin në Pamjen e Dendrës), dhe kur bëhet një rinisje, placeholder-et do të zbulohen në të njëjtat vlera.
Vlerat e ndara mund të shikohen me butonin Rendered në çdo instancë të detyrës. Kështu është për detyrën me dërgimin e e-mailit:

Dhe kështu është për detyrën me dërgimin e mesazhit:

Lista e plotë e makrosve të integruar për versionin më të fundit të disponueshëm është këtu:
Për më tepër, me ndihmën e plugin-eve, ne mund të shpallim makros tona të personalizuara, por kjo është një histori krejt tjetër.
Përveç gjërave të paracaktuar, ne mund të përdorim vlerat e variablave tanë (më sipër në kod unë e kam përdorur këtë). Le të krijojmë në Admin/Variables një çift gjërash:

Gati pĂ«r tâu pĂ«rdorur:
TelegramBotSendMessage(chat_id='{{ var.value.failures_chat }}')Në vlerë mund të jetë skalar, ose mund të jetë edhe JSON. Në rastin e JSON-it:
bot_config
{
"bot": {
"token": 881hskdfASDA16641,
"name": "Verter"
},
"service": "TG"
}thjesht përdorim rrugën për çelësin e kërkuar: {{ var.json.bot_config.bot.token }}.
Do të them vetëm një fjalë dhe do të tregoj një skrin që ka të bëjë me lidhet. Këtu gjithçka është elementare: në faqe Admin/Connections krijojmë lidhjen, vendosim atje emrat e përdoruesve/fjalëkalimet dhe parametra më specifik. Ja kështu:

FjalĂ«kalimat mund tĂ« shkruhen (mĂ« me kujdes se nĂ« variantin e paracaktuar), ose mund tĂ« mos caktohet lloji i lidhjes (ashtu siç bĂ«ra pĂ«r tg_main) â çështja Ă«shtĂ« se lista e llojeve Ă«shtĂ« e koduar nĂ« modelet Airflow dhe zgjerimi nuk lejon ndryshime nĂ« burim (nĂ«se ndonjĂ«herĂ« nuk kam gjetur diçka â ju lutem mĂ« korrigjoni), por nuk ka asgjĂ« qĂ« na ndalon tĂ« marrim kredencialet thjesht me emrin.
Gjithashtu, mund të bëni disa lidhje me një emër të vetëm: në këtë rast, metoda BaseHook.get_connection(), e cila na merr lidhjet sipas emrit, do të kthejë një rastësor nga disa të ngjashëm (do të ishte më logjike të bëhej Round Robin, por do ta lëmë këtë në ndërgjegjen e zhvilluesve të Airflow).
Variables dhe Connections, pa dyshim, janĂ« mjete tĂ« shkĂ«lqyera, por Ă«shtĂ« e rĂ«ndĂ«sishme tĂ« mos humbasĂ«sh balancimin: cilat pjesĂ« tĂ« rrjedhave tuaja i mbani nĂ« kodin tuaj dhe cilat â i jepni pĂ«r ruajtje nĂ« Airflow. Nga njĂ«ra anĂ«, ndĂ«rrimi i shpejtĂ« i vlerĂ«s, pĂ«r shembull, box-i i dĂ«rgesĂ«s, mund tĂ« jetĂ« i lehtĂ« pĂ«rmes UI. Nga ana tjetĂ«r, kjo Ă«shtĂ« nĂ« fund tĂ« fundit njĂ« kthim nĂ« klikimin me maus, nga i cili ne (unĂ«) do doja tĂ« shpenzoja.
Puna me lidhjet â Ă«shtĂ« njĂ« nga detyrat hook-esh. NĂ« pĂ«rgjithĂ«si, hook-at e Airflow janĂ« pika lidhĂ«se me shĂ«rbime dhe biblioteka tĂ« jashtme. P.sh., JiraHook do tĂ« hapĂ« pĂ«r ne njĂ« klient pĂ«r tĂ« bashkĂ«punuar me Jira (mund tĂ« lĂ«vizim detyra kĂ«tu-kĂ«tu), dhe me ndihmĂ«n e SambaHook mund tĂ« dĂ«rgoni njĂ« skedar lokal nĂ« smb-pikĂ«.
Analizojmë operatorin e personalizuar
Dhe ne u afruam ngushtësisht për të parë se si është bërë TelegramBotSendMessage
Kodi commons/operators.py me operatorin vetë:
from typing import Union
from airflow.operators import BaseOperator
from commons.hooks import TelegramBotHook, TelegramBot
class TelegramBotSendMessage(BaseOperator):
"""Dërgo mesazh në chat_id duke përdorur TelegramBotHook
Shembull:
>>> TelegramBotSendMessage(
... task_id='telegram_fail', dag=dag,
... tg_bot_conn_id='tg_bot_default',
... chat_id='{{ var.value.all_the_young_dudes_chat }}',
... message='{{ dag.dag_id }} dështoi :(',
... trigger_rule=TriggerRule.ONE_FAILED)
"""
template_fields = ['chat_id', 'message']
def __init__(self,
chat_id: Union[int, str],
message: str,
tg_bot_conn_id: str = 'tg_bot_default',
*args, **kwargs):
super().__init__(*args, **kwargs)
self._hook = TelegramBotHook(tg_bot_conn_id)
self.client: TelegramBot = self._hook.client
self.chat_id = chat_id
self.message = message
def execute(self, context):
print(f'Dërgo "{self.message}" në bisedën {self.chat_id}')
self.client.send_message(chat_id=self.chat_id,
message=self.message)Këtu, ashtu si gjithçka tjetër në Airflow, është shumë e thjeshtë:
- Trashëguam nga
BaseOperator, e cila implementon shumë gjëra specifike për Airflow (shikoni kur keni kohë) - Deklaruam fushat
template_fields, në të cilat Jinja do të kërkojë makro për përpunim. - Organizuam argumentet e duhura për
__init__(), vendosëm default-et atje ku duhet. - Nuk harrojmë as inicializimin e prindit.
- Hapëm hook-un përkatës
TelegramBotHook, morëm prej tij një objekt-klient. - E overrajduam (përkufizuam përsëri) metodën
BaseOperator.execute(), qĂ« Airflow do ta thĂ«rrasĂ« kur Ă«shtĂ« koha pĂ«r tĂ« filluar operatorin â aty do tĂ« realizojmĂ« veprimin kryesor, pa harruar tĂ« regjistrojmĂ«. (Regjistrohemi, pĂ«r shembull, direkt nĂ«stdoutdhestderrâ Airflow do ta kapĂ« gjithçka, do ta mbĂ«shtjellĂ« bukur, do ta vendosĂ« ku duhet.)
Të shikojmë çfarë kemi në commons/hooks.py. Pjesa e parë e skedarit, me vetë hook-un:
from typing import Union
from airflow.hooks.base_hook import BaseHook
from requests_toolbelt.sessions import BaseUrlSession
class TelegramBotHook(BaseHook):
"""Telegram Bot API hook
Shënim: shtoni një lidhje me Conn Type bosh dhe mos harroni
të plotësoni Extra:
{"bot_token": "YOuRAwEsomeBOtToKen"}
"""
def __init__(self,
tg_bot_conn_id='tg_bot_default'):
super().__init__(tg_bot_conn_id)
self.tg_bot_conn_id = tg_bot_conn_id
self.tg_bot_token = None
self.client = None
self.get_conn()
def get_conn(self):
extra = self.get_connection(self.tg_bot_conn_id).extra_dejson
self.tg_bot_token = extra['bot_token']
self.client = TelegramBot(self.tg_bot_token)
return self.clientS'di ashtë çfarë mund të shpjegohet këtu, thjesht do të theksoj disa pika të rëndësishme:
- Ne trashĂ«gojmĂ«, mendojmĂ« pĂ«r argumentet â nĂ« shumicĂ«n e rasteve do tĂ« jetĂ« vetĂ«m njĂ«:
conn_id; - Tejkalojmë metodat standarde: unë e kufizova
get_conn(), ku marrë parametrat e lidhjes sipas emrit dhe thjesht nxjerr seksioninextra(ky është një fushë për JSON), në të cilën unë (sipas udhëzimit tim!) vendosa tokenin e bosit Telegram:{"bot_token": "YOuRAwEsomeBOtToKen"}. - Krijoj një instancë të
TelegramBot, duke i dhënë gjithashtu tokenin e saktë.
Këtu është. Të marrë klientin nga huka mund të bëhet me TelegramBotHook().clent ose TelegramBotHook().get_conn().
Dhe pjesa e dytë e skedarit, ku bëj një mikro mbështjellje për Telegram REST API, për të mos sjellur të njëjtin për një metodë të vetme sendMessage.
class TelegramBot:
"""Telegram Bot API wrapper
Examples:
>>> TelegramBot('YOuRAwEsomeBOtToKen', '@myprettydebugchat').send_message('Hi, darling')
>>> TelegramBot('YOuRAwEsomeBOtToKen').send_message('Hi, darling', chat_id=-1762374628374)
"""
API_ENDPOINT = 'https://api.telegram.org/bot{}/'
def __init__(self, tg_bot_token: str, chat_id: Union[int, str] = None):
self._base_url = TelegramBot.API_ENDPOINT.format(tg_bot_token)
self.session = BaseUrlSession(self._base_url)
self.chat_id = chat_id
def send_message(self, message: str, chat_id: Union[int, str] = None):
method = 'sendMessage'
payload = {'chat_id': chat_id or self.chat_id,
'text': message,
'parse_mode': 'MarkdownV2'}
response = self.session.post(method, data=payload).json()
if not response.get('ok'):
raise TelegramBotException(response)
class TelegramBotException(Exception):
def __init__(self, *args, **kwargs):
super().__init__((args, kwargs))Qëllimi i duhur është të bashkosh të gjitha këto:
TelegramBotSendMessage,TelegramBotHook,TelegramBotâ nĂ« plugin, ta vendosĂ«sh nĂ« njĂ« depo tĂ« hapur dhe ta ofrosh si Open Source.
Përderisa po e shqyrtonim gjithçka, përditësimet tona të raporteve arritën që të dërgojnë një mesazh gabimi në kanalin tim. Po shkoj të kontrolloj se çfarë është sërish gabim...

Diçka është prishur në dagun tonë! Po ashtu kjo është ajo që prisnim? Pikërisht!
A do të derdhësh?
Ndiheni se kam harruar diçka? Duket se premtova të transferoj të dhënat nga SQL Server në Vertica dhe këtu u largova nga tema, sa bezdisës!
Krimi ishte i qëllimtë, thjesht kisha detyrim të shpjegoja disa terminologji për ju. Tani mund të vazhdojmë.
Plani ynë ishte si më poshtë:
- Të bëjmë dag
- Të gjenerojmë detyra
- Të shohim si duket gjithçka
- Të caktojmë numrat e seancave për mbushjet
- Të marrim të dhënat nga SQL Server
- Të vendosim të dhënat në Vertica
- Të grumbullojmë statistikat
Pra, për të gjithë këtë, kam bërë një shtesë të vogël në docker-compose.yml:
docker-compose.db.yml
version: '3.4'
x-mssql-base: &mssql-base
image: mcr.microsoft.com/mssql/server:2017-CU21-ubuntu-16.04
restart: always
environment:
ACCEPT_EULA: Y
MSSQL_PID: Express
SA_PASSWORD: SayThanksToSatiaAt2020
MSSQL_MEMORY_LIMIT_MB: 1024
services:
dwh:
image: jbfavre/vertica:9.2.0-7_ubuntu-16.04
mssql_0:
<<: *mssql-base
mssql_1:
<<: *mssql-base
mssql_2:
<<: *mssql-base
mssql_init:
image: mio101/py3-sql-db-client-base
command: python3 ./mssql_init.py
depends_on:
- mssql_0
- mssql_1
- mssql_2
environment:
SA_PASSWORD: SayThanksToSatiaAt2020
volumes:
- ./mssql_init.py:/mssql_init.py
- ./dags/commons/datasources.py:/commons/datasources.pyAtje shohim:
- Vertica si host
dwhme cilësimet më të zakonshme, - tri instanca SQL Server,
- plotsojmë bazat me disa të dhëna të fundit (nuk duhet të shikoni në
mssql_init.py!)
E kemi nisur gjithçka me një komandë pak më të komplikuar se më parë:
$ docker-compose -f docker-compose.yml -f docker-compose.db.yml up --scale worker=3ĂfarĂ« ka gjeneruar rastĂ«sia jonĂ«, mund tĂ« shikohet duke pĂ«rdorur pikĂ«n Profilizimi i tĂ« DhĂ«nave/KĂ«rkesa Ad Hoc:

Kryesore, mos e tregoni këtë analistëve
Të ndalemi në detaje seancat ETL nuk do ta bëj, aty është gjithçka e thjeshtë: krijojmë një bazë, në të një tabelë, e mbështjellim gjithçka me menaxherin e kontekstit, dhe tani veprojmë kështu:
with Session(task_name) as session:
print('Load', session.id, 'started')
# Load worflow
...
session.successful = True
session.loaded_rows = 15session.py
nga sys import stderr
klasa Session:
"""Seanca e punës ETL
Example:
me Session(task_name) si sesion:
print(session.id)
session.successful = True
session.loaded_rows = 15
session.comment = 'Mirë
"""
def __init__(self, connection, task_name):
self.connection = connection
self.connection.autocommit = True
self._task_name = task_name
self._id = None
self.loaded_rows = None
self.successful = None
self.comment = None
def __enter__(self):
return self.open()
def __exit__(self, exc_type, exc_val, exc_tb):
if any(exc_type, exc_val, exc_tb):
self.successful = False
self.comment = f'{exc_type}: {exc_val}n{exc_tb}'
print(exc_type, exc_val, exc_tb, file=stderr)
self.close()
def __repr__(self):
return (f'')
@property
def task_name(self):
return self._task_name
@property
def id(self):
return self._id
def _execute(self, query, *args):
with self.connection.cursor() as cursor:
cursor.execute(query, args)
return cursor.fetchone()[0]
def _create(self):
query = """
KRIJO TABELĂ NĂSE NUK EKZISTON sessions (
id SERIAL NUK NĂN NORMĂS
task_name VARCHAR(200) NUK NĂN NORMĂS,
started TIMESTAMPTZ NUK NĂN NORMĂS DEFAULT current_timestamp,
finished TIMESTAMPTZ DEFAULT current_timestamp,
successful BOOL,
loaded_rows INT,
comment VARCHAR(500)
);
"""
self._execute(query)
def open(self):
query = """
SHTO NĂ sessions (task_name, finished)
VLERAT (%s, NULL)
KTHIM id;
"""
self._id = self._execute(query, self.task_name)
print(self, 'hapur')
return self
def close(self):
if not self._id:
raise SessionClosedError('Seanca nuk është e hapur')
query = """
PĂRMIKSO sessions
SET
finished = DEFAULT,
successful = %s,
loaded_rows = %s,
comment = %s
KU
id = %s
KTHIM id;
"""
self._execute(query, self.successful, self.loaded_rows,
self.comment, self.id)
print(self, 'mbyllur',
', e suksesshme: ', self.successful,
', Ngarkuar: ', self.loaded_rows,
', komenti:', self.comment)
klasa SessionError(Exception):
kalon
klasa SessionClosedError(SessionError):
kalonKa ardhur koha të marrim të dhënat tona nga pesëmbëdhjetë tabela tona. Do ta bëjmë këtë me disa rreshta të thjeshtë:
source_conn = MsSqlHook(mssql_conn_id=src_conn_id, schema=src_schema).get_conn()
query = f"""
SELECT
id, start_time, end_time, type, data
FROM dbo.Orders
WHERE
CONVERT(DATE, start_time) = '{dt}'
"""
df = pd.read_sql_query(query, source_conn)- Me anë të hook-ut do marrim nga Airflow
pymssql-konnekt - NĂ« kĂ«rkesĂ« do vendosim njĂ« kufizim nĂ« formĂ«n e datĂ«s â kĂ«tĂ« do ta vendosĂ« shabloni.
- Tërheqim kërkesën tonë
pandas, e cila do tĂ« sjellĂ« pĂ«r neDataFrameâ do na nevojitet mĂ« vonĂ«.
Unë përdor zëvendësimin
{dt}në vend të parametrave të kërkesës%sjo pse unë jam një Buratino i keq, por sepsepandasnuk mund të përballojë mepymssqldhe i ofron këtij të funditparams: List, ndonëse ai shumë dëshirontuple.
Gjithashtu, vini re se zhvilluesipymssqlvendosi të mos e mbështesë më, dhe është koha të kalojmë nëpyodbc.
Le të shohim çfarë i shtoi Airflow argumenteve tona:

Nëse nuk ka të dhëna, nuk ka kuptim të vazhdojmë. Por gjithashtu është çuditshëm të quhet ngarkimi i suksesshëm. Po ashtu, nuk është një gabim. Ah, çfarë të bëjmë?! Ja çfarë:
if df.empty:
raise AirflowSkipException('Nuk ka rreshta për të ngarkuar')AirflowSkipException do të thotë Airflow se nuk ka gabime, por ne do ta kalojmë detyrën. Në ndërfaqe do të ketë një katror që nuk është gjelber dhe as i kuq, por me ngjyrë rozë.
Do t'i hedhim të dhënat tona për disa kolona:
df['etl_source'] = src_schema
df['etl_id'] = session.id
df['hash_id'] = hash_pandas_object(df[['etl_source', 'id']])Specifikisht:
- DB, nga e cila kemi marrë porositë,
- Identifikuesi i sesionit tonë të ngarkimit (do të jetë i ndryshëm për secilën detyrë),
- Hash nga burimi dhe identifikuesi i porosisĂ« â qĂ« nĂ« bazĂ«n pĂ«rfundimtare (ku gjithçka do tĂ« bashkohet nĂ« njĂ« tabelĂ«) tĂ« kemi njĂ« identifikues unik pĂ«r porosinĂ«.
Ka mbetur hapi i parafundit: të ngarkohet gjithçka në Vertica. Dhe, siç është e çuditshme, një nga mënyrat më efektive për ta bërë këtë është përmes CSV!
# Export data to CSV buffer
buffer = StringIO()
df.to_csv(buffer,
index=False, sep='|', na_rep='NUL', quoting=csv.QUOTE_MINIMAL,
header=False, float_format='%.8f', doublequote=False, escapechar='\')
buffer.seek(0)
# Push CSV
target_conn = VerticaHook(vertica_conn_id=target_conn_id).get_conn()
copy_stmt = f"""
COPY {target_table}({df.columns.to_list()})
FROM STDIN
DELIMITER '|'
ENCLOSED '"'
ABORT ON ERROR
NULL 'NUL'
"""
cursor = target_conn.cursor()
cursor.copy(copy_stmt, buffer)- Ne bëjmë një marrëveshje speciale
StringIO. pandasdo të përfshijë në tëDataFramenë formën erreshtat CSV.Do të hapim një lidhje me Vertica-n tonë të dashur.- Dhe tani me ndihmën e
- copy()
do t'i dërgojmë të dhënat tona direkt në Vertica!do t'i dërgojmë të dhënat tona direkt në Vertica!
Do të marrim nga driver-i sa rreshta janë ngarkuar, dhe do t'i themi menaxherit të sesionit që gjithçka është OK:
session.loaded_rows = cursor.rowcount
session.successful = TrueKëtu është gjithçka.
Në prodhim krijojmë tabelën e synuar manualisht. Këtu le të lejoj vetes një automatizëm të vogël:
create_schema_query = f'CREATE SCHEMA IF NOT EXISTS {target_schema};'
create_table_query = f"""
CREATE TABLE IF NOT EXISTS {target_schema}.{target_table} (
id INT,
start_time TIMESTAMP,
end_time TIMESTAMP,
type INT,
data VARCHAR(32),
etl_source VARCHAR(200),
etl_id INT,
hash_id INT PRIMARY KEY
);"""
create_table = VerticaOperator(
task_id='create_target',
sql=[create_schema_query,
create_table_query],
vertica_conn_id=target_conn_id,
task_concurrency=1,
dag=dag)Unë me ndihmën e
VerticaOperator()krijoj skemën e DB-së dhe tabelën (nëse ato nuk ekzistojnë, sigurisht). E rëndësishme është të vendosësh drejt varësitë:
për conn_id, schema në sql_server_ds:
load = PythonOperator(
task_id=schema,
python_callable=workflow,
op_kwargs={
'src_conn_id': conn_id,
'src_schema': schema,
'dt': '{{ ds }}',
'target_conn_id': target_conn_id,
'target_table': f'{target_schema}.{target_table}'},
dag=dag)
create_table >> loadPërmbledhje
â Ja, â tha minusi, â a nuk Ă«shtĂ« e vĂ«rtetĂ« se tani
Ti e ke kuptuar që ndodhem unë si kafsha më e tmerrshme në pyll?
Julia Donaldson, «Gruffalo»
Mendoj se po të organizonim një garë mes kolegëve të mi: kush e përgatit dhe e nis më shpejt ETL-procesin nga fillimi: ata me SSIS dhe miniku, dhe unë me Airflow⊠Pastaj do të krahasonim edhe lehtësinë e mirëmbajtjes⊠Ehm, mendoj se do të pajtohesh se do t'i kaloj në çdo front!
NĂ«se flasim pak mĂ« seriozisht, atĂ«herĂ« Apache Airflow â pĂ«rmes pĂ«rshkrimit tĂ« proceseve nĂ« formĂ«n e kodit programues â e bĂ«ri punĂ«n time shumĂ« mĂ« mĂ« tĂ« lehtĂ« dhe mĂ« tĂ« kĂ«ndshme.
Po ashtu, shkallĂ«zueshmĂ«ria e tij e pakufizuar: si nĂ« aspektin e plugins, ashtu edhe prirja pĂ«r tĂ« shkallĂ«zuar â ju jep mundĂ«sinĂ« tĂ« pĂ«rdorni Airflow nĂ« pothuajse çdo fushĂ«: qoftĂ« nĂ« ciklin e plotĂ« tĂ« mbledhjes, pĂ«rgatitjes dhe pĂ«rpunimit tĂ« tĂ« dhĂ«nave, qoftĂ« nĂ« raketat qĂ« lĂ«shohen (nĂ« Mars, sigurisht).
Pjesa përfundimtare, informative dhe reference
Gabimet që i kemi mbledhur për ju
start_date. Po, kjo Ă«shtĂ« njĂ« meme lokale. PĂ«rmes argumentit kryesor tĂ« DAG-utstart_datekalojnĂ« tĂ« gjitha. Shkurtimisht, nĂ«se caktojmĂ« nĂ«start_datedatĂ«n aktuale, dhe nĂ«schedule_intervalâ njĂ« ditĂ«, atĂ«herĂ« DAG-u do tĂ« nisĂ« nesĂ«r, mĂ« herĂ«t se.start_date = datetime(2020, 7, 7, 0, 1, 2)Dhe mĂ« shumĂ« asnjĂ« problem.
Me të lidhet dhe një tjetër gabim në ekzekutim:
Task is missing the start_date parameter, që zakonisht tregon se e keni harruar të lidhni me operatorin e DAG-ut.- Të gjitha në një makinë. Po, si bazat (e vetë Airflow dhe mbulesës sonë), ashtu edhe serveri web, planifikuesi dhe punëtorët. Dhe madje punonte. Por me kalimin e kohës numri i detyrave në shërbime u rrit, dhe kur PostgreSQL filloi të jepte përgjigje nëpërmjet indekseve për 20 ms në vend të 5 ms, ne e morëm dhe e trezoruam.
- LocalExecutor. Po, ne kemi qenë ende në të, dhe tashmë kemi arritur në skajin e greminës. LocalExecutor ka qenë i mjaftueshëm për ne deri tani, por tani ka ardhur koha për të zgjeruar me të paktën një punëtor, dhe do të duhet të mundohemi të kalojmë në CeleryExecutor. Dhe duke marrë parasysh se mund të punosh me të edhe në një makinë, asgjë nuk na ndalon nga përdorimi i Celery edhe në serverin që "natyrisht, kurrë nuk do të shkojë në prodhim, të premtoj!"
- Përdorim i papërfshirë mjetet e ndërtuara:
- Connections për ruajtjen e të dhënave të identifikimit të shërbimeve,
- SLA Misses për të reaguar ndaj detyrave që nuk përfunduan në kohë,
- XCom për ndarjen e metadatat (thashë metatë dhënat!) midis detyrave të dagut.
- Abuzimi me postĂ«n. ĂfarĂ« tĂ« thuash kĂ«tu? Ishin krijuar njoftime pĂ«r tĂ« gjitha pĂ«rsĂ«ritjet e detyrave tĂ« rĂ«na. Tani nĂ« Gmail-in tim tĂ« punĂ«s kam >90k letra nga Airflow, dhe ndĂ«rfaqja web e postĂ«s refuzon tĂ« marrĂ« dhe tĂ« fshijĂ« mĂ« shumĂ« se 100 njĂ«si pĂ«r herĂ«.
Më shumë pengesa:
Mjetet për automatizim më të madh
Për të punuar edhe më shumë me mendje dhe jo me duar, Airflow na përgatiti këtë:
- â ai ende ka status Experimental, qĂ« nuk e pengon tĂ« funksionojĂ«. Me kĂ«tĂ«, jo vetĂ«m qĂ« mund tĂ« merrni informacion mbi dagat dhe detyrat, por mund tĂ« ndaloni/nisni dagun, tĂ« krijoni DAG Run ose grup.
- â pĂ«rmes komandĂ«s nĂ« terminal janĂ« tĂ« ndAvailable shumĂ« mjete, tĂ« cilat jo vetĂ«m qĂ« janĂ« tĂ« pakĂ«ndshme pĂ«r t'u pĂ«rdorur pĂ«rmes WebUI, por edhe nuk ekzistojnĂ« fare. PĂ«r shembull:
backfillnevojiten për të rinisur instancat e detyrave.
Për shembull, vijnë analistët dhe thonë: «E keni, shoku, problem me të dhënat nga 1 deri më 13 janar! Rregulloni!». Dhe ti je si:airflow backfill -s '2020-01-01' -e '2020-01-13' orders- Mirëmbajtja e bazës:
initdb,resetdb,upgradedb,checkdb. run, i cili lejon të nisni një instancë të detyrës dhe t'i injoroni të gjitha varësitë. Për më tepër, mund ta nisni atë përmesLocalExecutor, madje edhe sikur të keni një klaster Celery.- Përgjithësisht, e njëjta gjë bënë
test, vetëm se nuk shkruan asgjë në bazë. connectionslejon krijimin masiv të lidhjeve nga shell.
- â njĂ« mĂ«nyrĂ« mjaft hardcore pĂ«r tĂ« bashkĂ«vepruar, e cila Ă«shtĂ« e destinuar pĂ«r plugina, dhe jo pĂ«r tĂ« grumbulluar me duar. Por kush na ndalon tĂ« shkojmĂ« nĂ«
/home/airflow/dags, të nxjerrimipythondhe të fillojmë të bëjmë çfarë deshi? Mund, për shembull, të eksportoni të gjitha lidhjet me këtë kod:nga airflow import settings nga airflow.models import Connection fushat = 'conn_id conn_type host port schema login password extra'.split() session = settings.Session() per conn në session.query(Connection).order_by(Connection.conn_id): d = {field: getattr(conn, field) për fushat} print(conn.conn_id, '=', d) - Lidhja me bazën e të dhënave të metadatas Airflow. Nuk e rekomandoj të shkruani në të, por merrni statet e detyrave për metrika të ndryshme specifike shumë më shpejt dhe më lehtë se përmes ndonjë API.
TĂ« themi, jo tĂ« gjitha detyrat tona janĂ« idempotente, dhe ndonjĂ«herĂ« ato mund tĂ« dĂ«shtojnĂ«, dhe kjo Ă«shtĂ« e pranueshme. Por disa dĂ«shtime â kjo Ă«shtĂ« tashmĂ« e dyshimtĂ«, dhe duhet tĂ« verifikohet.
Kujdes, SQL!
ME ekzekutimet e fundit SI ( SELECT task_id, dag_id, execution_date, state, row_number() PĂR ( PARTITION BY task_id, dag_id ORDER BY execution_date DESC) SI rn FROM public.task_instance WHERE execution_date > now() - INTERVAL '2' DITĂ ), failed AS ( SELECT task_id, dag_id, execution_date, state, RAST CASE kur rn = row_number() PĂR ( PARTITION BY task_id, dag_id ORDER BY execution_date DESC) THEN TRUE END SI last_fail_seq FROM last_executions WHERE state IN ('dĂ«shtuar', 'nĂ« pritje pĂ«r rikthim') ) SELECT task_id, dag_id, count(last_fail_seq) SI tĂ« pasuksesshĂ«m, count(CASE KUR last_fail_seq DHE state = 'dĂ«shtuar' ATĂHERĂ 1 END) SI dĂ«shtuar, count(CASE KUR last_fail_seq DHE state = 'nĂ« pritje pĂ«r rikthim' ATĂHERĂ 1 END) SI nĂ« pritje pĂ«r rikthim FROM failed GRUPI PĂR task_id, dag_id KUSH count(last_fail_seq) > 0
Linke
Natyrisht, dhjetë lidhjet e para nga rezultatet e Google përmbajnë dosjen e Airflow nga shënimet e mia.
- â natyrisht, duhet tĂ« filloni me dokumentacionin zyrtar, por kush i lexon udhĂ«zimet?
- â tĂ« paktĂ«n lexoni rekomandimet nga krijuesit.
- â fillimi: ndĂ«rfaqja pĂ«rdoruese nĂ« figura.
- â janĂ« mirĂ«shpjeguara bazat, nĂ«se (rastĂ«sisht!) nuk e keni kuptuar diçka nga unĂ«.
- â njĂ« udhĂ«zues i shkurtĂ«r pĂ«r konfigurimin e klasterit Airflow.
- â njĂ« artikel thuajse po aq interesant, pĂ«rveç se ka mĂ« shumĂ« formalizĂ«m dhe mĂ« pak shembuj.
- â pĂ«r funksionimin nĂ« bashkĂ«punim me Celery.
- â mbi idempotencĂ«n e detyrave, ngarkimin sipas ID-sĂ« nĂ« vend tĂ« datĂ«s, transformimet, strukturĂ«n e skedarĂ«ve dhe gjĂ«ra tĂ« tjera interesante.
- â varĂ«sitĂ« e detyrave dhe Rregulli i Aktivizimit, tĂ« cilat i pĂ«rmenda vetĂ«m sipas rastit.
- â si tĂ« kaloni disa nga "punon siç Ă«shtĂ« parashikuar" nĂ« planifikues, tĂ« ngarkoni tĂ« dhĂ«nat e humbura dhe tĂ« vendosni prioritetet e detyrave.
- â SQL pyetjet e dobishme pĂ«r metadaten e Airflow.
- â ka njĂ« seksion tĂ« dobishĂ«m pĂ«r krijimin e sensorĂ«ve personalizues.
- â njĂ« shĂ«nim interesant mbi ndĂ«rtimin e infrastrukturĂ«s nĂ« AWS pĂ«r ShkencĂ«n e tĂ« DhĂ«nave.
- â gabime tĂ« zakonshme (kur dikush gjithsesi nuk lexon udhĂ«zimet).
- â buzĂ«qeshni, si njerĂ«zit rregullojnĂ« ruajtjen e fjalĂ«kalimeve, kur mund tĂ« pĂ«rdorin thjesht Connections.
- â kalimi i paqartĂ« i DAG, kalimi i kontekstit nĂ« funksione, sĂ«rish pĂ«r varĂ«sitĂ«, dhe gjithashtu pĂ«r humbjen e nisjeve tĂ« detyrave.
- â mbi pĂ«rdorimin e
argumenteve standardedheparamsnĂ« shabllonat, si dhe pĂ«r Variables dhe Connections. - â njĂ« tregim mbi atĂ« se si pĂ«rgatitet planifikuesi pĂ«r Airflow 2.0.
- â njĂ« artikull disi i vjetĂ«r mbi implementimin e grupit tonĂ« nĂ«
docker-compose. - â detyra dinamike pĂ«rmes shablloneve dhe kalimit tĂ« kontekstit.
- â njoftime standarde dhe tĂ« personalizuara pĂ«rmes postĂ«s elektronike dhe Slack.
- â DegĂ«tim i detyrave, makros dhe XCom.
Dhe lidhjet e përdorura në artikull:
- â plĂ«hsitĂ« qĂ« mund tĂ« pĂ«rdoren nĂ« shabllon.
- â Gabime tĂ« zakonshme gjatĂ« krijimit tĂ« DAG-ve.
- â
docker-composepĂ«r eksperimente, debuggim dhe mĂ« shumĂ«. - â MbĂ«shtjellĂ«si Python pĂ«r Telegram REST API.
Burimi: habr.com




