Tere, mina olen Dmitri Logvinenko â andengineer grupi firmade «Vezot» analĂŒĂŒtika osakonnas.
Kojutame teile suurepĂ€rasest tööriistast ETL-protsesside arendamiseks â Apache Airflow. Kuid Airflow on nii universaalne ja mitmekesine, et tasub sellele tĂ€helepanu pöörata ka siis, kui te ei tegele andmevoogudega, vaid vajate aeg-ajalt mĂ”ne protsessi kĂ€ivitamist ja nende tĂ€itmise jĂ€lgimist.
Jah, ma ei rÀÀgi ainult, vaid nÀitan ka: programmi sisaldab palju koodi, ekraanipilte ja soovitusi.

Mida tavaliselt nÀed, kui googeldad sÔna Airflow / Wikimedia Commons
Sisukord
Sissejuhatus
Apache Airflow â ta on nagu Django:
- kirjutatud Pythonis,
- seal on suurepÀrane adminpaneel,
- piiramatu laiendatavusega,
â ainult parim, ja see on loodud hoopis muude eesmĂ€rkide jaoks, nimelt (nagu enne katset):
- ĂŒlesannete kĂ€ivitamine ja jĂ€lgimine piiramatul arvul masinatel (nii palju kui teie sĂŒdametunnistus ja Celery/Kubernetes lubavad)
- dĂŒnaamilise töövoo genereerimise abil, kasutades vĂ€ga lihtsat Python-koodi, mis on lihtne kirjutada ja mĂ”ista
- ning vÔimalusega siduda omavahel mis tahes andmebaasid ja API-d valmiskomponentide ning isetehtud pluginatega (mida on ÀÀrmiselt lihtne teha).
Kasutame Apache Airflow'i jÀrgmiselt:
- kogume andmeid erinevatest allikatest (mitmed SQL Serveri ja PostgreSQL instance, erinevad API-d rakenduste mÔÔdikute jaoks, isegi 1C) DWH-sse ja ODS-sse (meie puhul Vertica ja Clickhouse).
- nagu arenenud
cron, mis kÀivitab andmete konsolideerimise protsesse ODS-is ja jÀlgib ka nende haldamist.
Kuni hiljuti katab meie vajadusi ĂŒks vĂ€ike server 32 tuuma ja 50 GB mĂ€luga. Airflow'is töötab sel juhul:
- ĂŒle 200 dag'i (töötavad workflows, kuhu oleme ĂŒlesandeid tĂ€itnud),
- igaĂŒhes keskmiselt 70 ĂŒlesannet,
- killedakse seda (jĂ€llegi keskmiselt) ĂŒks kord tunnis..
Ja sellest, kuidas me laienesime, kirjutan alla, aga nĂŒĂŒd mÀÀratleme ĂŒber-ĂŒlesande, mida me lahendame:
Kolmes SQL Serverit, igaĂŒhel 50 andmebaasi â ĂŒhe projekti instantsid, seega on neil struktuur peaaegu sama (peaaegu kĂ”ikjal, muha-haha), mis tĂ€hendab, et igas on tabel Orders (Ă”nneks saab sellise nimega tabelit panna igasse Ă€risse). Me vĂ”tame andmed, lisades teenindusvĂ€ljad (allikaserver, allikandmebaas, ETL-ĂŒlesande identifikaator) ning ĂŒritame nad sĂŒĂŒdistusteta visata, ĂŒtleme, Verticasse.
LĂ€hme!
Peamine, praktiline osa (ja natuke teooriat)
Miks see meile (ja teile) vajalik on
Kui puud olid suured ja mina olin lihtne SQL-hooldaja ĂŒhes Venemaa jaeketis, heitsime ETL-protsesside ehk andmevoogude kallale kahe kĂ€esolevate vahenditega:
- Informatica Power Center â ÀÀrmiselt keeruline sĂŒsteem, ÀÀrmiselt tootlik, oma riistvara ja versioneerimisega. Kasutasin vĂ”ib-olla 1% selle vĂ”imalustest. Miks? No, esmalt, see liides on kuskil nullindatest vaimselt tapnud meid. Teiseks, see vidin on loodud ÀÀrmiselt keerukate protsesside, intensiivse komponentide taaskasutamise ja teiste vĂ€ga-oluliste ettevĂ”tte-nippide jaoks. Me jĂ€tame vahele, et see maksab sama palju kui Airbus A380 tiib/ aasta.
Olge ettevaatlik, ekraanipilt vÔib inimestele, kes on alla 30, natuke haiget teha.

- SQL Server Integration Server â seda tööriista oleme kasutanud oma siseprojektide voogudes. Tegelikult kasutame me juba SQL Serverit ja selle ETL tööriistade mittekasutamine oleks olnud ebaaus. KĂ”ik on selles hĂ€sti: liides on ilus, ja töötlemisaruanded... Aga me ei armasta tarkvaratooteid selle pĂ€rast, oh ei. Versioonide haldamine
dtsx(mis on XML, millel on salvestamisel segased sĂ”lmed) on meil vĂ”imalik, aga mis sellest kasu? Kuidas luua ĂŒlesannete pakett, mis tĂ”mbab sada tabelit ĂŒhest serverist teise? Mis sada, kahekĂŒmne tabeliga kaotab teil sĂ”rm, mis klĂ”psab rottnuppu. Kuid see nĂ€eb kindlasti stiilsem vĂ€lja:
Oleme absoluutsetelt otsinud vÀljapÀÀse. Asi oli peaaegu peaaegu jÔudnud kohandatud SSIS-pakettide generaatorini...
... ja siis leidis mind uus töö. Ja seal tabas mind Apache Airflow.
Kui ma sain teada, et ETL-protsesside kirjeldused on lihtsalt Python'i kood, ei osanud ma rÔÔmust tantsida. Nii andmevood lĂ€bisid versioonihaldust ja diffimist, ning tabelite koondamine ĂŒhise struktuuriga sadadest andmebaasidest ĂŒhte sihtkohta muutus Python'i koodi kettaks 13â ekraanil.
Klastri kokkupanemine
Ărge muretsege siin lasteaiatuks, rÀÀgime hoopis asjadest, mis on ilmselged, nagu Airflow installimine, teie valitud andmebaas, Celery ja muud teemad, mis on dokumentides kirjeldatud.
Selleks, et me vÔiksime kohe katsetama asuda, olen ma koostanud docker-compose.yml mille sisu on jÀrgmine:
- KĂ€ivitame enda Airflow: Scheduler, Webserver. Seal töötab ka Flower, et jĂ€lgida Celery ĂŒlesandeid (kuna see on juba pakitud
apache/airflow:1.10.10-python3.7, ja me pole selle vastu); - PostgreSQL, kuhu Airflow kirjutab oma logifailid (plaanijate andmed, tĂ€itmise statistika jne), ja Celery mĂ€rkis lĂ”petatud ĂŒlesandeid;
- Redis, mis toimib Celery jaoks ĂŒlesannete vahendajana;
- Celery worker, mis tĂ”eliselt tegeleb ĂŒlesannete tĂ€itmisega.
- Katalooge
./dagskasutame meie DAGi kirjeldusfailide salvestamiseks. Need tuvastatakse reaalajas, seega ei ole vaja kogu steki iga vÀiksema muudatuse jÀrel edastada.
MĂ”nes kohas on kood nĂ€idetes osaliselt esitatud (et mitte teksti ĂŒle koormata), ja mĂ”nes kohas muudetakse seda protsessi kĂ€igus. Terveid ja töötavaid koodinĂ€iteid saab vaadata hoidlas. .
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 as a Celery broker
broker:
image: redis:6.0.5-alpine
# DB for the Airflow metadata
airflow-db:
image: postgres:10.13-alpine
environment:
- POSTGRES_USER=airflow
- POSTGRES_PASSWORD=airflow
- POSTGRES_DB=airflow
volumes:
- ./db:/var/lib/postgresql/data
# Main container with Airflow Webserver, Scheduler, Celery Flower
airflow:
<<: *airflow-base
environment:
<
-c " sleep 10 &&
pip install --user -r /requirements.txt &&
/entrypoint initdb &&
(/entrypoint webserver &) &&
(/entrypoint flower &) &&
/entrypoint scheduler"
ports:
# Celery Flower
- 5555:5555
# Airflow Webserver
- 8080:8080
# Celery worker, will be scaled using `--scale=n`
worker:
<<: *airflow-base
environment:
<
-c " sleep 10 &&
pip install --user -r /requirements.txt &&
/entrypoint worker"
depends_on:
- airflow
- airflow-db
- brokerMĂ€rkused:
- Kompositsiooni koostamisel toetus mulle palju tuntud mudelid. â kindlasti vaadake. VĂ”ib-olla ei vajagi te elus midagi muud.
- KÔik Airflow seadistused on saadaval mitte ainult lÀbi
airflow.cfg, vaid ka keskkonnaparameetrite kaudu (tĂ€nu arendajatele), mille kaudu ma laialdaselt kasutasin. - Muidugi ei ole see production-ready: ma ei seadnud konteineritele teadetega ĂŒhendust ja ei aja turvalisusega jamas. Aga miinimum, mis sobib meie katsetamiseks, on loomulikult olemas.
- Pange tÀhele, et:
- DAG-failide kaust peab olema juurdepÀÀsetav nii ajastajale kui ka töötajatele.
- Sama kehtib ka kĂ”ikide kolmandate osapoolte raamatukogude kohta â need peavad olema installitud masinatesse, kus on ajastaja ja töötajad.
NĂŒĂŒd on see lihtne:
$ docker-compose up --scale worker=3PĂ€rast seda, kui kĂ”ik on ĂŒles tĂ”usnud, saab vaadata veebiliideseid:
- Airflow:
- Flower:
Peamised mÔisted
Kui te ei saanud aru nendest "DAGidest", siis siin on lĂŒhike sĂ”nastik:
- Ajastaja â kĂ”ige tĂ€htsam tegelane Airflow's, kes jĂ€lgib, et töötehnikud teevad tööd ja mitte inimene: jĂ€lgib graafikut, vĂ€rskendab DAG-e, kĂ€ivitab ĂŒlesandeid.
Ăldiselt, vanades versioonides olid tal mĂ€luga (ei, mitte amneesia, vaid leketega) probleemid ja konfiguratsioonides jĂ€i isegi alles pĂ€randparameeter.
run_durationâ tema taaskĂ€ivitamisintervall. Kuid nĂŒĂŒd on kĂ”ik korras. - DAG (niisugune "dag") â "suunatud aktsĂŒkliline graaf", kuid selline mÀÀratlemine ĂŒtleb vĂ€he kellelegi, kuigi tegelikult on see konteiner omavahel suhtlevate ĂŒlesannete jaoks (vt allpool) vĂ”i analoog paketile SSIS ja töövoogudele Informatica-s.
Lisaks DAG-idele vÔivad olla veel subdag-id, kuid me ei jÔua neid tÔenÀoliselt kÀsitleda.
- DAG Run â initsialiseeritud DAG, millele on antud oma
execution_date. Ăhe DAG-i DAG-runnid vĂ”ivad kenasti töötada paralleelselt (kui olete oma ĂŒlesanded idempotentseteks teinud). - Operator â need on koodijupid, mis vastutavad konkreetse tegevuse tĂ€itmise eest. On kolm tĂŒĂŒpi operaatorit:
- action, nagu nÀiteks meie lemmik
PythonOperator, mis on vĂ”imeline tĂ€itma igasugust (kehtivat) Python-koodi; - transfer, mis viivad andmeid ĂŒhest kohast teise, nĂ€iteks
MsSqlToHiveTransfer; - sensor , mis vĂ”imaldab reageerida vĂ”i peatada DAG-i edasise tĂ€itmise mingisuguse sĂŒndmuse toimumiseni.
HttpSensorsuudab kĂŒsida mÀÀratud lĂ”pp-punkti ja kui ootab vajalikku vastust, alustada ĂŒlekannetGoogleCloudStorageToS3Operator. Uudishimu kĂŒsib: âmiks? KĂ”ike saab ju teha otse operaatori sees!â JĂ€rgnevalt, et mitte koormata tĂ¶Ă¶ĂŒlesannete panga riputanud operaatoritega. Sensor kĂ€ivitub, kontrollib ja sureb jĂ€rgmise katse ootamiseks.
- action, nagu nÀiteks meie lemmik
- Ălesanne â kuulutatud operaatorid, sĂ”ltumata tĂŒĂŒbist ja seotud DAGiga, tĂ”stetakse ĂŒlesande tasemele.
- Ălesande instants â kui peaplanner otsustab, et ĂŒlesanded on valmis heidma vĂ”itlejate töödele (otse kohal, kui me kasutame
LocalExecutorvĂ”i kaugnĂ”lva puhulCeleryExecutor), mÀÀrab ta neile konteksti (st muutujate komplekti â tĂ€itmise parameetrid), rakendab kĂ€su vĂ”i pĂ€ringu malle ning paneb need pangale.
Ălesannete genereerimine
Esmalt mÀÀratlege meie DAGi ĂŒldine skeem ja seejĂ€rel sĂŒveneme ĂŒha rohkem detailidesse, sest rakendame mĂ”ningaid ebatavalisi lahendusi.
Nii et kÔige lihtsamas vormis nÀeb selline DAG vÀlja nii:
impordi {timedelta, datetime}
from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from 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)
for conn_id, schema in sql_server_ds:
PythonOperator(
task_id=schema,
python_callable=workflow,
provide_context=True,
dag=dag)Alustame:
- Esimesena impordime vajalikud teegid ja veel mÔned asjad;
sql_server_dsâ see onList[namedtuple[str, str]], mis sisaldab Airflow Connections-i ĂŒhenduste nimesid ja andmebaase, millest me oma tabelit vĂ”tame;dagâ meie DAG-i deklareerimine, mis peab kindlasti olemaglobals(), muidu Airflow ei leia seda. DAG-ile tuleb ka öelda:- kuidas seda nimetada
ordersâ see nimi ilmub hiljem veebiliideses, - et see hakkab tööle alates 8. juulist keskööst,
- ja et selle ajakavad on umbes iga 6 tunni jÀrel (kÔvadele kutidele on siin 'timedelta()' asemel lubatud
-string, vĂ€hem kĂ”vadele kutidele - lause nagucron@daily0 0 0/6 ? * * *workflow()teeb pĂ”hiosa tööst, aga mitte nĂŒĂŒd. Hetkel lihtsalt vĂ€ljastame meie konteksti logisse.);
- kuidas seda nimetada
Ja nĂŒĂŒd lihtne maagia ĂŒlesannete loomise osas:kĂ”nnime allikate jĂ€rgi;- algatame
- , mis tĂ€idab meie tĂŒhi funktsioon
- initsialiseerime
PythonOperator, mis tĂ€idab meie tĂŒhikutJa nĂŒĂŒd lihtne maagia ĂŒlesannete loomise osas:. Ăra unusta anda unikaalne (DAGi piires) ĂŒlesande nimi ning siduda see DAGiga. Lippprovide_contexttoob omakorda funktsiooni tĂ€iendavad argumendid, mille me ettevaatlikult kogume**context.
Siinkohal lÔpetame. Mida me oleme saanud:
- uus DAG veebiliideses,
- kahesaja viiskĂŒmmend ĂŒlesannet, mis töötavad paralleelselt (kui Airflow, Celery ja serverite seaded seda lubavad).
Noh, peaaegu saime.

Kes paigaldab sÔltuvused?
Kogu seda asja lihtsustada, lĂŒkkasin ma docker-compose.yml töötlemisse requirements.txt kĂ”ikidesse nodidesse.
NĂŒĂŒd lĂ€heb asi kiireks:

Hallid ruudud â ĂŒlesande instantsid, mida planeerija on töötlenud.
Ootame veidi, ĂŒlesanded haaravad töötajad:

Rohelised, loomulikult, â edukalt lĂ€bitud. Punased â mitte nii edukalt.
Muide, meie tootmises ei ole mingit kausta
./dags, mis sĂŒnkroniseerub masinate vahel â kĂ”ik DAGid asuvadgitmeie Gitlabis, ja Gitlab CI seab uuendused masinatesse, kui need on kokku pandudmaster.
Natuke Flowerist
Kuni töötajad töötlevad meie tĂŒhje ĂŒlesandeid, tuletame meelde veel ĂŒhe tööriista, mis vĂ”ib meile mĂ”ningaid asju nĂ€idata â Flower.
Esimene leht, kus on kokkuvÔtete info töötajate nodide kohta:

KĂ”ige sisukam leht ĂŒlesannetega, mis on töös:

KÔige igavam leht meie maaklerite seisuga:

KĂ”ige silmatorkavam leht â ĂŒlesannete seisugraafik ja nende tĂ€itmise ajad:

Laadime alla, mis ei ole laaditud
Nii, kĂ”ik ĂŒlesanded on töös, saame kannatanud vĂ€lja viia.

Ja kannatanuid oli pĂ€ris mitu â erinevatel pĂ”hjustel. Kui Airflowd Ă”igesti kasutada, siis need ruudud nĂ€itavad, et andmed pole kindlasti kohale jĂ”udnud.
Tuleb vaadata logi ja taaskĂ€ivitada kukkunud ĂŒlesandeinstantsid.
Klikkides mÔnel ruudul, nÀeme meie jaoks saadaval olevaid tegevusi:

Saame vĂ”tta ja teha kukkunutele Clear. See tĂ€hendab, et unustame, et seal midagi takerdus, ja sama ĂŒlesandeinstants suundub planeerijale.

Selge, et teiste punaste ruutudega hiirega nii kĂ€ituda ei ole kuigi inetu â mitte seda me Airflowlt ootame. Loomulikult on meil massihĂ€vitusrelv: Browse/Task Instances

Valime kÔik korraga ja nullime, vajutame Ôiget punkti:

PÀrast puhastamist nÀevad meie taksod vÀlja nii (nad juba ootavad, millal planeerija nad plaanib):

Ăhendused, hooikud ja muud muutujad
On viimane aeg vaadata jÀrgmist DAG-i, 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='DWH Reports updated',
html_content=dedent("""Tere, head inimesed, raportid on uuendatud"""),
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, Ă€rka ĂŒles, me kukutasime {{ dag.dag_id }}
"""),
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]Kas oleme kĂ”ik kunagi aru saanud aruandeid uuendama? Siin see on jĂ€lle: meil on loetelu allikatest, kust andmed vĂ”tta; on loetelu, kuhu need panna; Ă€rme unustame signaale anda, kui kĂ”ik on sĂŒndinud vĂ”i purunenud (aga see ei kehti meie kohta, eks?).
Vaadakem uuesti faili ja vaatame uusi mÔistetavaid asju:
from commons.operators import TelegramBotSendMessageâ mis takistab meid oma operaatorite loomisel, ja me kasutasime seda Ă€ra, tehes vĂ€ikese mĂ€hise sĂ”numite saatmiseks Razblokirovanny'sse. (Selle operaatori kohta rÀÀgime veel allpool);default_args={}â DAG vĂ”ib jagada samu argumente kĂ”igile oma operaatoritele;to='{{ var.value.all_the_kings_men }}'â vĂ€litomeil ei ole seda kĂ”vakooditud, vaid dĂŒnaamiliselt genereeritav Jinja ja muutujaga, kus on e-posti aadresside loend, mille ma ettevaatlikult paninAdmin/Variables;trigger_rule=TriggerRule.ALL_SUCCESSâ operaatori kĂ€ivitamise tingimus. Meie juhul saadetakse kiri juhtidele ainult juhul, kui kĂ”ik sĂ”ltuvused on töötanud edu;tg_bot_conn_id='tg_main'â argumendidconn_idvĂ”tavad endasse ĂŒhenduste identifikaatorid, mille me loomeAdmin/Connections;trigger_rule=TriggerRule.ONE_FAILEDâ sĂ”numid Telegramis saadetakse ainult siis, kui on ebaĂ”nnestunud ĂŒlesandeid;task_concurrency=1â keelame ĂŒhe ĂŒlesande mitme taski instantsi samaaegse kĂ€itamise. Vastasel juhul saame mitu ĂŒlesande instantsi,VerticaOperator(mis vaatavad ĂŒhte tabelit);report_update >> [email, tg]â kĂ”ikVerticaOperatorkoondatakse e-kirjade ja sĂ”numite saatmiseks, nagu jĂ€rgneb:

Kuna aga teavitajate operaatoritel on erinevad kĂ€ivitamistingimused, töötab vaid ĂŒks. Tree View'is nĂ€eb see vĂ€lja vĂ€hem lĂ€bipaistev:

RÀÀgin natuke makrode ja nende sĂ”prade â muutujate.
Makrod on Jinja-sisu kohandajad, mis vÔivad asetada erinevat kasulikku teavet operaatorite argumentidesse. NÀiteks nii:
SELECT
id,
payment_dtm,
payment_type,
client_id
FROM orders.payments
WHERE
payment_dtm::DATE = '{{ ds }}'::DATE{{ ds }} laiendab konteksti muutujate sisu execution_date vormingus YYYY-MM-DD: 2020-07-14. Parim on see, et konteksti muutujad on kinnitatud kindla taski instantsi (ruudu) kĂŒlge Tree View'is, ja kui instantsi uuesti kĂ€ivitada, lahetakse kohandajad samadele vÀÀrtustele.
MÀÀratud vÀÀrtusi saab vaadata igas taski instantsis nupul Rendered. Nii on e-kirjade saatmise ĂŒlesande puhul:

Ja nii on sĂ”numi saatmise ĂŒlesande puhul:

Viimane koodiversioon sisaldab kÔiki integreeritud makrode loendit, mis on saadaval siin:
Lisaks saame pluginade abil luua oma makrosid, kuid sellest rÀÀgime hiljem.
Eeldefineeritud asjade kÔrval saame sisestada ka oma muutuja vÀÀrtusi (nagu olen juba koodis teinud). Loome mÔned: Admin/Variables paar asja:

Nii, nĂŒĂŒd saab kasutada:
TelegramBotSendMessage(chat_id='{{ var.value.failures_chat }}')VÀÀrtuses vÔib olla skalaarsed andmed vÔi isegi JSON. JSON-i puhul:
bot_config
{
"bot": {
"token": 881hskdfASDA16641,
"name": "Verter"
},
"service": "TG"
}kasutame lihtsalt Ôige vÔtme teed: {{ var.json.bot_config.bot.token }}.
Ătlen lihtsalt ĂŒhe sĂ”na ja nĂ€itan ĂŒhte ekraanipilti seoses ĂŒhendustega. Siin on kĂ”ik lihtne: lehe peal Admin/Connections loome ĂŒhenduse, salvestame sinna oma kasutajanimesid/paroolid ja spetsiifilisemad parameetrid. Nii:

Paroolid saab krĂŒpteerida (tĂ€iendavalt, kui vĂ”rrelda vaikeseadetega) vĂ”i ei pea ka ĂŒhenduse tĂŒĂŒpi mainima (nagu ma tegin) tg_main) â asi on selles, et Airflow mudelites on tĂŒĂŒpide nimekiri sisse ehitatud ja seda ei saa muudetud allikateta laiendada (kui ma midagi valesti otsisin â palun parandage mind), kuid me ei takista midagi, et saada akrediteeringud lihtsalt nime jĂ€rgi.
Samuti saab luua mitu ĂŒhendust sama nimega: sellisel juhul meetod BaseHook.get_connection(), mis toob meile ĂŒhendused nime jĂ€rgi, annab juhusliku mitme sarnase seas (loogilisem oleks olnud kasutada Round Robin'i, kuid jĂ€tame selle Airflow arendajate sĂŒdametunnistusele).
Muutujad ja ĂŒhendused on kindlasti suurepĂ€rased tööriistad, kuid on oluline hoida tasakaalu: millised osad teie voogudest te hoiate koodis ja milliseid â annate Airflowâle sĂ€ilitamiseks. Ăhelt poolt vĂ”ib vÀÀrtuse, nĂ€iteks postkasti, kiirelt muutmine UI kaudu olla mugav. Teiselt poolt on see siiski tagasipöördumine hiireklikkide juurde, millest me (mina) tahtsime vabaneda.
Ăhendustega töötamine on ĂŒks hook'ideĂŒlesandeid. Ăldiselt on Airflow hook'id side punktid kolmandate teenuste ja raamatukogudega. NĂ€iteks, JiraHook avatakse meie jaoks klient, et suhelda Jira'ga (saame ĂŒlesandeid edasi-tagasi liigutada), ja SambaHook abil saab pushida kohaliku faili smb-punkti.
Kohandatud operaatori analĂŒĂŒs
Ja me oleme lÀhenenud sellele, et vaadata, kuidas on tehtud TelegramBotSendMessage
Kood commons/operators.py isegi operaatori juurde:
from typing import Union
from airflow.operators import BaseOperator
from commons.hooks import TelegramBotHook, TelegramBot
class TelegramBotSendMessage(BaseOperator):
"""Saada sÔnum chat_id-le, kasutades TelegramBotHook'i
NĂ€ide:
>>> 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 }} ebaÔnnestus :(',
... 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'Saada "{self.message}" chat' + f'{self.chat_id}')
self.client.send_message(chat_id=self.chat_id,
message=self.message)Siin, nagu kÔik muu Airflow's, on kÔik vÀga lihtne:
- Oleme pÀrinud
BaseOperator, mis rakendab ĂŒsna palju Airflow'le spetsiifilisi asju (vaadake rahus) - Oleme kuulutanud vĂ€lja vĂ€ljad
template_fields, kus Jinja otsib makrosid töötlemiseks. - Oleme korraldanud Ôiged argumendid
__init__(), mÀÀrasid vaikimisi, kus vaja. - EeljÀtku initsialiseerimist ei unustatud.
- Avatud vastav hook.
TelegramBotHook, saime temalt klient-objekti. - Ălekirjutasime meetodi.
BaseOperator.execute(), mida Airflow kutsub, kui on aega operaatori kĂ€ivitamiseks â selles realiseerime pĂ”hitegevuse, unustamata logida. (Logime, muide, kohe.stdoutjastderrâ Airflow kĂ”ik pĂŒĂŒab kinni, pakib ilusti kokku ja paigutab Ă”igesse kohta.)
Vaadakem, mis meil on failis commons/hooks.py. Faili esimene osa, mis sisaldab hooki:
from typing import Union
from airflow.hooks.base_hook import BaseHook
from requests_toolbelt.sessions import BaseUrlSession
class TelegramBotHook(BaseHook):
"""Telegram Bot API hook
MĂ€rkus: lisage ĂŒhendus tĂŒhja Conn Type'iga ja Ă€rge unustage
tÀita 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.clientMa isegi ei tea, mida siin seletada, lihtsalt toon vÀlja mÔned olulised punktid:
- Kasutame pĂ€randit, mĂ”tleme argumentide ĂŒle â enamikul juhtudel on see ĂŒks:
conn_id; - Ăksikute standardmeetodite ĂŒletamine: piirdun
get_conn(), kus ma saan ĂŒhenduse parameetrid nime jĂ€rgi ja lihtsalt tĂ”mban sektsiooniextra(see vĂ€li on JSON-i jaoks), kuhu ma (kui ma ise juba selgitasin!) panin Telegrami boti tokeni:{"bot_token": "YOuRAwEsomeBOtToKen"}. - Loome meie
TelegramBot, edastades talle konkreetse tokeni.
Ja ongi kĂ”ik. Saame kliendi ĂŒlesandest lĂ€bi TelegramBotHook().clent vĂ”i TelegramBotHook().get_conn().
Ja teises osas faile, kus ma teen mikropakkumise Telegram REST API jaoks, et mitte vedada sama ĂŒhe meetodi jaoks 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))Ăige tee on kĂ”ik see kokku panna:
TelegramBotSendMessage,TelegramBotHook,TelegramBotâ plugina, panna avalikku hoidlasse ja jagada avatud lĂ€htekoodiga.
Kuna me kĂ”ike seda uurisime, olid meie raporti vĂ€rskendused edukalt ilmu t seekord kokku kukkunud ja saatsid mulle kanalis vea teate. Pean minema kontrollima, mis seekord valesti lĂ€ksâŠ

Meie dags on katki! Ja ei olnud just seda, mida me ootasime? Just nimelt!
Kas hakkad valama?
Kas tunnete, et jĂ€tan midagi tĂ€helepanuta? Tundub, et lubasin andmed SQL Serverist Vertica'le ĂŒle kanda, ja nĂŒĂŒd lĂ€ksin teemadest kĂ”rvale, kurat!
Kuritegu oli tahtlik, ma pidin teile selgitama mĂ”ningaid termineid. NĂŒĂŒd saame liikuda edasi.
Meie plaan oli jÀrgmine:
- Luua dag
- Genererida ĂŒlesanded
- Vaadata, kui ilus kÔik vÀlja nÀeb
- Kinnituselamistele jagada sessiooninumbreid
- Andmete vÔtmine SQL Serverist
- Andmete paigutamine Verticasse
- Statistika kogumine
Nii et selle kÔik kÀivitamiseks tegin vÀikese tÀienduse meie 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.pySeal me tÔstame:
- Vertica kui host
dwhkÔige vaikimisi seadistustega, - kolm SQL Serveri instantsi,
- tĂ€idame andmebaasid viimaste andmetega (Ă€rge lĂŒkake vaadata
mssql_init.py!)
KÀivitame kogu selle hea veidi keerulisema kÀsu abil kui eelmisel korral:
$ docker-compose -f docker-compose.yml -f docker-compose.db.yml up --scale worker=3Mida meie imetore juhuslik generaator genereeris, saab kasutada punkti kaudu Andmete profiilimine/Ad Hoc pÀring:

Peamine, et seda analĂŒĂŒtikutele ei nĂ€idata
Detailidesse laskuda ETL-seansid ei hakka, seal on kĂ”ik triviaalne: loome andmebaasi, sinna tabeli, mĂ€hkime kĂ”ik konteksti haldajasse ja nĂŒĂŒd teeme nii:
with Session(task_name) as session:
print('Load', session.id, 'started')
# Laadimise töövoog
...
session.successful = True
session.loaded_rows = 15session.py
from sys import stderr
class Session:
"""ETL töövoo sessioon
NĂ€ide:
with Session(task_name) as session:
print(session.id)
session.successful = True
session.loaded_rows = 15
session.comment = 'Tore töö'
"""
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 = """
LOONI TABEL, KUI EI OLE JUBA olemas sessioonid (
id SERIAL NOT NULL PRIMARY KEY,
task_name VARCHAR(200) NOT NULL,
started TIMESTAMPTZ NOT NULL DEFAULT current_timestamp,
finished TIMESTAMPTZ DEFAULT current_timestamp,
successful BOOL,
loaded_rows INT,
comment VARCHAR(500)
);
"""
self._execute(query)
def open(self):
query = """
INSERT INTO sessions (task_name, finished)
VALUES (%s, NULL)
RETURNING id;
"""
self._id = self._execute(query, self.task_name)
print(self, 'avatud')
return self
def close(self):
if not self._id:
raise SessionClosedError('Sessioon ei ole avatud')
query = """
UPDATE sessions
SET
finished = DEFAULT,
successful = %s,
loaded_rows = %s,
comment = %s
WHERE
id = %s
RETURNING id;
"""
self._execute(query, self.successful, self.loaded_rows,
self.comment, self.id)
print(self, 'suletud',
', edukas: ', self.successful,
', Laetud: ', self.loaded_rows,
', kommentaar:', self.comment)
class SessionError(Exception):
pass
class SessionClosedError(SessionError):
passOn aeg teha meie andmete vÔtmisest meie poolest sada tabelit. Teeme seda lihtsate ridade abil:
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)- Kasutame Airflow'd, et saada
pymssql-ĂŒhendust - Lisame pĂ€ringusse kuupĂ€eva piirangu â funktsioon lisab selle ĆĄabloonide kaudu.
- Saadame meie pÀringu
pandas, mis toob meileDataFrameâ see tuleb meile kasuks edaspidi.
Kasutame asendust
{dt}pÀringu parameetri asemel%smitte sellepÀrast, et ma oleks kuri Buratino, vaid sellepÀrast, etpandasei oska hakkama saadapymssqlja annab viimasele edasiparams: List, kuigi see vÀga tahabtuple.
Olge tĂ€helepanelik, et arendajapymssqlotsustas seda enam mitte toetada, ja on viimane aeg minna ĂŒlepyodbc.
Vaadakem, millised argumendid Airflow meie funktsioonidesse sisestas:

Kui andmeid polnud, siis pole mĂ”tet edasi minna. Kuid arvestada, et ĂŒleslaadimine oli edukas, on ka kummaline. Kuid see ei ole ka viga. Ah, mida teha?! Niimoodi:
if df.empty:
raise AirflowSkipException('No rows to load')AirflowSkipException Airflow ĂŒtleb, et viga ei ole, ja ĂŒlesanne jĂ€etakse vahele. Liideses ei ole rohelise ega punase ruudu asemel roosa vĂ€rv.
Lisame oma andmetele mitu veergu:
df['etl_source'] = src_schema
df['etl_id'] = session.id
df['hash_id'] = hash_pandas_object(df[['etl_source', 'id']])Konkreetsemalt:
- Andmebaas, kust me tellimused saime,
- Meie laadimise seansi identifikaator (see erineb kĂ”igil ĂŒlesannetel),
- Ressursi ja tellimuse identifikaatorist saadud hash â et lĂ”pp-andmebaasis (kus kĂ”ik kogutakse ĂŒhte tabelisse) oleks ainulaadne tellimuse identifikaator.
JÀÀnud on eelviimane samm: laadida kĂ”ik Verticasse. Ja, kummalisel kombel, ĂŒks kĂ”ige efektsematest viisidest selleks on CSV kaudu!
# 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)- Loome spetsiaalse vastuvÔtja
StringIO. pandasmis kenasti kogub meieDataFramekujulCSV-ridu.- Avame ĂŒhenduse meie lemmik Vertica hookiga.
- Ja nĂŒĂŒd kasutame
copy()saadame meie andmed otse Verticasse!
Draiverist saame, kui palju ridu laaditi, ja ĂŒtleme seansi juhile, et kĂ”ik on OK:
session.loaded_rows = cursor.rowcount
session.successful = TrueJa kÔik.
Tootmises loome sihttabeli kÀsitsi. Siin julgen lubada vÀikest automatiseerimist:
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)Ma olen abiks
VerticaOperator()loome andmebaasi skeemi ja tabeli (kui neid pole, muidugi). Peamine on Ôigesti seada sÔltuvused:
for conn_id, schema in 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 >> loadTeeme kokkuvÔtteid
â Noh, â ĂŒtles hiir, â ei ole tĂ”si, et nĂŒĂŒd
Sa veendusid, et metsas olen mina kÔige hirmsam metsaline?
Julia Donaldson, âGruffaloâ
Ma arvan, et kui me oma kolleegidega korraldaksime konkurentsi: kes suudab kiiremini koostada ja kĂ€ivitada ETL-protsessi nullist: nemad oma SSIS ja hiirega ning mina Airflowâga⊠Ja seejĂ€rel vĂ”rdleksime hooldamise mugavust⊠Oh, arvan, et nĂ”ustute, et ma tĂ”enĂ€oliselt ĂŒletan neid igal rindel!
Kui vĂ”tta asja veidi tĂ”sisemalt, siis Apache Airflow â tĂ€nu protsesside kirjeldamisele programmeerimiskoodina â tegi minu töö kaugelt mugavamaks ja meeldivamaks.
Tema piiramatud laienemisvĂ”imalused: nii pistikprogrammide osas kui ka skaleeritavuse poolest â annavad teile vĂ”imaluse rakendada Airflowâd praktiliselt igas valdkonnas: olgu see siis andmete kogumise, töötlemise ja ettevalmistamise tĂ€islĂŒka vĂ”i rakettide kĂ€ivitamine (Marsile, muidugi).
LÔpposa, viidatud ja informatiivne
KĂŒhvlid, mis me teie eest kogusime
start_date. Jah, see on juba kohalike meemide teemade ring. Peamise dĂŒnaamika argumendi kaudustart_datelĂ€bivad kĂ”ik. LĂŒhidalt, kui mÀÀratastart_datepraegune kuupĂ€ev jaschedule_intervalâ ĂŒks pĂ€ev, siis DAG kĂ€ivitub homme mitte varem.start_date = datetime(2020, 7, 7, 0, 1, 2)Ja rohkem ei mingeid probleeme.
Sellega on seotud ka ĂŒks teine tĂ€itmisviga:
Task is missing the start_date parameter, mis enamasti tĂ€hendab, et unustasite seondada DAG operaatoriga.- KĂ”ik ĂŒhel masinal. Jah, nii Airflow andmebaasid kui ka meie kattekiht, veebiserver, ajakavaja ja töötajad. Ja see töötas isegi. Kuid aja jooksul kasvas teenuste ĂŒlesannete arv ning kui PostgreSQL hakkas andma vastuse indeksi jĂ€rgi 20 ms asemel 5 ms, siis me vĂ”tsime selle ja viibisime Ă€ra.
- LocalExecutor. Jah, me kasutame seda endiselt ja oleme jĂ”udnud ÀÀre poole. LocalExecutor on seni olnud piisav, kuid nĂŒĂŒd on aeg vĂ€hemalt ĂŒhe töötaja vĂ”rra laieneda ja tuleb pingutada, et ĂŒle minna CeleryExecutorile. Kuna sellega saab töötada ka ĂŒhel masinal, ei ole miski takistuseks kasutada Celery't isegi serveris, mis "loomulikult ei lĂ€he kunagi tootmisse, ausĂ”na!"
- Kasutamata sisseehitatud tööriistu:
- Ăhendused teenuste autentimisandmete salvestamiseks,
- SLA puudumised ĂŒlesannete jaoks, mis ei töötanud Ă”igel ajal,
- XCom metaandmete vahetamiseks (ma ĂŒtlesin metaandmed!) ĂŒlesannete vahel DAG'is.
- E-kirjade liialdamine. Mis siin muud öelda? Oleme seadnud hoiatused kĂ”igi ebaĂ”nnestunud ĂŒlesannete korduste kohta. NĂŒĂŒd on mu töö Gmailis >90k kirja Airflow'ilt ja e-kirjade veebirakendus keeldub vĂ”tma ja kustutama rohkem kui 100 korraga.
Rohkem takistusi:
Veel suurema automatiseerimise vahendid
Kuna soovime veelgi rohkem mÔtlema hakata, mitte kÀtega töötada, on Airflow meie jaoks valmistanud jÀrgmise:
- â ta jĂ€tkuvalt omab Ekspressi staatust, mis ei takista selle tööleminekut. Selle abil on vĂ”imalus mitte ainult saada teavet DAGide ja ĂŒlesannete kohta, vaid ka peatada/algatada DAGi, luua DAGi jooksu vĂ”i hulki.
- â kĂ€surealt on kergesti kĂ€ttesaadavad paljud tööriistad, mis ei ole mitte ainult ebamugavad WebUI's, vaid ei ole seal ĂŒldse saadaval. NĂ€iteks:
backfillon vajalik ĂŒlesannete instantside uuesti kĂ€ivitamiseks.
NĂ€iteks kui analĂŒĂŒtikud tulevad ja ĂŒtlevad: «Teie andmeanalĂŒĂŒs on halb ja vajab parandamist 1. kuni 13. jaanuarini! Parandage, parandage, parandage, parandage!» Siis vĂ”tad ja teed:airflow backfill -s '2020-01-01' -e '2020-01-13' orders- Andmebaasi hooldus:
initdb,resetdb,upgradedb,checkdb. run, mis vĂ”imaldab kĂ€ivitada ĂŒhe ĂŒlesande instantsi, jĂ€ttes kĂ”ik sĂ”ltuvused tĂ€helepanuta. Veelgi enam, seda on vĂ”imalik kĂ€ivitada lĂ€biLocalExecutor, isegi kui teil on Celery-klaster.- Umbes sama asja teeb
test, ainult et ei kirjuta midagi andmebaasi. connectionsvĂ”imaldab massiliselt luua ĂŒhendusi shell'ist.
- â ĂŒsna hardcore viis suhtlemiseks, mis on mĂ”eldud pistikprogrammide jaoks, mitte kĂ€sitsi tööks. Aga kes meid kĂŒll takistab minemast
/home/airflow/dags, kĂ€ivitadaipythonja hakkama saama? NĂ€iteks on vĂ”imalik kĂ”ik ĂŒhendused eksportida sellise koodiga:from airflow import settings from airflow.models import Connection fields = 'conn_id conn_type host port schema login password extra'.split() session = settings.Session() for conn in session.query(Connection).order_by(Connection.conn_id): d = {field: getattr(conn, field) for field in fields} print(conn.conn_id, '=', d) - Airflow'i metaandmebaasiga ĂŒhendamine. Soovitan sinna kirjutada mitte, kuid erinevate spetsiifiliste mÔÔdikute jaoks ĂŒlesannete olekute kĂ€tte saamine on kindlasti kiirem ja lihtsam kui lĂ€bi ĂŒkskĂ”ik millise API.
Ătleme, et kĂ”ik meie ĂŒlesanded ei ole idempotentsed, ja need vĂ”ivad mĂ”nikord ebaĂ”nnestuda, ja see on normaalne. Kuid mitu kokkuvarisemist on juba kahtlane ning tuleks kontrollida.
Ole ettevaatlik, SQL!
WITH last_executions AS ( SELECT task_id, dag_id, execution_date, state, row_number() OVER ( PARTITION BY task_id, dag_id ORDER BY execution_date DESC) AS rn FROM public.task_instance WHERE execution_date > now() - INTERVAL '2' DAY ), failed AS ( SELECT task_id, dag_id, execution_date, state, CASE WHEN rn = row_number() OVER ( PARTITION BY task_id, dag_id ORDER BY execution_date DESC) THEN TRUE END AS last_fail_seq FROM last_executions WHERE state IN ('failed', 'up_for_retry') ) SELECT task_id, dag_id, count(last_fail_seq) AS unsuccessful, count(CASE WHEN last_fail_seq AND state = 'failed' THEN 1 END) AS failed, count(CASE WHEN last_fail_seq AND state = 'up_for_retry' THEN 1 END) AS up_for_retry FROM failed GROUP BY task_id, dag_id HAVING count(last_fail_seq) > 0
Lingid
Ja ja, ja esiteks kĂŒmme linki Google'i otsingust, mis viivad minu Airflow kaustade sisule.
- â muidugi tuleks alustada ametlikust dokumentatsioonist, aga kes neid juhiseid ikka loeb?
- â vĂ€hemalt loo autorite soovitusi loe.
- â esialgne tutvustus: kasutajaliides piltidena
- â peamised mĂ”isted on hĂ€sti kirjeldatud, juhul kui (Ă€kki!) sa ei saanud minust kĂ”igest aru.
- â lĂŒhike juhend Airflow-kliendi seadistamiseks.
- â peaaegu sama huvitav artikel, ainult et ametlikumat juttu on rohkem ja nĂ€iteid vĂ€hem.
- â koostööst Celeryga.
- â idempotentsuse, ĂŒlesannete laadimise ID jĂ€rgi mitte kuupĂ€eva, transformatsioonide, failistruktuuri ja muu huvitava kohta.
- â ĂŒlesannete sĂ”ltuvused ja Trigger Rule, millest ma ainult mĂ€rkisin.
- â kuidas ĂŒletada teatud «see töötab nagu kavandatud» planeerija puhul, laadida kadunud andmeid ja seada ĂŒlesannete prioriteedid.
- â kasulikud SQL-pĂ€ringud Airflow'i metaandmetele.
- â on kasulik jaotis kohandatud sensori loomisest.
- â huvitav lĂŒhike mĂ€rkus AWS-i andmeteaduse infrastruktuuri rajamisest.
- â levinud vead (kui keegi siiski ei loe juhiseid).
- â naeratage, kuidas inimesed paroolide salvestamist teevad, kuigi saaks lihtsalt kasutada Connections.
- â DAG-i vaikimisi edastamine, konteksti edastamine funktsioonides, taas sĂ”ltuvustest, samuti tööde kĂ€ivituste vahelejĂ€tmisest.
- â kasutamise kohta
vaikimisi argumendidjaparamsĆĄabloonides, samuti Variables ja Connections. - â jutustus sellest, kuidas ajakava Airflow 2.0-ks ette valmistatakse.
- â veidi aegunud artikkel meie klastrite kasutuselevĂ”tust
docker-compose. - â dĂŒnaamilised ĂŒlesanded ĆĄabloonide ja konteksti edastamise abil.
- â standardsete ja kohandatud teavitustega e-posti ja Slacki kaudu.
- â Ălesannete harud, makrod ja XCom.
Ja lingid, mis on artiklis kasutusel:
- â ĆĄablonites kasutamiseks saadaval olevad kohtade hoidjad.
- â Levinud vead dagide loomisel.
- â
docker-composeeksperimentideks, silumiseks ja muuks. - â Python-mĂ€his Telegram REST API jaoks.
Allikas: habr.com




