
Hallo, Habr! In diesem Artikel möchte ich über ein bemerkenswertes Werkzeug zur Entwicklung von Batch-Prozessen zur Datenverarbeitung sprechen, zum Beispiel in der Infrastruktur eines Unternehmens-DWH oder Ihres DataLake. Es geht um Apache Airflow (im Folgenden Airflow). Es wird ungerechtfertigt wenig Aufmerksamkeit auf Habr geschenkt, und im Hauptteil werde ich versuchen, Sie davon zu überzeugen, dass Airflow mindestens eine Überlegung wert ist, wenn es um die Wahl eines Planers für Ihre ETL/ELT-Prozesse geht.
Zuvor habe ich eine Reihe von Artikeln zum Thema DWH geschrieben, als ich bei der Tinkoff Bank gearbeitet habe. Jetzt bin ich Teil des Teams der Mail.Ru Group und beschäftige mich mit der Entwicklung einer Plattform zur Datenanalyse im Bereich Gaming. Tatsächlich werden wir mit unserem Team hier über unsere Plattform zur Datenanalyse berichten, sobald Neuigkeiten und interessante Lösungen erscheinen.
Prolog
Fangen wir an. Was ist Airflow? Es ist eine Bibliothek (oder ) zur Entwicklung, Planung und Überwachung von Arbeitsabläufen. Das Hauptmerkmal von Airflow ist, dass zur Beschreibung (Entwicklung) der Prozesse Code in Python verwendet wird. Daraus ergeben sich zahlreiche Vorteile für die Organisation Ihres Projekts und die Entwicklung: Im Grunde genommen ist Ihr (zum Beispiel) ETL-Projekt einfach ein Python-Projekt, und Sie können es so organisieren, wie es Ihnen passt, unter Berücksichtigung der Infrastruktur, der Teamgröße und anderer Anforderungen. Das Werkzeug ist einfach zu bedienen. Verwenden Sie zum Beispiel PyCharm + Git. Das ist großartig und sehr praktisch!
Nun betrachten wir die grundlegenden Entitäten von Airflow. Wenn Sie deren Wesen und Zweck verstehen, werden Sie die Architektur der Prozesse optimal organisieren. Wahrscheinlich ist die Hauptentität der Directed Acyclic Graph (im Folgenden DAG).
DAG
DAG ist eine sinnvolle Gruppierung Ihrer Aufgaben, die Sie in einer genau festgelegten Reihenfolge nach einem festgelegten Zeitplan ausführen möchten. Airflow bietet eine benutzerfreundliche Weboberfläche zur Arbeit mit DAGs und anderen Entitäten:

Ein DAG könnte so aussehen:

Der Entwickler legt bei der Gestaltung eines DAGs eine Reihe von Operatoren fest, auf denen die Aufgaben innerhalb des DAGs basieren. Hier stoßen wir auf eine weitere wichtige Entität: den Airflow Operator.
Operatoren
Ein Operator ist eine Entität, auf deren Grundlage Instanzen von Aufgaben erstellt werden, in denen beschrieben wird, was während der Ausführung einer Aufgabeninstanz passieren wird. beinhalten bereits eine Sammlung von einsatzbereiten Operatoren. Beispiele:
- BashOperator — Operator zum Ausführen von Bash-Befehlen.
- PythonOperator — Operator zum Aufrufen von Python-Code.
- EmailOperator — Operator zum Versenden von E-Mails.
- HTTPOperator — Operator zur Verarbeitung von HTTP-Anfragen.
- SqlOperator — Operator zum Ausführen von SQL-Code.
- Sensor — Operator, der auf das Eintreten eines Ereignisses wartet (z.B. Erreichen eines bestimmten Zeitpunkts, Erscheinen einer bestimmten Datei, Datensatz in einer Datenbank, Antwort aus einer API usw.).
Es gibt spezifischere Operatoren: DockerOperator, HiveOperator, S3FileTransferOperator, PrestoToMysqlOperator, SlackOperator.
Sie können auch Operatoren entwickeln, die auf Ihre speziellen Anforderungen ausgerichtet sind, und diese in Ihrem Projekt verwenden. Zum Beispiel haben wir MongoDBToHiveViaHdfsTransfer erstellt, einen Operator zum Export von Dokumenten aus MongoDB nach Hive, sowie mehrere Operatoren für die Arbeit mit : CHLoadFromHiveOperator und CHTableLoaderOperator. Grundsätzlich, sobald in einem Projekt häufig verwendeter Code entsteht, der auf grundlegenden Operatoren basiert, kann man in Erwägung ziehen, diesen in einen neuen Operator zu bündeln. Das vereinfacht die weitere Entwicklung und erweitert Ihre Bibliothek an Operatoren im Projekt.
Jetzt müssen alle diese Instanzen von Aufgaben ausgeführt werden, und nun geht es um den Scheduler.
Scheduler
Der Task-Scheduler in Airflow basiert auf . Celery ist eine Python-Bibliothek, die es ermöglicht, eine Warteschlange sowie asynchrone und verteilte Ausführung von Aufgaben zu organisieren. Seitens Airflow werden alle Aufgaben in Pools unterteilt. Pools werden manuell erstellt. In der Regel dienen sie dazu, die Last bei der Arbeit mit der Quelle zu begrenzen oder die Aufgaben innerhalb von DWH zu typisieren. Pools können über die Web-Oberfläche verwaltet werden:

Jeder Pool hat eine Begrenzung der Anzahl an Slots. Bei der Erstellung eines DAG wird ihm ein Pool zugewiesen:
ALERT_MAILS = Variable.get("gv_mail_admin_dwh")
DAG_NAME = 'dma_load'
OWNER = 'Vasya Pupkin'
DEPENDS_ON_PAST = True
EMAIL_ON_FAILURE = True
EMAIL_ON_RETRY = True
RETRIES = int(Variable.get('gv_dag_retries'))
POOL = 'dma_pool'
PRIORITY_WEIGHT = 10
start_dt = datetime.today() - timedelta(1)
start_dt = datetime(start_dt.year, start_dt.month, start_dt.day)
default_args = {
'owner': OWNER,
'depends_on_past': DEPENDS_ON_PAST,
'start_date': start_dt,
'email': ALERT_MAILS,
'email_on_failure': EMAIL_ON_FAILURE,
'email_on_retry': EMAIL_ON_RETRY,
'retries': RETRIES,
'pool': POOL,
'priority_weight': PRIORITY_WEIGHT
}
dag = DAG(DAG_NAME, default_args=default_args)
dag.doc_md = __doc__Der Pool, der auf DAG-Ebene festgelegt wurde, kann auf Aufgabenebene überschrieben werden.
Für die Planung aller Aufgaben in Airflow ist ein separater Prozess verantwortlich — der Scheduler. Der Scheduler kümmert sich im Wesentlichen um die gesamte Mechanik der Einreihung von Aufgaben zur Ausführung. Eine Aufgabe durchläuft mehrere Phasen, bevor sie zur Ausführung kommt:
- In DAG wurden die vorherigen Aufgaben erledigt, die neue kann in die Warteschlange gestellt werden.
- Die Warteschlange wird je nach Priorität der Aufgaben sortiert (Prioritäten können ebenfalls verwaltet werden), und wenn im Pool ein freier Slot vorhanden ist, kann die Aufgabe bearbeitet werden.
- Wenn ein freier Celery-Worker verfügbar ist, wird die Aufgabe an ihn weitergeleitet; die Arbeit, die Sie im Auftrag programmiert haben, beginnt unter Verwendung des jeweiligen Operators.
Ganz einfach.
Der Scheduler arbeitet über alle DAGs hinweg sowie über alle Aufgaben innerhalb der DAGs.
Um den Scheduler mit dem DAG arbeiten zu lassen, muss das DAG einen Zeitplan erhalten:
dag = DAG(DAG_NAME, default_args=default_args, schedule_interval='@hourly')Es gibt eine Reihe von fertigen Presets: @once, @hourly, @daily, @weekly, @monthly, @yearly.
Es können auch Cron-Ausdrücke verwendet werden:
dag = DAG(DAG_NAME, default_args=default_args, schedule_interval='*\/10 * * * *')Ausführungsdatum
Um zu verstehen, wie Airflow funktioniert, ist es wichtig zu wissen, was das Ausführungsdatum für ein DAG ist. In Airflow hat ein DAG eine Dimension des Ausführungsdatums, d. h. basierend auf dem Zeitplan des DAGs werden Instanzen von Aufgaben für jedes Ausführungsdatum erstellt. Und für jedes Ausführungsdatum können die Aufgaben wiederholt ausgeführt werden – oder beispielsweise kann das DAG gleichzeitig in mehreren Ausführungsdaten laufen. Dies wird hier anschaulich dargestellt:

Leider (oder vielleicht auch zum Glück, das hängt von der Situation ab), wenn die Implementierung einer Aufgabe im DAG geändert wird, erfolgt die Ausführung in früheren Ausführungsdaten bereits unter Berücksichtigung der Anpassungen. Das ist gut, wenn Daten in vergangenen Zeiträumen mit einem neuen Algorithmus neu berechnet werden müssen, aber schlecht, weil die Reproduzierbarkeit des Ergebnisses verloren geht (natürlich steht es Ihnen frei, aus Git die benötigte Version der Quelle zurückzuholen und einmalig das zu berechnen, was notwendig ist, so wie es benötigt wird).
Generierung von Aufgaben
Die Implementierung eines DAGs ist Code in Python, daher haben wir eine sehr bequeme Möglichkeit, die Codegröße beim Arbeiten mit beispielsweise shardierten Quellen zu reduzieren. Nehmen wir an, Sie haben drei MySQL-Shards als Quelle und müssen in jeden zugreifen, um Daten abzurufen. Dabei unabhängig und parallel. Der Python-Code in einem DAG könnte folgendermaßen aussehen:
connection_list = lv.get('connection_list')
export_profiles_sql = '''
SELECT
id,
user_id,
nickname,
gender,
{{params.shard_id}} as shard_id
FROM profiles
'''
for conn_id in connection_list:
export_profiles = SqlToHiveViaHdfsTransfer(
task_id='export_profiles_from_' + conn_id,
sql=export_profiles_sql,
hive_table='stg.profiles',
overwrite=False,
tmpdir='\/data\/tmp',
conn_id=conn_id,
params={'shard_id': conn_id[-1:], },
compress=None,
dag=dag
)
export_profiles.set_upstream(exec_truncate_stg)
export_profiles.set_downstream(load_profiles)Der DAG sieht folgendermaßen aus:

Dabei können Sie einen Shard hinzufügen oder entfernen, indem Sie einfach die Einstellung anpassen und das DAG aktualisieren. Praktisch!
Es ist auch möglich, komplexere Codegenerierung zu verwenden, zum Beispiel mit Datenquellen in Form von Datenbanken zu arbeiten oder die Tabellenstruktur, den Arbeitsalgorithmus mit der Tabelle zu beschreiben und, unter Berücksichtigung der Besonderheiten der DWH-Infrastruktur, den Prozess zum Laden von N Tabellen in Ihr Speicher zu generieren. Oder beispielsweise mit einer API, die keine Aufrufe mit einer Liste von Parametern unterstützt, können Sie aus dieser Liste N Aufgaben im DAG generieren, die Parallelität der API-Anfragen durch einen Pool begrenzen und die benötigten Daten aus der API abrufen. Flexibel!
Repository
In Airflow gibt es ein eigenes Backend-Repository, eine Datenbank (kann MySQL oder Postgres sein, bei uns ist es Postgres), in der die Zustände von Aufgaben, DAGs, Verbindungswerte, globale Variablen usw. gespeichert sind. Hier möchte ich erwähnen, dass das Repository in Airflow sehr einfach (ungefähr 20 Tabellen) und praktisch ist, wenn Sie einen eigenen Prozess darauf aufbauen möchten. Es erinnert an die 100500 Tabellen im Informatica-Repository, die man lange studieren musste, um zu verstehen, wie man eine Anfrage aufbaut.
Überwachung
Angesichts der Einfachheit des Repositories können Sie selbst einen für Sie bequemen Prozess zur Überwachung der Aufgaben aufbauen. Wir verwenden ein Notizbuch in Zeppelin, um den Zustand der Aufgaben zu überprüfen:

Es kann auch das Web-Interface von Airflow sein:

Der Code von Airflow ist offen, daher haben wir bei uns eine Benachrichtigung in Telegram hinzugefügt. Jede aktive Instanz einer Aufgabe, wenn ein Fehler auftritt, sendet Spam in die Telegram-Gruppe, in der das gesamte Entwicklungsteam und der Support sind.
Wir erhalten über Telegram eine schnelle Reaktion (falls erforderlich) und über Zeppelin einen Überblick über die Aufgaben in Airflow.
Insgesamt
Airflow ist in erster Linie Open Source, und man sollte keine Wunder erwarten. Seien Sie bereit, Zeit und Mühe zu investieren, um eine funktionierende Lösung zu entwickeln. Das Ziel ist erreichbar, glauben Sie mir, es lohnt sich. Geschwindigkeit der Entwicklung, Flexibilität, einfache Hinzufügung neuer Prozesse – Sie werden es mögen. Natürlich muss man viel Aufmerksamkeit auf die Organisation des Projekts und die Stabilität von Airflow selbst legen: Wunder gibt es nicht.
Momentan arbeitet unser Airflow täglich an etwa 6.500 Aufgaben. Sie sind im Charakter recht unterschiedlich. Es gibt Aufgaben zum Laden von Daten in das Haupt-DWH aus vielen verschiedenen und sehr spezifischen Quellen, es gibt Aufgaben zur Berechnung von Dashboards innerhalb des Haupt-DWH, es gibt Aufgaben zur Veröffentlichung von Daten in ein schnelles DWH, es gibt viele, viele verschiedene Aufgaben - und Airflow verarbeitet sie alle Tag für Tag. Wenn wir in Zahlen sprechen, sind das 2,3 Tausend ELT-Aufgaben unterschiedlicher Komplexität innerhalb des DWH (Hadoop), etwa 2,5 Hundert Datenbanken Quellen, das ist ein Team von vier ETL-Entwicklern, die sich auf die ETL-Datenverarbeitung im DWH und die ELT-Datenverarbeitung innerhalb des DWH aufteilen und natürlich noch einen Administrator, der sich um die Infrastruktur des Dienstes kümmert.
Zukunftspläne
Die Anzahl der Prozesse wächst unweigerlich, und das Hauptthema, mit dem wir uns im Bereich der Airflow-Infrastruktur beschäftigen werden, ist das Skalieren. Wir möchten einen Airflow-Cluster aufbauen, ein paar Worker-Knoten für Celery bereitstellen und einen sich selbst duplizierenden Kopf mit Aufgabenplanung und Repository einrichten.
Epilog
Das ist natürlich längst nicht alles, was ich über Airflow erzählen möchte, aber die Hauptpunkte habe ich versucht zu beleuchten. Der Appetit kommt beim Essen, probieren Sie es aus - und es wird Ihnen gefallen 🙂
Quelle: habr.com
