Tere, mina olen Dmitri Logvinenko â andmete insener ettevĂ”tte «VezŃŃ» analĂŒĂŒsiosakonnas.
RÀÀgin teile suurepĂ€rasest tööriistast ETL-protsesside arendamiseks â Apache Airflow. Kuid Airflow on nii universaalne ja mitmekesine, et tasub sellele tĂ€helepanu pöörata isegi siis, kui te ei tegele andmevoogudega, vaid vajate aeg-ajalt mĂ”nede protsesside kĂ€ivitamist ja nende teostamise jĂ€lgimist.
Ja jah, ma mitte ainult ei rÀÀgi, vaid ka nÀitan: programmis on palju koodi, ekraanipilte ja soovitusi.

Mida tavaliselt nÀed, kui googled sÔna Airflow / Wikimedia Commons
Sisukord
Sissejuhatus
Apache Airflow â ta on nagu Django:
- kirjutatud Pythonis,
- on suurepÀrane adminpaneel,
- piiramatu laiendatavus,
â ainult parem, ja loodud hoopis teiste eesmĂ€rkide saavutamiseks, nimelt (kui on kirjutatud enne katti):
- ĂŒlesannete kĂ€ivitamine ja jĂ€lgimine piiramatul arvul masinatel (nii palju kui teie lubab Celery/Kubernetes ja teie sĂŒdametunnistus)
- dĂŒnaamilise töövoo genereerimine vĂ€ga lihtsalt kirjutatava ja tajutava Python-koodi kaudu
- ja vĂ”imalus siduda omavahel ĂŒkskĂ”ik millised andmebaasid ja API-d nii valmis komponentide kui ka isetehtud pluginatega (mida on ÀÀrmiselt lihtne teha).
Me kasutame Apache Airflow'i nii:
- kogume andmeid erinevatest allikatest (palju SQL Serveri ja PostgreSQL instantsse, erinevaid API-sid rakenduste mÔÔdikute jaoks, isegi 1C) DWH-sse ja ODS-i (meie puhul on need Vertica ja Clickhouse).
- nagu arenenud
cron, mis kÀivitab andmete konsolideerimise protsesse ODS-is ja jÀlgib nende hooldust.
Kuni hiljutise ajani katab meie vajadusi ĂŒks vĂ€ike server 32 tuuma ja 50 GB RAM-iga. Airflow'is töötab sel juhul:
- ĂŒle 200 DAG-i (nende töövood, kuhu oleme ĂŒlesandeid tĂ€itnud),
- igaĂŒhes keskmiselt 70 ĂŒlesannet,
- selle headuse kÀivitamine (ka keskmiselt) korra tunni jooksul.
Ja sellest, kuidas me laienesime, kirjutan ma allpool, kuid praegu mÀÀratlegem ĂŒber-ĂŒlesanne, mida me lahendama hakkame:
On each of the three source SQL Servers, there are 50 databases â instances of one project, so their structure is identical (almost everywhere, muah-ha-ha), which means each has a table called Orders (after all, a table with that name can fit into any business). We pull data, adding metadata fields (source server, source database, ETL task identifier) and naively drop them into, say, Vertica.
LĂ€hme!
Peamine, praktiline (ja natuke teoreetiline) osa
Miks see meile (ja teile) vajalik on
When the trees were tall, and I was just a simple SQL-worker in a Russian retail company, we rolled out ETL processes aka data streams using the two tools available to us:
- Informatica Power Center â an extremely sprawling system, highly productive, with its own hardware, its own versioning. I used maybe 1% of its capabilities. Why? Well, first, this interface is from the noughties and mentally weighed down on us. Secondly, this thing is designed for incredibly complex processes, fierce component reuse, and other very-important-enterprise-features. Let's not even mention that it costs as much as the wing of an Airbus A380/year.
Caution, the screenshot may hurt people under 30 a bit.

- SQL Server Integration Server â we used this fellow in our internal project streams. Well, in fact: we already use SQL Server, and not using its ETL tools would be unreasonable. Everything about it is good: the interface is beautiful, and the execution reports⊠But thatâs not why we love software products, oh no. We can version it
dtsx(which is an XML with nodes mixed when saving); but what's the use? And to create a task package that will move a hundred tables from one server to another? Youâd drop your index finger on the mouse button before moving twenty. But it definitely looks more stylish:
We were definitely looking for ways out. It even almost came to a self-made SSIS package generatorâŠ
⊠and then I found a new job. And on that job, I encountered Apache Airflow.
When I learned that describing ETL processes is just simple Python code, I almost danced with joy. This is how data streams underwent versioning and diffusion, and pouring tables with a uniform structure from a hundred databases into one target became a task for Python code on one and a half to two 13" screens.
Kogume klastrit
Ărge korraldage lastaeda ja rÀÀkige tĂ€iesti ilmsetest asjadest, nagu Airflow'i, teie valitud andmebaasi, Celery ja muude asjade seadistamisest, millest dokumentatsioonis rÀÀgitakse.
Kuna me saame kohe katsetama hakata, olen ma koostanud docker-compose.yml kuna:
- KĂ€ivitame tegelikult Airflow: Scheduler, Webserver. Seal kĂ€ivitub ka Flower Celery ĂŒlesannete jĂ€lgimiseks (sest see on juba lĂŒkatud
apache/airflow:1.10.10-python3.7, ja me ei ole selle vastu); - PostgreSQL, kuhu Airflow kirjutab oma haldusandmed (planeerija andmed, tĂ€itmise statistika jne), ja Celery tĂ€histab lĂ”petatud ĂŒlesandeid;
- Redis, mis toimib Celery ĂŒlesannete vahendajana;
- Celery worker, mis hakkab ĂŒlesandeid tegelikult tĂ€itma.
- Kausta
. /dagssalvestame meie DAGide kirjeldusfailid. Need tuuakse automaatselt ĂŒles, seega pole vaja iga kord kogu steki tĂ€iesti ĂŒles tĂ”sta.
Kohati on kood nÀidetes osaliselt toodud (et mitte teksti ummistada) ja mÔned kohad muudetakse protsessi kÀigus. TÀielikke töötavaid koodinÀiteid saab vaadata hoidlas. .
docker-compose.yml
versioon: '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:
<<: *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
# Airflow Webserver
- 8080:8080
# Celery worker, will be scaled using `--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
- brokerMĂ€rkused:
- Ma kompositsiooni koostamisel toitusin peamiselt tuntud pildist â kindlasti vaadake. VĂ”ib-olla ei vaja te elus midagi muud.
- KÔik Airflow seaded on saadaval mitte ainult lÀbi
airflow.cfg, vaid ka keskkonnamuutujate kaudu (tÀnu arendajatele), mida ma hÀbivÀÀrselt kasutasin. - Muidugi ei ole see production ready: ma ei seadnud konteineritele heartbeate, ma ei vaevanud end turvalisusega. Kuid miinimum, mis sobib meie katsetele, on loodud.
- Pange tÀhele, et:
- DAGide kaust peab olema kergesti ligipÀÀsetav nii ajakava koostajale kui ka töötajatele.
- Sama kehtib kĂ”igi kolmandate osapoolte teekide kohta â need peavad kĂ”ik olema installitud ajakava koostaja ja töötajate masinatesse.
NĂŒĂŒd aga lihtsalt:
$ docker-compose up --scale worker=3Kui kĂ”ik on kĂ€ima lĂŒkatud, saab veebiliideseid vaatama minna:
- Airflow:
- Flower:
PÔhimÔisted
Kui te ei saanud nendest "dagidest" midagi aru, siis siin on lĂŒhike sĂ”nastik:
- Scheduler â peamine tegelane Airflow's, kes kontrollib, et robotid töötaksid, mitte inimene: jĂ€lgib ajakava, uuendab dag'e, kĂ€ivitab ĂŒlesandeid.
Tegelikult olid vanade versioonide puhul tal mĂ€luprobleemid (ei, mitte amneesia, vaid lekked) ja konfiguratsioonides jĂ€i isegi alles ĂŒks legassi parameeter
run_durationâ tema taaskĂ€ivitamise intervall. Aga nĂŒĂŒd on kĂ”ik korras. - DAG (ehk "dag") â "suunatud ahelate graaf", kuid see mÀÀratlemine ei ĂŒtle kellelegi suurt midagi, aga sisuliselt on see konteiner omavahel suhtlevatele ĂŒlesannetele (vt allpool) vĂ”i analoog Paketile SSIS-is ja Töökargile Informatica-s.
Lisaks dag'idele vÔivad olla ka subdag'id, aga tÔenÀoliselt me needeni ei jÔua.
- DAG Run â initsialiseeritud dag, millele on antud oma
execution_date. Ăhe dag'i dag'runnid vĂ”ivad tĂ€iesti töötada paralleelselt (kui olete muidugi teinud oma ĂŒlesanded idempotentseteks). - Operator â see on koodi tĂŒkk, mis vastutab mingi konkreetse toimingu tĂ€itmise eest. On kolm tĂŒĂŒpi operaatorit:
- tegevus, nagu nÀiteks meie lemmik
PythonOperator, mis on vĂ”imeline tĂ€itma mis tahes (kehtivat) Python-koodi; - ĂŒlekandmine, mis viivad andmeid ĂŒhelt poolt teisele, nĂ€iteks
MsSqlToHiveTransfer; - sensor lubab reageerida vĂ”i peatada edasise dag'i tĂ€itmise, kuni mingi sĂŒndmus ei ole toimunud.
HttpSensorvĂ”ib kĂŒsida mÀÀratud lĂ”pp-punkti, ja kui ootab vajaliku vastuse, kĂ€ivitab körvĂŒlekandeGoogleCloudStorageToS3Operator. Uudishimu kĂŒsib: "miks? Kas poleks lihtsam teha kordusi otse operaatoris!" Ja siis, et mitte koormata ĂŒlesannete kogumit seisatavate operaatoritega. Sensor kĂ€ivitub, kontrollib ja sureb kuni jĂ€rgmise katseni.
- tegevus, nagu nÀiteks meie lemmik
- Task â vĂ€ljakuulutatud operaatorid, olenemata tĂŒĂŒbist, ja dag'iga seotud ĂŒlesannete mÀÀramine tĂ”stetakse ĂŒlesande tasemele.
- Task instance â kui peaplaneerija otsustas, et ĂŒlesanded on valmis lĂ€hetamiseks tote jooksutavatele töötajatele (otse kohapeal, kui kasutame
LocalExecutorvÔi kaug-noodi puhulCeleryExecutor), mÀÀrab ta neile konteksti (st muutuja-sÀtte komplekti), rakendab kÀskude vÔi pÀringute malle ja kogub need kokku.
Genereerime ĂŒlesandeid
Esmalt mÀÀratleme meie dag'i ĂŒldise skeemi, seejĂ€rel sĂŒveneme jĂ€rjest enam detailidesse, sest rakendame mitmeid mitte triviaalsetes lahendustes.
Nii et kÔige lihtsamas vormis nÀeb selline dag vÀlja nii:
import 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)Hakkame siis ĂŒheskoos vĂ€lja selgitama:
- Esmalt impordime vajalikud teegid ja veel mÔned muud asjad;
sql_server_dsâ see onList[namedtuple[str, str]]Airflow Connections'i ĂŒhenduste nimed ja andmebaasid, millest me meie tabelit toome;dagon meie DAG'i deklaratsioon, mis peab olemaglobals(), muidu Airflow seda ei leia. DAG'ile tuleb ka öelda:- kuidas seda kutsuda
orderssee nimi ilmub hiljem veebiliideses, - et see hakkab töötama alates 8. juuli keskööst,
- ja kÀivitama ta peaks umbes iga 6 tunni jÀrel (karmimatele on siin 'timedelta()' asemel lubatud
-string,kergematele - vÀljend nagucron@daily0 0 0/6 ? * * *workflow()teeb peamise töö, aga mitte praegu. Praegu viskame lihtsalt meie konteksti logisse.);
- kuidas seda kutsuda
NĂŒĂŒd lihtne maagia ĂŒlesannete loomisel:jalutame meie allikate kaudu;- initsialiseerime
- , mis tÀidab meie anomaalia,
- . Ărge unustage mÀÀrata ĂŒlesande ainulaadne (DAG'i ulatuses) nimi ja seondada see DAG'iga. Lipp
PythonOperatorprovide_contextNĂŒĂŒd lihtne maagia ĂŒlesannete loomisel:toob omakorda funktsiooni lisaargumente, mille me hoolikalt kogume, kasutades**contextSellega on praegu kĂ”ik. Mida me saime:uus DAG veebiliideses,.
poolteist sada ĂŒlesannet, mis töötavad paralleelselt (kui Airflow, Celery ja serverite seadistused seda lubavad).
- Noh, peaaegu saime.
- Kes seab sÔltuvused paika?
Selle asja lihtsustamiseks lĂŒlitasin selle

kÀsitlemiseks
kĂ”ikidesse sĂ”lmedesse. docker-compose.yml NĂŒĂŒd siis lĂ€ks lahti: requirements.txt Hallid ruudud - ĂŒlesande instantsid, mida planeerija on töödelnud.
Ootame natuke, ĂŒlesanded haaravad töötajad:

Rohelised, nagu mÔistetav, - edukalt lÔpule viidud. Punased - mitte nii edukalt.
Muide, meie tootel ei ole mingit kausta,

, mis sĂŒsynkroniseerib masinate vahel - kĂ”ik DAG'id asuvad
meie Gitlab'is, ja Gitlab CI paigutab uuendused masinatesse, kui toimub liitmine
. /dagsKui töötajad töötlevad meie tĂŒhiseid ĂŒlesandeid, tuletame meelde veel ĂŒhte tööriista, mis meile midagi nĂ€idata vĂ”ib - Flower.gitEsimene leht, kus on kokkuvĂ”tvad andmed sĂ”lmedest - töötajatest:master.
Natuke Flowerist
KĂ”ige jaotatum leht ĂŒlesannetest, mis on tööle saadetud:
Esimene leht, mis sisaldab kokkuvÔtlikku teavet töötajate sÔlmpunktide kohta:

Rikkaim leht ĂŒlesannetega, mis on saadetud töötlemiseks:

KÔige igavam leht meie maakleri seisundi kohta:

KĂ”ige vĂ€rvilisem leht â ĂŒlesannete seisundi graafikute ja nende tĂ€itmise ajaga:

Laadime alla alalaadimise
Nii et kĂ”ik ĂŒlesanded on töötanud, saame haavatud Ă€ra viia.

Haavatuid oli ĂŒllatavalt palju â erinevatel pĂ”hjustel. Kui Airflowd kasutatakse Ă”igesti, siis need ruudud nĂ€itavad, et andmed ei ole kindlasti kohale jĂ”udnud.
Peame vaatama logisid ja taaskĂ€ivitama nurjunud ĂŒlesande nĂ€idiseid.
Klikkides igal ruudul, nÀeme meie kÀsutuses olevaid toiminguid:

Saame vĂ”tta ja teha nurjunule Clear. See tĂ€hendab, et unustame, et seal on mingi probleem ja sama ĂŒlesande nĂ€idis suundub planeerijale.

Selge, et ei ole vĂ€ga inimsĂ”bralik kĂ”iki punaseid ruute hiirega nii kĂ€sitleda â seda me Airflowlt ei oota. Loomulikult on meil massihĂ€vitusrelv: Browse/Task Instances

Valime kÔik korraga ja nullime, valides Ôige punkti:

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

Ăhendused, hookid ja muud muutujad
On paras 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("""Lugupeetud, 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, meil on {{ dag.dag_id }} kukkunud
"""),
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 kÔik on kunagi teinud raportite uuendamist? See on jÀlle see: on nimekiri allikatest, kust andmeid vÔtta; on nimekiri, kuhu need panna; Àrge unustage mÀrku anda, kui kÔik juhtus vÔi lÀks rikki (noh, see ei puuduta meid, ei).
Vaatame taas faili ja uurime uusi arusaamatuid asju:
from commons.operators import TelegramBotSendMessageâ ei ole midagi, mis takistaks meil oma operaatorite loomist, mida me tegime, luues vĂ€ikese mĂ€hise sĂ”numite saatmiseks Unblockitud. (Sellest operaatorist rÀÀgime ka hiljem);default_args={}â DAG vĂ”ib jagada samu argumendi kĂ”igile oma operaatoritele;to='{{ var.value.all_the_kings_men }}'â vĂ€litomeil ei ole seda kĂ”vasti mÀÀratletud, vaid see genereeritakse dĂŒnaamiliselt Jinja ja e-mailide nimekirja muutuja abil, mille ma hoolikalt asetasinAdmin/Variables;trigger_rule=TriggerRule.ALL_SUCCESSâ operaatori kĂ€ivitamise tingimus. Meie puhul saadetakse kiri juhtidele ainult siis, kui kĂ”ik sĂ”ltuvused töötasid edukalt;tg_bot_conn_id='tg_main'â argumendidconn_idvĂ”tavad endasse ĂŒhenduste identifikaatorid, mille me loomeAdmin/Connections;trigger_rule=TriggerRule.ONE_FAILEDâ Telegrami sĂ”numid saadetakse ainult, kui on katkestatud ĂŒlesandeid;task_concurrency=1â keelame sama ĂŒlesande mitme ĂŒlesande eksemplari samaaegse kĂ€ivitamise. Vastasel juhul saame mitu korragaVerticaOperator(mis vaatavad ĂŒhte tabelit);report_update >> [email, tg]â kĂ”ikVerticaOperatorkohtuvad kirja ja sĂ”numi saatmisel, just nii:

Kuid kuna teavitajate operaatoritel on erinevad kĂ€ivitamistingimused, töötab ainult ĂŒks. Puustruktuuris nĂ€eb see vĂ€lja vĂ€hem selge:

RÀÀgin paar sĂ”na makrode ja nende sĂ”prade â muutujad.
Makrod on Jinja-plekistrid, mis vÔivad esitada 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 }} lahti muudetakse konteksti muutuja sisuks execution_date Port 2222 YYYY-MM-DD: 2020-07-14. KĂ”ige parem on see, et konteksti muutujaid seotakse kindlalt konkreetse ĂŒlesande eksemplariga (ruudukesega Puustruktuuris), ja uuesti kĂ€ivitamisel avanevad plekistrid samadele vÀÀrtustele.
MÀÀratud vÀÀrtusi saab vaadata nuppuga Rendered igal ĂŒlesande eksemplaril. NĂ€iteks on see ĂŒlesandel, mis saadab kirja:

Ja see ĂŒlesandel, mis saadab sĂ”numi:

TĂ€ielik nimekiri sisseehitatud makrodest viimase saadaval oleva versiooni jaoks on siin:
Veelgi enam, pluginade abil saame kuulutada oma makrosid, kuid see on juba hoopis teine lugu.
Lisaks eelnevalt mÀÀratud asjadele saame kasutada oma muutujaid (muutsin neid ĂŒlaltoodud koodis). Loome Admin/Variables paar asja:

KĂŒll, nĂŒĂŒd saab hakata kasutama:
TelegramBotSendMessage(chat_id='{{ var.value.failures_chat }}')MÔisted vÔivad olla skalaarsed vÔi JSON. JSON-i puhul:
bot_config
{
"bot": {
"token": 881hskdfASDA16641,
"name": "Verter"
},
"service": "TG"
}kasutame lihtsalt vajalikku vÔtme teed: {{ var.json.bot_config.bot.token }}.
KĂ”nestan vaid ĂŒhe sĂ”na ja nĂ€itan ĂŒhte ekraanipilti seoses ĂŒhendustega. Siin on kĂ”ik elementaarne: lehe peal Admin/Connections loome ĂŒhenduse, paneme sinna oma kasutajanimed/paroolid ja veel spetsiifilisemad seaded. Nii:

Paroolid saab krĂŒpteerida (rohkem kui vaikimisi juhul) vĂ”i pole mÀÀratud ĂŒhenduse tĂŒĂŒpi (nagu tegin tg_main) â asi on selles, et tĂŒĂŒbid on Airflow mudelites sisse ehitatud ja ilma lĂ€htekoodidesse sekkumata neid ei saa muuta (kui midagi on vale, siis palun parandage mind), kuid krĂŒpteerida nime jĂ€rgi pole meil midagi takistust.
Ja lisaks vĂ”ime luua mitu sama nimega ĂŒhendust: sel juhul meetod BaseHook.get_connection(), mis toob meile ĂŒhendused nime jĂ€rgi, annab juhusliku ĂŒhe mitme sarnase seast (oleks loogilisem teha Round Robin, kuid jĂ€tame selle Airflow arendajate mureks).
Muutujad ja ĂŒhendused on kahtlemata suurepĂ€rased tööriistad, kuid oluline on mitte kaotada tasakaalu: millised teie voogude osad sĂ€ilitad konkreetselt koodis ning millised usaldad Airflowâd. Ăhelt poolt on mugav kiirelt muuta vÀÀrtust, nĂ€iteks postkasti, lĂ€bi UI. Teiselt poolt on see ikka naasmine hiireklikkide juurde, millest tahtsime (mina) lahti saada.
Ăhendustega töötamine on ĂŒks ĂŒlesanne hookide paralleelseks tĂ€itmiseks.. Ăldiselt on Airflow hook'id ĂŒhenduspunktid kolmandate teenuste ja raamatukogudega. NĂ€iteks JiraHook avab meile kliendi suhtlemiseks Jira'ga (saame ĂŒlesandeid liigutada edasi-tagasi), ja SambaHook vĂ”imaldab me pushida kohaliku faili smb-punkti.
TĂ€itame kohandatud operaatorit
Ja oleme lÀhenenud sellele, et vaadata, kuidas on tehtud TelegramBotSendMessage
Kood commons/operators.py oma operaatori kohta:
importida 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'Saadan "{self.message}" chati {self.chat_id}')
self.client.send_message(chat_id=self.chat_id,
message=self.message)Siin, nagu kÔik Airflow's, on kÔik vÀga lihtne:
- Oleme pÀrinud
BaseOperator, mis rakendab palju Airflow-spetsiifilisi asju (vaadake kunagi, kui on aega) - Oleme kuulutanud vÀlja vÀljad
template_fields, kus Jinja otsib makrosid töötlemiseks. - Oleme korraldanud Ôiged argumendid
__init__(), seadnud vaikimisi vÀÀrtused sinna, kus on vajalik. - Ka eelneva initsialiseerimist ei unustanud.
- Oleme avanud vastava hook'i
TelegramBotHook, saanud sellelt kliendi objekti. - Olemegi ĂŒle kirjutanud (override) meetodi
BaseOperator.execute(), mida Airflow kutsub vĂ€lja, kui on aeg operaatorit tööle panna â just seal me teeme pĂ”hitegevuse, unustamata logida. (Me logime, muide, otsestdoutjastderrâ Airflow teeb kĂ”ik pĂŒĂŒdmiseks ja pakib ilusasti kokku, paneb Ă”igesse kohta.)
Vaatame, mis meil on commons/hooks.py. Faili esimene osa, koos hook'iga:
importida 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 ei tea isegi, mida siit selgitada, lihtsalt mÀrgin olulised punktid:
- PĂ€rime, mĂ”tleme argumentidele â enamasti on neid ĂŒks:
conn_id; - Ăle kirjutame standardsed meetodid: ma piirdusin
get_conn(), kus ma saan ĂŒhenduse parameetrid nime jĂ€rgi ja lihtsalt tĂ”stan sektsiooni.extra(see code for JSON), where I put the Telegram bot token according to my own instructions:{"bot_token": "YOuRAwEsomeBOtToKen"}. - Creating an instance of our
TelegramBot, giving it a specific token.
That's it. You can get the client from the hook using TelegramBotHook().client vÔi TelegramBotHook().get_conn().
And the second part of the file, where I wrap the Telegram REST API to avoid carrying the same for just one method 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))The correct approach is to put all of this in:
TelegramBotSendMessage,TelegramBotHook,TelegramBotâ a plugin, store it in a public repository, and release it as Open Source.
While we were studying all this, our report updates successfully piled up and sent me a message about an error in the channel. I will go check what went wrong again...

Something broke in our DAG! Wasn't this what we were waiting for? Exactly!
Kas sa teed seda?
Do you feel like I've missed something? I promised to transfer data from SQL Server to Vertica, and here I went off-topic, how careless!
This misdeed was intentional; I simply had to explain some terminology to you. Now we can move on.
Our plan was as follows:
- Create a DAG
- Generate tasks
- See how everything looks nice
- Assign session numbers to uploads
- Fetch data from SQL Server
- Store data in Vertica
- Compile statistics
So, to run all this, I made a small addition to our 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 tÔstame:
- Vertica kui host
dwhkÔige vaikimisi seadistustega, - kolm SQL Serveri eksemplari,
- tÀiendame andmebaase viimasel ajal mÔne andmega (Àrge mingil juhul vaadake sisse
mssql_init.py!)
KÀivitame kogu selle hea kraami veidi keerukama kÀsuga kui eelmisel korral:
$ docker-compose -f docker-compose.yml -f docker-compose.db.yml up --scale worker=3Mis meie imetore genereerija tootis, on vÔimalik, kasutades punkti Data Profiling/Ad Hoc Query:

Peamine, et seda analĂŒĂŒtikutele ei nĂ€idata
Detailid ETL-seanssides ma ei peatuks, seal on kĂ”ik triviaalne: loome andmebaasi, sinna tabeli, katame kĂ”ik konteksti juhiga, ja nĂŒĂŒd teeme nii: with Session(task_name) as session: print('Load', session.id, 'started')# Load workflow ...session.successful = True session.loaded_rows = 15
session.pysession.py
from sys import stderr
class Session:
"""ETL töövoo seanss
NĂ€ide:
with Session(task_name) as session:
print(session.id)
session.successful = True
session.loaded_rows = 15
session.comment = 'Hea 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 = """
CREATE TABLE IF NOT EXISTS sessions (
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('Seanss 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,
', Laaditud: ', self.loaded_rows,
', kommentaar:', self.comment)
class SessionError(Exception):
pass
class SessionClosedError(SessionError):
passKÀtte on jÔudnud aeg andmed kokku koguda meie poolteise saja tabeli hulgast. Teeme seda vÀga lihtsalt ridadega:
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)- HĂŒĂŒdme abil saame Airflow'st
pymssql-ĂŒhenduse - KĂŒsitlusse lisame kuupĂ€eva piirangu - selle edastab meile mallitegija.
- KÀivitame meie pÀringu
pandas, mis tĂ”mbab meileDataFrameâ see tuleb meile hiljem kasuks.
Kasutame asendust
{dt}kĂŒsi parameetri asemel%smitte sellepĂ€rast, et ma olen kuri Buratino, vaid pigem sellepĂ€rast, etpandasei suuda hakkama saadapymssqlja paneb viimaseleparams: Loend, kuigi see vĂ€ga tahabtulp.
Pange tÀhele, et arendajapymssqlotsustas teda enam mitte toetada, ja on aeg kolidapyodbc.
Vaatame, millega Airflow tÀitis meie funktsioonide argumendid:

Kui andmeid ei olnud, siis pole mĂ”tet jĂ€tkata. Aga ka selle ĂŒle lugemine, et laadimine Ă”nnestus, on kummaline. Kuid see pole viga. A-a-a, mida teha?! Siin on, mida:
if df.empty:
raise AirflowSkipException('Ridasid ei ole laadimiseks')AirflowSkipException ĂŒtleb Airflow, et viga ei ole, ja ĂŒlesanne jÀÀb vahele. Liideses ei ole roheline ja ei punane ruut, vaid roosa.
Lisame meie andmetele mÔned veerud:
df['etl_source'] = src_schema
df['etl_id'] = session.id
df['hash_id'] = hash_pandas_object(df[['etl_source', 'id']])TĂ€pselt:
- DB, kust me tellimused vÔtsime,
- Meie laadimisessiooni identifikaator (see on erinev iga ĂŒlesande jaoks),
- Allika ja tellimuse identifikaatori hash â et lĂ”pp-andmebaasis (kus kĂ”ik valatakse ĂŒhte tabelisse) oleks meil ainulaadne tellimuse identifikaator.
JÀÀnud on eelviimane samm: laadida kĂ”ik Vertica'sse. Ja, nagu ei oleks uskumatu, on ĂŒks efektiivseid viise teha seda â lĂ€bi 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)- Teeme spetsiaalse vastuvÔtja
StringIO. pandasmis lahendab meieDataFramevĂ€lise lisandinaCSV-rid.- Avame ĂŒhenduse meie lemmik Vertica hook'iga.
- Ja nĂŒĂŒd saame
copy()saata meie andmed otse Vertica'sse!
Draiverist vĂ”tame, kui palju ridu laaditi, ja ĂŒtleme sessiooni juhile, et kĂ”ik on OK:
session.loaded_rows = cursor.rowcount
session.successful = TrueSee on kÔik.
Produksioonis loome sihttableti kÀsitsi. Siin lubasin endale vÀikese automaatika:
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 kasutan
VerticaOperator()loomiseks andmebaasi skeemi ja tabeli (kui need veel puuduvad, loomulikult). 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Ôtte
â No nii, â ĂŒtles hiirepoeg, â ei ole tĂ”si, et nĂŒĂŒd
Kas sa oled veendunud, et metsa kÔige hirmsam loom olen mina?
Julia Donaldson, âGruffaloâ
MÔtle, kui me kolleegidega korraldaksime konkursi: kes suudab kiiremini luua ja kÀivitada ETL-protsessi nullist: nemad oma SSIS-i ja hiirega ja mina Airflow'iga... Ja siis vÔrdleksime veel hooldamise mugavust... Oh, ma arvan, et oled nÔus, et ma möödun neist igas osas!
Kui rÀÀkida natuke tĂ”sisemalt, siis Apache Airflow â tĂ€nu protsesside kirjeldamisele programmikoodina â tegi mu töö mĂ€rksa mugavamaks ja meeldivamaks.
Selle piiramatud laienemisvĂ”imalused: nii pluginate osas kui ka skaleerimisvĂ”imes â annavad vĂ”imaluse rakendada Airflowâd praktiliselt igas valdkonnas: olgu see siis andmete kogumise, ettevalmistamise ja töötlemise tĂ€issĂŒkkel vĂ”i isegi rakettide kĂ€ivitamine (loomulikult Marsile).
Viimane, teaduslik-informatiivne osa
Kohad, mille me teie jaoks kokku korjasime
start_date. Jah, see on juba kohalik meem. Peamise DAG-argumenti 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 pole mingeid probleeme.
Sellega on seotud veel ĂŒks tĂ€itmisviga:
Task is missing the start_date parameter, mis enamikul juhtudest ĂŒtleb, et unustasid DAG-i operaatoriga siduda.- KĂ”ik ĂŒhel masinal. Jah, ja andmebaasid (nii Airflow enda kui ka meie ĂŒmbruse), ja veebiserver, ja ajastaja, ja töötajad. Ja see isegi toimis. Aga aja jooksul kasvas teenuste ĂŒlesannete arv, ja kui PostgreSQL hakkas vastama indeksilt 20 ms asemel 5 ms, vĂ”tsime selle ja viisime minema.
- LocalExecutor. Jah, me kasutame seda endiselt ja oleme juba ÀÀre peal. LocalExecutori jaoks on meil siiani piisavalt, kuid nĂŒĂŒd on aeg vĂ€hemalt ĂŒhe töötaja vĂ”rra laiendada, ja tuleb pingutada, et ĂŒle minna CeleryExecutori peale. Ja arvestades, et sellega saab töötada ka ĂŒhel masinal, ei takista miski Celery kasutamist isegi mitte serveris, mis "loomulikult ei tule kunagi tootmisse, tĂ”otame!"
- Kasutamata sisekasutuse:
- Connections teenuste autentimisandmete haldamiseks,
- SLA Misses ĂŒlesannete jaoks, mis ei töötanud Ă”igeaegselt,
- XCom metatekstide vahetamiseks (ma ĂŒtlesin metateavet!) DAG-i ĂŒlesannete vahel.
- Posti kuritarvitamine. Mis siin ikka öelda? Olin seadnud teavitused kĂ”ikide korduvate ebaĂ”nnestunud ĂŒlesannete jaoks. NĂŒĂŒd on minu töö Gmailis ĂŒle 90 000 e-kirja Airflow'lt ning postiteenuse veebivĂ€ljaanne keeldub kustutamast rohkem kui 100 korraga.
Rohkem peidetud kivisid:
Veel suurema automatiseerimise vahendid
Selleks, et saaksime veel rohkem mÔelda, mitte ainult kÀtega töötada, on Airflow meie jaoks ette valmistanud jÀrgmised vÔimalused:
- â tal on endiselt Eksperimentaalse staatuse, mis ei takista selle toimimist. Selle abil on vĂ”imalik mitte ainult saada teavet DAG-ide ja ĂŒlesannete kohta, vaid ka peatada/taasalustada DAG-i, luua DAG Run vĂ”i baas.
- â kĂ€surealt on saadaval palju tööriistu, mis ei ole lihtsalt ebamugavad veebiliidese kaudu kasutada, vaid pole seal isegi saadaval. NĂ€iteks:
backfillvajalik, et kĂ€ivitada ĂŒlesannete instantside uuesti kĂ€ivitamine.
NĂ€iteks tulid analĂŒĂŒtikud ja ĂŒtlesid: âTeie, seltsimees, andmetes on jama 1. kuni 13. jaanuarini! Parandage!â. Ja siis sa ĂŒtled: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 kĂ”rvale. Veelgi enam, seda saab kĂ€ivitada lĂ€biLocalExecutor, isegi kui sul on Celery klastri.- Umbes sama asja teeb
test, ainult et see ei kirjuta andmebaasi. connectionsvĂ”imaldab massiliselt luua ĂŒhendusi shellist.
- on ĂŒsna keeruline suhtlemisviis, mis on mĂ”eldud pistikprogrammidena, mitte selle kĂ€sitsi mudimiseks. Aga kes meid takistab minemast
/home/airflow/dags, kĂ€ivitamaipythonja hakkama siin mĂ€ngima? NĂ€iteks vĂ”ib kĂ”iki ĂŒhendusi eksportida jĂ€rgmise 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) - Ăhendus Airflow'i metadatuuri andmebaasi. Ma ei soovita sinna kirjutada, kuid ĂŒlesannete olekute saamine erinevate spetsiifiliste metrikate jaoks on oluliselt kiirem ja lihtsam kui lĂ€bi ĂŒhegi API.
Ătleme nii, et kaugeltki kĂ”ik meie ĂŒlesanded ei ole idempotentsed ning vĂ”ivad mĂ”nikord ebaĂ”nnestuda, ja see on normaalne. Aga mitu ebaĂ”nnestumist on juba kahtlane ja tuleks kontrollida.
Ole ettevaatlik, SQL!
VIIMASTE_EXECUTIONIDEGA AS ( VALI ĂŒlesande_id, dag_id, tĂ€itmise_aeg, olek, rida_numbrina() ĂLE ( JAOTUSEKS ĂŒlesande_id, dag_id KORRALDA tĂ€itmise_aeg LANGUS) NII rn FROM public.task_instance KUS tĂ€itmise_aeg > nĂŒĂŒd() - INTERVALL '2' PĂEVA ), ebaĂ”nnestunud AS ( VALI ĂŒlesande_id, dag_id, tĂ€itmise_aeg, olek, JUHTUM KUI rn = rida_numbrina() ĂLE ( JAOTUSEKS ĂŒlesande_id, dag_id KORRALDA tĂ€itmise_aeg LANGUS) SIIS TĂENE LĂPP JN kui ebaĂ”nnestunud KUS olek IN ('ebaĂ”nnestunud', 'ootab_kordamist') ) VALI ĂŒlesande_id, dag_id, loe(ebaĂ”nnestunud) AS ebaĂ”nnestunud, loe(JUHTUM KUI ebaĂ”nnestunud JA olek = 'ebaĂ”nnestunud' SIIS 1 LĂPP) AS ebaĂ”nnestunud, loe(JUHTUM KUI ebaĂ”nnestunud JA olek = 'ootab_kordamist' SIIS 1 LĂPP) AS ootab_kordamist FROM ebaĂ”nnestunud GRUPEERI ĂŒlesande_id, dag_id OLLES loe(ebaĂ”nnestunud) > 0
Viidatud lingid
Ja loomulikult esimesed kĂŒmme linki Google'i otsingust, mis viivad minu Airflow kaustadesse.
- â loomulikult peaks alustama ametlikust dokumentatsioonist, aga kes neid juhiseid ĂŒldiselt loeb?
- â no vĂ€hemalt loe loo autoreid nende soovitusi.
- â algus: kasutajaliides piltides
- â hĂ€sti kirjeldatud pĂ”hikontseptsioonid, kui sa (juhtumisi!) ei saanud aru, mis ma ĂŒtlen.
- â lĂŒhike juhend Airflow klastrite seadistamiseks.
- â peaaegu sama huvitav artikkel, kuid formaalsust on rohkem ja nĂ€iteid vĂ€hem.
- â töös koos Celeryga.
- â idempotentsete ĂŒlesannete, ID kaudu laadimise ning muude huvitavate asjade kohta.
- â ĂŒlesannete sĂ”ltuvused ja Trigger Rule, mida mainisin ainult möödaminnes.
- â kuidas ĂŒletada mĂ”ningaid 'töötab nagu kavandatud' plaanija puhul, laadida kaotatud andmeid ja seada ĂŒlesannetele prioriteete.
- â kasulikud SQL-pĂ€ringud Airflow metainfo kohta.
- â seal on kasulik jaotis kohandatud sensori loomise kohta.
- â huvitav lĂŒhike mĂ€rkus andeteaduse infrastruktuuri ĂŒlesehitamisest AWS-is.
- â levinud vead (kui keegi ikka ei loe juhiseid).
- â naeratage, kuidas inimesed paroolide hoidmise juures jalge alla saadavad, kuigi saate lihtsalt kasutada Ăhendusi.
- â varjatud DAG-i edastus, konteksti edastamine funktsioonis, taas sĂ”ltuvustest, ja ka ĂŒlesannete kĂ€ivitamise vahelejĂ€tmisest.
- â kasutamise kohta
vaikimisi argumentidejaparamsmalle ning ka muutujaid ja ĂŒhendusi. - â jutt sellest, kuidas plaani koostaja Airflow 2.0 jaoks ette valmistatakse.
- â veidi vananenud artikkel meie klastri kasutusele vĂ”tmisest
version: '1' services: simplesample-sonar: image: sonarqube:lts ports: - 9001:9000 - 9092:9092 network_mode: bridge. - â dĂŒnaamilised ĂŒlesanded mallide ja konteksti edastamise abil.
- â standardsed ja kohandatud teated e-posti ja Slacki kaudu.
- â Ălesannete lĂ”hestamine, makrod ja XCom.
Ja lingid, mis on artiklis kasutatud:
- â saadaval mallides kasutamiseks kohandatud kohad.
- â Levinud vead DAGide loomisel.
- â
version: '1' services: simplesample-sonar: image: sonarqube:lts ports: - 9001:9000 - 9092:9092 network_mode: bridgeeksperimentideks, silumiseks ja muuks. - â Pythoni pakett Telegrami REST API jaoks.
Allikas: habr.com




