Hallo, ich bin Dmitrij Logvinenko â Data Engineer in der Analytik-Abteilung der Unternehmensgruppe âVezetâ.
Ich möchte Ihnen ein fantastisches Werkzeug zur Entwicklung von ETL-Prozessen vorstellen â Apache Airflow. Doch Airflow ist so vielseitig und facettenreich, dass es sich auch fĂŒr Sie lohnt, einen Blick darauf zu werfen, selbst wenn Sie sich nicht mit Datenströmen beschĂ€ftigen, sondern gelegentlich Prozesse starten und deren AusfĂŒhrung ĂŒberwachen mĂŒssen.
Und ja, ich werde nicht nur darĂŒber sprechen, sondern auch zeigen: Es gibt viel Code, Screenshots und Empfehlungen im Programm.

Was gewöhnlich erscheint, wenn man nach dem Wort Airflow googelt / Wikimedia Commons
Inhaltsverzeichnis
EinfĂŒhrung
Apache Airflow â er ist genau wie Django:
- geschrieben in Python,
- mit einem hervorragenden Admin-Panel,
- ohne EinschrÀnkungen erweiterbar,
â aber besser, und er wurde ganz fĂŒr andere Zwecke entwickelt, nĂ€mlich (wie im Vorwort geschrieben):
- Start und Ăberwachung von Aufgaben auf einer unbegrenzten Anzahl von Maschinen (so viele wie es Celery/Kubernetes und Ihr Gewissen zulassen)
- mit dynamischer Workflow-Generierung aus leicht verstÀndlichem Python-Code
- und der Möglichkeit, beliebige Datenbanken und APIs mit sowohl vorgefertigten Komponenten als auch selbstgebauten Plugins zu verbinden (was Ă€uĂerst einfach ist).
Wir verwenden Apache Airflow wie folgt:
- wir sammeln Daten aus verschiedenen Quellen (z. B. mehrere Instanzen von SQL Server und PostgreSQL, verschiedene APIs mit Anwendungsmetriken, sogar 1C) in DWH und ODS (bei uns sind das Vertica und Clickhouse).
- als fortschrittliches
cron, das Datenkonsolidierungsprozesse im ODS ausfĂŒhrt und deren Wartung ĂŒberwacht.
Bis vor kurzem deckte unsere Infrastruktur einen kleinen Server mit 32 Kernen und 50 GB RAM ab. In Airflow laufen dabei:
- mehr als 200 DAGs (tatsÀchlich Workflows, in denen wir Aufgaben platziert haben),
- in jedem durchschnittlich 70 Aufgaben,
- dieses System wird ( ebenfalls im Durchschnitt) einmal pro Stunde ausgefĂŒhrt..
Wie wir uns erweitert haben, werde ich spĂ€ter schreiben, aber lassen Sie uns jetzt die ĂŒbergeordnete Aufgabe definieren, die wir lösen wollen:
Es gibt drei SQL Server, jeder mit 50 Datenbanken â Instanzen eines Projekts, sodass die Struktur nahezu identisch ist (fast ĂŒberall, muhahaha). Das bedeutet, dass jede eine Tabelle mit dem Namen âOrdersâ enthĂ€lt (zum GlĂŒck kann man eine solche Tabelle in jedes GeschĂ€ft einfĂŒgen). Wir extrahieren die Daten und fĂŒgen Feldinformationen hinzu (Quellserver, Quell-Datenbank, ETL-Task-ID) und werfen sie naiv in, sagen wir, Vertica.
Los geht's!
Der Hauptteil, praktisch (und ein wenig theoretisch)
Warum das fĂŒr uns (und fĂŒr Sie) wichtig ist
Als die BĂ€ume groĂ waren und ich ein einfacher SQL-Bearbeiter in einem russischen EinzelhĂ€ndler war, haben wir ETL-Prozesse aka Datenströme mit zwei verfĂŒgbaren Mitteln durchgefĂŒhrt:
- Informatica Power Center â ein Ă€uĂerst komplexes System, extrem leistungsstark, mit eigenen Hardware-Anforderungen und eigener Versionierung. Ich habe gerade einmal 1 % ihrer Möglichkeiten genutzt. Warum? Nun, erstens hat diese BenutzeroberflĂ€che, die irgendwo aus den Nuller Jahren stammt, psychisch auf uns gelastet. Zweitens ist dieses Ding fĂŒr extrem komplexe Prozesse, intensive Wiederverwendung von Komponenten und andere sehr wichtige Enterprise-Features ausgelegt. DarĂŒber, dass es so viel kostet wie ein FlĂŒgel eines Airbus A380 pro Jahr, lassen wir uns lieber nicht aus.
Achtung, ein Screenshot könnte jĂŒngeren Menschen unter 30 etwas schaden.

- SQL Server Integration Server â mit diesem Werkzeug haben wir in unseren internen Projektströmen gearbeitet. Aber mal im Ernst: SQL Server nutzen wir bereits, und es wĂ€re unklug, seine ETL-Tools nicht zu verwenden. Alles daran ist gut: die OberflĂ€che ist ansprechend, die AusfĂŒhrungsberichte⊠Doch dafĂŒr lieben wir Softwareprodukte nicht, oh nein, nicht dafĂŒr. Versionskontrolle vornehmen,
dtsx(das eine XML-Datei ist, in der beim Speichern die Knoten vermischt werden) können wir, aber was bringt das? Ein Paket von Aufgaben zu erstellen, das hundert Tabellen von einem Server auf einen anderen ĂŒbertrĂ€gt? Ach, was hundert â nach zwanzig Exemplaren wird Ihr Zeigefinger beim Klicken auf die Maustaste schmerzen. Aber es sieht zweifellos moderner aus:
Wir haben definitiv nach Auswegen gesucht. Es ging sogar so weit, fast einen selbstgeschriebenen SSIS-Paket-Generator zu erstellenâŠ
... und dann fand mich ein neuer Job. Und dort traf mich Apache Airflow.
Als ich erfuhr, dass die Beschreibung von ETL-Prozessen einfach Python-Code ist, hĂ€tte ich fast einen Freudentanz aufgefĂŒhrt. So wurden Datenströme versioniert und diffundiert, und das ZusammenfĂŒgen von Tabellen mit einheitlicher Struktur aus Hunderten von Datenbanken in ein Ziel wurde mit Python-Code auf einem 13-Zoll-Bildschirm zum Kinderspiel.
Cluster aufbauen
Lassen Sie uns nicht ins Kinderzimmer zurĂŒckfallen und ĂŒber offensichtliche Dinge sprechen, wie die Installation von Airflow, Ihrer Datenbank, Celery und anderen Aspekten, die in der Dokumentation beschrieben sind.
Damit wir sofort mit den Experimenten beginnen können, habe ich einige Ideen skizziert. docker-compose.yml in dem:
- Wir werden also Airflow: Scheduler, Webserver. Dort wird auch Flower zur Ăberwachung der Celery-Aufgaben laufen (da es bereits in
apache/airflow:1.10.10-python3.7, und wir haben nichts dagegen); - PostgreSQL, in den Airflow seine Betriebsinformationen (Scheduler-Daten, AusfĂŒhrungsstatistiken usw.) schreiben wird, wĂ€hrend Celery die abgeschlossenen Aufgaben vermerkt;
- Redis, das als Aufgabenbroker fĂŒr Celery dienen wird;
- Celery-Worker, der sich um die eigentliche AusfĂŒhrung der Aufgaben kĂŒmmert.
- In den Ordner
./dagswerden wir unsere Dateien mit der Beschreibung der DAGs ablegen. Diese werden spontan erfasst, sodass wir das gesamte System nach jeder Kleinigkeit nicht neu starten mĂŒssen.
An einigen Stellen ist der Code in den Beispielen nicht vollstĂ€ndig angegeben (um den Text nicht zu ĂŒberladen), und an anderen Stellen wird er wĂ€hrend des Prozesses modifiziert. VollstĂ€ndige, funktionierende Codebeispiele können im Repository angesehen werden. .
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 als Celery-Broker
broker:
image: redis:6.0.5-alpine
# DB fĂŒr die Airflow-Metadaten
airflow-db:
image: postgres:10.13-alpine
environment:
- POSTGRES_USER=airflow
- POSTGRES_PASSWORD=airflow
- POSTGRES_DB=airflow
volumes:
- ./db:/var/lib/postgresql/data
# Hauptcontainer mit 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, wird mit `--scale=n` skaliert
worker:
<<: *airflow-base
environment:
<
-c " sleep 10 &&
pip install --user -r /requirements.txt &&
/entrypoint worker"
depends_on:
- airflow
- airflow-db
- brokerHinweise:
- In der Komposition von Compose habe ich mich stark auf ein bekanntes Bild gestĂŒtzt. â schauen Sie unbedingt vorbei. Vielleicht benötigen Sie im Leben nichts anderes.
- Alle Einstellungen von Airflow sind nicht nur ĂŒber
airflow.cfg, sondern auch ĂŒber Umgebungsvariablen zugĂ€nglich (Dank an die Entwickler), was ich ausgiebig genutzt habe. - NatĂŒrlich ist es nicht produktionsbereit: Ich habe absichtlich keine Heartbeats auf die Container gesetzt und mich nicht um die Sicherheit gekĂŒmmert. Aber ich habe ein Minimum geschaffen, das fĂŒr unsere Experimente geeignet ist.
- Bitte beachten Sie, dass:
- Der DAG-Ordner muss sowohl fĂŒr den Scheduler als auch fĂŒr die Worker zugĂ€nglich sein.
- Das Gleiche gilt fĂŒr alle Drittanbieterbibliotheken â sie mĂŒssen auf den Maschinen mit dem Scheduler und den Workern installiert sein.
Nun, ganz einfach:
$ docker-compose up --scale worker=3Nachdem alles gestartet ist, können Sie sich die Web-OberflÀchen ansehen:
- Airflow:
- Flower:
Wichtige Konzepte
Wenn Sie von all diesen âDAGsâ nichts verstanden haben, hier ist ein kurzer Glossar:
- Scheduler â der wichtigste Akteur in Airflow, der sicherstellt, dass die Roboter arbeiten und nicht der Mensch: er ĂŒberwacht den Zeitplan, aktualisiert die DAGs und startet die Aufgaben.
FrĂŒher gab es in den Ă€lteren Versionen Probleme mit dem Speicher (nein, nicht Amnesie, sondern Lecks) und in den Konfigurationen blieb sogar ein Legacy-Parameter zurĂŒck.
run_durationâ das Intervall fĂŒr seinen Neustart. Aber jetzt lĂ€uft alles gut. - DAG (auch als âDAGâ bekannt) â âDirected Acyclic Graphâ. Diese Definition sagt den meisten Leuten wenig, aber letztlich handelt es sich um einen Container fĂŒr miteinander interagierende Tasks (siehe unten) oder um das Pendant zu Package in SSIS und Workflow in Informatica.
ZusÀtzlich zu DAGs kann es auch Sub-DAGs geben, aber wir werden wahrscheinlich nicht dazu kommen.
- DAG Run â ein initialisierter DAG, dem ein
execution_datezugewiesen ist. Die Runs eines DAGs können durchaus parallel laufen (sofern Sie Ihre Tasks natĂŒrlich idempotent gestaltet haben). - Operator â Teile des Codes, die fĂŒr die AusfĂŒhrung einer bestimmten Aktion verantwortlich sind. Es gibt drei Typen von Operatoren:
- Aktion, wie zum Beispiel unseren Lieblings-
PythonOperator, der in der Lage ist, jeden (gĂŒltigen) Python-Code auszufĂŒhren; - Ăberweisung, die Daten von einem Ort zum anderen transportieren, zum Beispiel,
MsSqlToHiveTransfer; - sensor , das Ihnen ermöglicht, auf Ereignisse zu reagieren oder die DurchfĂŒhrung des DAGs zu verlangsamen, bis ein bestimmtes Ereignis eintritt.
HttpSensorkann den angegebenen Endpunkt ansprechen, und wenn die benötigte Antwort kommt, den Transfer startenGoogleCloudStorageToS3Operator. Ein neugieriger Geist könnte fragen: âWarum? Man könnte die Wiederholungen doch direkt im Operator machen!â Um jedoch den Task-Pool nicht mit hĂ€ngenden Operatoren zu belasten. Der Sensor startet, ĂŒberprĂŒft und stirbt bis zum nĂ€chsten Versuch.
- Aktion, wie zum Beispiel unseren Lieblings-
- Task zurĂŒckgibt. â deklarierte Operatoren unabhĂ€ngig vom Typ und dem an den DAG angehĂ€ngten steigen auf den Rang eines Tasks.
- Task-Instanz â wenn der General-Planer entscheidet, dass die Tasks bereit sind, in den Kampf gegen die Worker-Executoren geschickt zu werden (sofort, wenn wir den
LocalExecutoroder auf einen entfernten Node im Falle desCeleryExecutor), weist er ihnen einen Kontext zu (d. h. ein Set von Variablen â AusfĂŒhrungsparametern), entfaltet die Vorlagen fĂŒr Befehle oder Anfragen und lagert diese im Pool.
Tasks generieren
ZunÀchst skizzieren wir das Gesamtkonzept unseres DAGs, und dann werden wir immer tiefer in die Details eintauchen, da wir einige nicht triviale Lösungen anwenden.
Also, in seiner einfachsten Form wĂŒrde ein solcher DAG so aussehen:
from datetime 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)Lassen Sie uns klÀren:
- ZunÀchst importieren wir die notwendigen Bibliotheken und einige weitere Dinge;
sql_server_dsâ istList[namedtuple[str, str]]mit den Namen der Verbindungen aus Airflow Connections und den Datenbanken, aus denen wir unsere Tabelle abrufen werden;dagâ die Deklaration unseres DAGs, die unbedingt inglobals(), sonst findet Airflow ihn nicht. Dem DAG muss auch gesagt werden:- wie er heiĂt
ordersâ dieser Name wird spĂ€ter im Web-Interface angezeigt, - dass er um Mitternacht am achten Juli gestartet wird,
- und dass er etwa alle 6 Stunden ausgefĂŒhrt werden sollte (fĂŒr die Coolen hier kann anstelle von
timedelta()auchcron-String0 0 0/6 ? * * *verwendet werden, fĂŒr die weniger Coolen â ein Ausdruck wie@daily);
- wie er heiĂt
workflow()wird die Hauptarbeit erledigen, aber nicht jetzt. Momentan geben wir einfach unseren Kontext ins Log aus.- Und jetzt die einfache Magie der Aufgaben Erstellung:
- wir durchlaufen unsere Quellen;
- initialisieren
PythonOperator, die unser Platzhalter ausfĂŒhren wirdworkflow(). Vergessen Sie nicht, einen einzigartigen (innerhalb des DAGs) Task-Namen anzugeben und den DAG selbst zu verknĂŒpfen. Flagprovide_contextwird in der Folge in die Funktion zusĂ€tzlicher Argumente hineingespeist, die wir sorgfĂ€ltig mit Hilfe von**context.
Das war's vorerst. Was haben wir erhalten:
- ein neuer DAG im Web-Interface,
- eineinhalb Hundert Tasks, die parallel ausgefĂŒhrt werden (sofern es die Konfigurationen von Airflow, Celery und die ServerkapazitĂ€ten zulassen).
Nun, fast haben wir es erhalten.

Wer kĂŒmmert sich um die AbhĂ€ngigkeiten?
Um all das zu vereinfachen, habe ich in die docker-compose.yml Verarbeitung requirements.txt auf allen Knoten integriert.
Jetzt geht's los:

Graue KĂ€stchen â Task-Instanzen, die vom Planer verarbeitet wurden.
Wir warten einen Moment, wĂ€hrend die Worker die Aufgaben ĂŒbernehmen:

Die grĂŒnen sind, wie zu erwarten, â erfolgreich abgeschlossen. Die roten â nicht so erfolgreich.
Ăbrigens, auf unserem Produktivsystem gibt es keinen Ordner,
./dagsder zwischen den Maschinen synchronisiert â alle DAGs liegen aufgitunserem GitLab, und GitLab CI verteilt Updates auf die Maschinen beim Merge inmaster.
Ein bisschen ĂŒber Flower
WĂ€hrend die Worker an unseren Platzhalter-Tasks arbeiten, erinnern wir uns an ein weiteres Tool, das uns einige Informationen zeigen kann â Flower.
Die erste Seite mit zusammenfassenden Informationen zu den Worker-Knoten:

Die am stĂ€rksten gefĂŒllte Seite mit den Aufgaben, die in Bearbeitung sind:

Die langweiligste Seite mit dem Status unseres Brokers:

Die auffĂ€lligste Seite â mit den Diagrammen zum Status der Aufgaben und deren AusfĂŒhrungszeiten:

Noch nicht vollgeladene Elemente nachladen
Also, alle Aufgaben wurden bearbeitet, wir können die Verletzten abtransportieren.

Und es gab viele Verletzte â aus verschiedenen GrĂŒnden. Bei korrekter Nutzung von Airflow zeigen diese Quadrate, dass die Daten definitiv nicht angekommen sind.
Es muss das Protokoll angesehen und die abgestĂŒrzten Task-Instanzen neu gestartet werden.
Durch Klicken auf eines der Quadrate sehen wir die verfĂŒgbaren Aktionen:

Wir können einfach die Clear-Option fĂŒr die abgestĂŒrzte Instanz auswĂ€hlen. Das heiĂt, wir vergessen, dass etwas schief gelaufen ist, und dieselbe Task-Instanz geht zurĂŒck zum Planer.

Es ist klar, dass es nicht sehr human ist, mit der Maus durch alle roten Quadrate zu fahren â das erwarten wir nicht von Airflow. NatĂŒrlich haben wir eine Massenzerstörungswaffe: Durchsuchen/Task-Instanzen

WĂ€hlen wir alles auf einmal aus und setzen wir den richtigen Punkt zurĂŒck:

Nach der RĂŒcksetzung sehen unsere Tasks so aus (sie können es kaum erwarten, bis der Scheduler sie plant):

Verbindungen, Hooks und andere Variablen
Es ist an der Zeit, den nÀchsten DAG zu betrachten, 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-Berichte aktualisiert',
html_content=dedent("""Sehr geehrte Damen und Herren, die Berichte wurden aktualisiert"""),
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("""
ĐаŃаŃа, wach auf! Wir haben {{ dag.dag_id }} verloren
"""),
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]Hat nicht jeder irgendwann schon einmal Berichte aktualisiert? Hier sind wir wieder: Es gibt eine Liste von Quellen, von denen wir Daten abrufen; es gibt eine Liste, wo wir sie ablegen; und wir vergessen nicht, ein Signal zu senden, wenn alles geklappt hat oder etwas kaputtgegangen ist (na ja, das gilt nicht fĂŒr uns, oder?).
Lass uns noch einmal die Datei durchgehen und uns die neuen, unklaren Dinge anschauen:
from commons.operators import TelegramBotSendMessageâ nichts hindert uns daran, eigene Operatoren zu erstellen, was wir auch getan haben, indem wir eine kleine HĂŒlle fĂŒr das Senden von Nachrichten in Entsperrt gebaut haben. (Ăber diesen Operator werden wir uns weiter unten unterhalten);default_args={}â der DAG kann denselben Argumenten an alle seine Operatoren verteilen;to='{{ var.value.all_the_kings_men }}'â das Feldzuwird bei uns nicht hartkodiert, sondern dynamisch mit Hilfe von Jinja und einer Variablen mit einer Liste von E-Mails gebildet, die ich sorgfĂ€ltig inAdmin/Variables;trigger_rule=TriggerRule.ALL_SUCCESSâ die Bedingung zum Starten des Operators. In unserem Fall wird die Nachricht an die Chefs nur gesendet, wenn alle AbhĂ€ngigkeiten erfolgreich abgeschlossen wurden. erfolgreich.;tg_bot_conn_id='tg_main'â die Argumenteconn_idnehmen die Verbindungsidentifikatoren an, die wir inAdmin/Connections;trigger_rule=TriggerRule.ONE_FAILEDâ Nachrichten in Telegram werden nur gesendet, wenn Aufgaben fehlgeschlagen sind;task_concurrency=1â wir verbieten den gleichzeitigen Start mehrerer Task-Instanzen desselben Tasks. Andernfalls erhalten wir den gleichzeitigen Start mehrererVerticaOperator(die auf dieselbe Tabelle zugreifen);report_update >> [email, tg]â alleVerticaOperatorwerden in der Versendung von E-Mails und Nachrichten zusammenkommen, so:

Da jedoch bei den Benachrichtigungsoperators unterschiedliche Startbedingungen gelten, wird nur einer arbeiten. Im Baumdiagramm sieht alles etwas weniger ĂŒbersichtlich aus:

Ich möchte ein paar Worte ĂŒber Makros und ihre Freunde â Variablen.
Makros sind Jinja-Platzhalter, die verschiedene nĂŒtzliche Informationen in die Argumente der Operatoren einfĂŒgen können. Zum Beispiel so:
SELECT
id,
payment_dtm,
payment_type,
client_id
FROM orders.payments
WHERE
payment_dtm::DATE = '{{ ds }}'::DATE{{ ds }} wird in den Inhalt der Kontextvariable erweitert execution_date im Format YYYY-MM-DD: 2020-07-14. Das Angenehme ist, dass Kontextvariablen an eine bestimmte Task-Instanz (das Quadrat im Baumdiagramm) gebunden sind, und beim Neustart werden die Platzhalter in dieselben Werte aufgelöst.
Die zugewiesenen Werte können mit der SchaltflÀche Rendered bei jeder Task-Instanz angezeigt werden. So sieht es bei der Task zum Versenden von E-Mails aus:

Und so bei der Task zum Versenden von Nachrichten:

Die vollstĂ€ndige Liste der integrierten Makros fĂŒr die letzte verfĂŒgbare Version finden Sie hier:
DarĂŒber hinaus können wir mit Plugins eigene Makros deklarieren, aber das ist eine ganz andere Geschichte.
Neben den vordefinierten Elementen können wir auch Werte unserer Variablen einsetzen (wie ich oben im Code bereits gezeigt habe). Lassen Sie uns in Admin/Variables ein paar Elemente erstellen:

Alles klar, Sie können verwenden:
TelegramBotSendMessage(chat_id='{{ var.value.failures_chat }}')Im Wert kann ein Skalar sein, oder es kann auch JSON enthalten. Im Fall von JSON:
bot_config
{
"bot": {
"token": 881hskdfASDA16641,
"name": "Verter"
},
"service": "TG"
}verwenden Sie einfach den Pfad zum gewĂŒnschten SchlĂŒssel: {{ var.json.bot_config.bot.token }}.
Ich werde nur ein Wort sagen und einen Screenshot ĂŒber Verbindungenzeigen. Hier ist alles einfach: Auf der Seite Admin/Connections erstellen wir eine Verbindung, fĂŒgen unsere Logins/Passwörter und spezifischere Parameter hinzu. So geht's:

Passwörter können verschlĂŒsselt werden (grĂŒndlicher als in der Standardversion), oder Sie können den Verbindungstyp weglassen (wie ich es fĂŒr tg_main gemacht habe)) â das Problem ist, dass die Liste der Typen in den Airflow-Modellen fest verankert ist und ohne Eingriff in den Quellcode nicht verĂ€ndert werden kann (falls ich etwas nicht gefunden habe, bitte korrigiert mich), aber es steht uns nichts im Wege, die Zugangsdaten einfach nach Namen zu erhalten.
AuĂerdem kann man mehrere Verbindungen mit demselben Namen erstellen: In diesem Fall wird die Methode BaseHook.get_connection(), die uns Verbindungen nach Namen bereitstellt, eine zufĂ€llige Auswahl aus mehreren Namensvettern zurĂŒckgeben (es wĂ€re logischer, Round Robin zu implementieren, aber das ĂŒberlasse ich den Entwicklern von Airflow).
Variablen und Verbindungen sind ohne Zweifel groĂartige Werkzeuge, aber es ist wichtig, die Balance nicht zu verlieren: welche Teile Ihrer Workflows speichern Sie im Code und welche ĂŒberlassen Sie Airflow. Auf der einen Seite kann es praktisch sein, Werte wie eine E-Mail-Adresse schnell ĂŒber das UI zu Ă€ndern. Auf der anderen Seite ist das doch ein RĂŒckschritt zum Klick-Klick, den wir (ich) vermeiden wollten.
Die Arbeit mit Verbindungen ist eine der Aufgaben Hooks. Generell sind Hooks in Airflow Verbindungspunkte zu externen Diensten und Bibliotheken. Zum Beispiel wird JiraHook uns einen Client fĂŒr die Interaktion mit Jira bereitstellen (wir können Aufgaben hin und her bewegen), und mit SambaHook kann man eine lokale Datei auf smb hochladen.-Punkt.
Einen benutzerdefinierten Operator analysieren
Und wir sind ganz nah daran, einen Blick darauf zu werfen, wie TelegramBotSendMessage
Code commons/operators.py mit dem eigentlichen Operator:
from typing import Union
from airflow.operators import BaseOperator
from commons.hooks import TelegramBotHook, TelegramBot
class TelegramBotSendMessage(BaseOperator):
"""Sendet eine Nachricht an chat_id ĂŒber TelegramBotHook
Beispiel:
>>> 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 }} fehlgeschlagen :(',
... 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'Sende "{self.message}" an den Chat {self.chat_id}')
self.client.send_message(chat_id=self.chat_id,
message=self.message)Hier, wie alles andere in Airflow, ist alles sehr einfach:
- Abgeleitet von
BaseOperator, der viele Airflow-spezifische Dinge realisiert (sehen Sie sich das in Ruhe an) - Wir haben die Felder
template_fields, in denen Jinja nach Makros zur Verarbeitung suchen wird, deklariert. - Wir haben die richtigen Argumente fĂŒr
__init__() organisiert., haben die Standardwerte dort gesetzt, wo es nötig ist. - Die Initialisierung des Elternteils haben wir ebenfalls nicht vergessen.
- Wir haben den entsprechenden Hook geöffnet
TelegramBotHook, und das Objekt-Client von ihm erhalten. - Wir haben die Methode
BaseOperator.execute()ĂŒberschrieben, die Airflow aufrufen wird, wenn es an der Zeit ist, den Operator auszufĂŒhren â darin implementieren wir die Hauptaktion, wobei wir nicht vergessen, uns zu authentifizieren. (Ăbrigens protokollieren wir direkt instdoutundstderrâ Airflow wird alles abfangen, schön verpacken und alles an den richtigen Ort sortieren.)
Lassen Sie uns ansehen, was wir in commons/hooks.pyhaben. Der erste Teil der Datei, mit dem Hook selbst:
from typing import Union
from airflow.hooks.base_hook import BaseHook
from requests_toolbelt.sessions import BaseUrlSession
class TelegramBotHook(BaseHook):
"""Telegram Bot API Hook
Hinweis: FĂŒgen Sie eine Verbindung mit leeren Conn Type hinzu und vergessen Sie nicht,
Extra auszufĂŒllen:
{"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.clientIch weià nicht einmal, was ich hier erklÀren soll, ich möchte nur auf wichtige Punkte hinweisen:
- Wir erben, denken ĂŒber die Argumente nach â in den meisten FĂ€llen wird es nur eines sein:
conn_id; - Ăbliche Methoden neu definieren: ich habe mich eingeschrĂ€nkt
get_conn(), in dem ich die Verbindungsparameter nach Name abfrage und lediglich den Abschnitt herausholeextra(dieses Feld fĂŒr JSON), in das ich (nach meiner eigenen Anweisung!) das Token des Telegram-Bots gelegt habe:{"bot_token": "YOuRAwEsomeBOtToKen"}. - Ich erstelle eine Instanz unseres
TelegramBot, indem ich ihm das spezifische Token ĂŒbergebe.
Das ist alles. Den Client aus dem Hook erhÀlt man mit TelegramBotHook().clent oder TelegramBotHook().get_conn().
Und der zweite Teil der Datei, in dem ich eine Mikroverpackung fĂŒr die Telegram REST API mache, um nicht die gleiche nur fĂŒr eine Methode sendMessage.
class TelegramBot:
"""Telegram Bot API Wrapper
Beispiele:
>>> TelegramBot('YOuRAwEsomeBOtToKen', '@myprettydebugchat').send_message('Hallo, Schatz')
>>> TelegramBot('YOuRAwEsomeBOtToKen').send_message('Hallo, Schatz', 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))Der richtige Weg ist, all dies zu kombinieren:
TelegramBotSendMessage,TelegramBotHook,TelegramBotâ in das Plugin, in ein öffentliches Repository zu legen und als Open Source zur VerfĂŒgung zu stellen.
WĂ€hrend wir das alles untersucht haben, sind unsere Berichtsaktualisierungen erfolgreich gescheitert und haben mir eine Fehlermeldung in den Kanal gesendet. Ich werde schauen, was diesmal nicht stimmt...

In unserem DAG ist etwas kaputt! War das nicht das, worauf wir gewartet haben? Genau!
Wirst du auch einschenken?
Habe ich etwas ĂŒbersehen? Ich hatte versprochen, Daten von SQL Server nach Vertica zu migrieren, und habe mich dann einfach vom Thema abgewandt, wie ungezogen!
Die Tat war absichtlich, ich musste Ihnen unbedingt ein paar Begriffe erklÀren. Jetzt können wir weitermachen.
Unser Plan war folgender:
- Einen DAG erstellen
- Tasks generieren
- ĂberprĂŒfen, wie schön alles aussieht
- Sessions fĂŒr die LadevorgĂ€nge zuweisen
- Daten aus SQL Server abrufen
- Daten in Vertica speichern
- Statistiken sammeln
Um all dies zu starten, habe ich eine kleine ErgÀnzung zu unserem 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.pyDort heben wir an:
- Vertica als Host
dwhmit den Standardkonfigurationen, - drei SQL Server Instanzen,
- die Datenbanken mit einigen letzten Daten fĂŒllen (bitte schauen Sie auf keinen Fall in die
mssql_init.py!)
Wir starten alles mit einem etwas komplizierteren Befehl als beim letzten Mal:
$ docker-compose -f docker-compose.yml -f docker-compose.db.yml up --scale worker=3Was unser Wunder-Randomizer generiert hat, kann man nutzen, indem man Punkt Datenprofilierung/Ad-hoc-Abfragen:

Wichtig ist, das nicht den Analysten zu zeigen
Detailliert darauf eingehen ETL-Sitzungen werde ich nicht, da alles trivial ist: wir erstellen eine Datenbank, darin eine Tabelle, umhĂŒllen alles mit einem Kontextmanager, und nun machen wir Folgendes:
with Session(task_name) as session:
print('Laden', session.id, 'gestartet')
# Workflow laden
...
session.successful = True
session.loaded_rows = 15session.py
from sys import stderr
class Session:
"""ETL-Workflow-Sitzung
Beispiel:
mit Session(task_name) als session:
print(session.id)
session.successful = True
session.loaded_rows = 15
session.comment = 'Gut gemacht'
"""
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, 'geöffnet')
return self
def close(self):
if not self._id:
raise SessionClosedError('Sitzung ist nicht geöffnet')
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, 'geschlossen',
', erfolgreich: ', self.successful,
', Geladen: ', self.loaded_rows,
', Kommentar:', self.comment)
class SessionError(Exception):
pass
class SessionClosedError(SessionError):
passEs ist an der Zeit unsere Daten abzuholen aus unseren mehr als hundert Tabellen. Das erledigen wir mit ein paar einfachen Zeilen:
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)- Mit dem Hook ziehen wir aus Airflow
pymssql-Verbindung - Im Anfrage fĂŒgen wir das Datum als Bedingung hinzu â der Template-Engine wird das ĂŒbergeben.
- Wir fĂŒttern unsere Anfrage
pandas, die uns die Daten abrufen wirdDataFrameâ sie wird uns spĂ€ter nĂŒtzlich sein.
Ich benutze die Ersetzung
{dt}anstelle des Anfrageparameters%snicht, weil ich ein böser Buratino bin, sondern weilpandaser nicht mitpymssqlzurechtkommt und dem Letzterenparams: List, obwohl er es sehr gerne möchte.Tupel.
Beachten Sie auch, dass der Entwicklerpymssqlbeschlossen hat, ihn nicht lĂ€nger zu unterstĂŒtzen, und es Zeit ist, zu wechseln zupyodbc.
Sehen wir uns an, womit Airflow die Argumente unserer Funktionen gefĂŒllt hat:

Wenn keine Daten vorhanden sind, macht es keinen Sinn, fortzufahren. Aber es wĂ€re auch merkwĂŒrdig, die EinfĂŒgung als erfolgreich zu betrachten. Aber das ist kein Fehler. A-a-a, was nun?! So:
if df.empty:
raise AirflowSkipException('Keine Zeilen zum Laden')AirflowSkipException Airflow wird sagen, dass es keinen Fehler gibt, und wir den Task ĂŒberspringen. Im Interface wird es kein grĂŒnes oder rotes Feld geben, sondern in Pink.
Wir fĂŒgen unseren Daten einige Spalten hinzu:
df['etl_source'] = src_schema
df['etl_id'] = session.id
df['hash_id'] = hash_pandas_object(df[['etl_source', 'id']])angegeben. Genauer gesagt:
- Die Datenbank, aus der wir die Bestellungen entnommen haben,
- Die ID unserer Lade-Session (diese wird unterschiedlich sein fĂŒr jeden Task),
- Der Hash vom Quellen- und Bestell-Identifikator â damit wir in der finalen Datenbank (wo alles in einer Tabelle zusammengefĂŒhrt wird) einen einzigartigen Bestell-Identifier haben.
Es bleibt nur noch ein Schritt: Alles in Vertica laden. Und tatsĂ€chlich ist einer der effektivsten Wege, dies zu tun â ĂŒber 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)- Wir erstellen einen speziellen EmpfÀnger
StringIO. pandaswird freundlich unsereDataFramein Form vonCSV-Zeilen speichern.- Wir öffnen eine Verbindung zu unserem bevorzugten Vertica ĂŒber einen Hook.
- Und jetzt werden wir mit
copy()unsere Daten direkt an Vertica senden!
Wir holen von dem Treiber, wie viele Zeilen geladen wurden, und sagen dem Session-Manager, dass alles in Ordnung ist:
session.loaded_rows = cursor.rowcount
session.successful = TrueDas war's.
In der Produktion erstellen wir die Zieltabelle manuell. Hier habe ich mir eine kleine Automatisierung erlaubt:
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)Ich bin mit Hilfe von
VerticaOperator()erstelle ich das DB-Schema und die Tabelle (sofern sie noch nicht existieren, natĂŒrlich). Wichtig ist, die AbhĂ€ngigkeiten richtig zu setzen:
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 >> loadZusammenfassung
â Nun, â sagte das MĂ€uschen, â nicht wahr, jetzt
Bist du ĂŒberzeugt, dass ich im Wald das furchterregendste Tier bin?
Julia Donaldson, âDer GrĂŒffeloâ
Ich denke, wenn meine Kollegen und ich einen Wettbewerb veranstalten wĂŒrden, wer am schnellsten einen ETL-Prozess von Grund auf entwickelt: sie mit ihren SSIS und mir mit Airflow⊠Und danach wĂŒrden wir auĂerdem die Wartungsfreundlichkeit vergleichen⊠Oh, ich denke, ihr werdet zustimmen, dass ich sie in jeder Hinsicht ĂŒbertreffen werde!
Wenn es etwas ernster wird, hat Apache Airflow â durch die Beschreibung von Prozessen in Form von Programmcode â meine Arbeit deutlich bequemer und angenehmer gestaltet.
Seine unbegrenzte Erweiterbarkeit â sowohl bei Plugins als auch hinsichtlich der Skalierbarkeit â ermöglicht es Ihnen, Airflow praktisch in jedem Bereich anzuwenden: sei es im vollstĂ€ndigen Zyklus der Datensammlung, -aufbereitung und -verarbeitung oder beim Start von Raketen (natĂŒrlich fĂŒr den Mars).
Der abschlieĂende, informative Teil
Die Stolpersteine, die wir fĂŒr Sie gesammelt haben
start_date. Ja, das ist schon ein lokales Meme. Durch das Hauptargument des DAGsstart_datewerden alle durchlaufen. Kurz gesagt, wenn Sie in derstart_dateaktuellen Datumsangabe angeben und imschedule_intervalâ einen Tag, wird der DAG morgen frĂŒhestens gestartet.start_date = datetime(2020, 7, 7, 0, 1, 2)Und keine weiteren Probleme.
Damit ist auch ein weiterer AusfĂŒhrungsfehler verbunden:
Task is missing the start_date parameter, der meist bedeutet, dass Sie vergessen haben, den DAG-Operator zuzuordnen.- Alles auf einem einzigen GerĂ€t. Ja, sowohl die Datenbanken (von Airflow selbst und unserer Schicht), der Webserver, der Scheduler und die Worker. Und es hat sogar funktioniert. Aber im Laufe der Zeit wuchs die Anzahl der Aufgaben in den Services, und als PostgreSQL begann, Antworten ĂŒber den Index in 20 ms statt in 5 ms zu geben, haben wir ihn einfach entfernt.
- LocalExecutor. Ja, wir nutzen es immer noch und stehen nun am Rand des Abgrunds. Der LocalExecutor hat uns bisher ausgereicht, aber jetzt ist es an der Zeit, mindestens einen Worker hinzuzufĂŒgen, und wir mĂŒssen uns anstrengen, um auf den CeleryExecutor umzusteigen. Da man mit diesem auch auf einer Maschine arbeiten kann, gibt es nichts, was uns davon abhĂ€lt, Celery selbst auf einem Server zu verwenden, der ânatĂŒrlich niemals in die Produktion gehen wird, das schwöre ich!â
- Nicht verwendet integrierte Mittel:
- Verbindungen zum Speichern von Service-Anmeldeinformationen,
- SLA-VerstöĂe zum Handeln von Aufgaben, die nicht rechtzeitig verarbeitet wurden,
- XCom zum Austausch von Metadaten (ich habe âMetaâ gesagt Daten!) zwischen den Aufgaben des DAGs.Missbrauch von E-Mails. Was soll man dazu sagen? Es wurden Benachrichtigungen fĂŒr alle Wiederholungen fehlgeschlagener Aufgaben eingerichtet. Jetzt habe ich ĂŒber 90.000 Emails von Airflow in meinem Arbeits-Gmail, und die WeboberflĂ€che des E-Mail-Clients weigert sich, mehr als 100 StĂŒck auf einmal zu löschen.
- Weitere Fallstricke:
Apache Airflow Fallstricke
Mittel zur weiteren Automatisierung
Um sicherzustellen, dass wir noch mehr mit dem Kopf und nicht mit den HĂ€nden arbeiten, hat Airflow Folgendes fĂŒr uns vorbereitet:
- â es hat immer noch den Status Experimental, was seiner FunktionalitĂ€t jedoch nicht im Wege steht. Damit können Sie nicht nur Informationen zu DAGs und Tasks abrufen, sondern auch DAGs anhalten oder starten, einen DAG Run erstellen oder einen Pool verwalten.
- â ĂŒber die Kommandozeile stehen viele Werkzeuge zur VerfĂŒgung, die entweder umstĂ€ndlich ĂŒber die WebUI zu bedienen sind oder dort gar nicht vorhanden sind. Beispielsweise:
backfillwird benötigt, um Task-Instanzen erneut zu starten.
Sagen wir, Analysten kommen und sagen: âIhre Daten vom 1. bis 13. Januar sind nicht in Ordnung! Reparieren Sie das!â Und Sie denken sich:airflow backfill -s '2020-01-01' -e '2020-01-13' orders- Datenbankpflege:
initdb,resetdb,upgradedb,checkdb. run, welches es ermöglicht, eine einzelne Task-Instanz zu starten, ohne sich um AbhĂ€ngigkeiten kĂŒmmern zu mĂŒssen. DarĂŒber hinaus kann es ĂŒberLocalExecutor, auch wenn Sie einen Celery-Cluster haben.- Etwa das gleiche tut
test, nur dass es nicht in die Datenbank schreibt. connectionsermöglicht das massenhafte Erstellen von Verbindungen aus der Shell.
- â ist eine ziemlich anspruchsvolle Interaktionsmethode, die fĂŒr Plugins gedacht ist, nicht fĂŒr das Herumprobieren. Aber wer hĂ€lt uns davon ab, zu
/home/airflow/dags, zu startenipythonund einfach loszulegen? Man kann beispielsweise alle Verbindungen mit folgendem Code exportieren: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) - Verbindung zur Airflow-Metadatenbank. Ich empfehle nicht, in diese zu schreiben, aber es ist deutlich schneller und einfacher, die ZustĂ€nde von Tasks fĂŒr verschiedene spezifische Metriken abzurufen, als ĂŒber eine der APIs.
Sagen wir, nicht alle unsere Tasks sind idempotent und können manchmal fehlschlagen, was in Ordnung ist. Aber mehrere Fehler sind schon verdÀchtig, da sollte man nachforschen.
Vorsicht, 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
Links
NatĂŒrlich sind die ersten zehn Links in der Google-Suche Inhalte aus meinem Airflow-Ordner.
- â NatĂŒrlich sollte man mit der offiziellen Dokumentation beginnen, aber wer liest schon Anleitungen?
- â Lesen Sie zumindest die Empfehlungen der Entwickler.
- â Der Anfang: eine grafische Ăbersicht der BenutzeroberflĂ€che
- â Gut erklĂ€rte Grundbegriffe, falls Sie (ausnahmsweise!) etwas bei mir nicht verstanden haben.
- â Ein kurzer Leitfaden zur Konfiguration eines Airflow-Clusters.
- â Fast ein ebenso interessanter Artikel, nur mit mehr Formalismus und weniger Beispielen.
- â Arbeiten im Zusammenspiel mit Celery.
- â Ăber Idempotenz von Tasks, das Laden nach ID statt Datum, Transformationen, Dateistrukturen und andere interessante Aspekte.
- â Task-AbhĂ€ngigkeiten und Trigger-Regeln, die ich nur am Rande erwĂ€hnt habe.
- â Wie man gewisse âes funktioniert wie geplantâ beim Scheduler ĂŒberwindet, verlorene Daten lĂ€dt und die PrioritĂ€ten der Tasks festlegt.
- â nĂŒtzliche SQL-Abfragen fĂŒr die Metadaten von Airflow.
- â es gibt einen hilfreichen Abschnitt ĂŒber die Erstellung eines benutzerdefinierten Sensors.
- â eine interessante kurze Notiz zum Aufbau einer Infrastruktur auf AWS fĂŒr Data Science.
- â gĂ€ngige Fehler (wenn jemand die Anleitungen nicht liest).
- â lĂ€cheln Sie, wie Leute das Speichern von Passwörtern umgangen haben, obwohl sie einfach Connections verwenden könnten.
- â implizite Ăbertragung von DAGs, KontextĂŒbergabe in Funktionen, erneut ĂŒber AbhĂ€ngigkeiten und das Ăberspringen von Task-AusfĂŒhrungen.
- â zur Verwendung von
StandardargumentenundParameterin Vorlagen sowie zu Variablen und Verbindungen. - â eine ErzĂ€hlung darĂŒber, wie der Scheduler auf Airflow 2.0 vorbereitet wird.
- â ein etwas veralteter Artikel ĂŒber das Deployment unseres Clusters in
version: '1' services: simplesample-sonar: image: sonarqube:lts ports: - 9001:9000 - 9092:9092 network_mode: bridge. - â dynamische Aufgaben mithilfe von Vorlagen und KontextĂŒbergabe.
- â Standard- und benutzerdefinierte Benachrichtigungen ĂŒber E-Mail und Slack.
- â Task-Verzweigungen, Makros und XCom.
Und die in diesem Artikel verwendeten Links:
- â Platzhalter, die in Vorlagen verwendet werden können.
- â Verbreitete Fehler beim Erstellen von DAGs.
- â
version: '1' services: simplesample-sonar: image: sonarqube:lts ports: - 9001:9000 - 9092:9092 network_mode: bridgefĂŒr Experimente, Debugging und mehr. - â Python-Wrapper fĂŒr die Telegram REST API.
Quelle: habr.com




