Apache Airflow
Kostenlos
Apache Airflow ist eine Open-Source-Workflow-Orchestrierungsplattform, die Aufgabenabhängigkeiten und Planungslogik über Python DAG definiert. Es wird häufig in Datenpipelines, ML-Pipelines und der Automatisierung der Cloud-Infrastruktur eingesetzt.
ApacheAirflow
Kernparameter und Statistiken
Apache Airflow ist kein „KI-Tool“, sondern die Infrastruktur der KI-Datenpipeline – es ist für die Orchestrierung der Abhängigkeiten und die Planungslogik von Aufgaben wie Modelltraining, Datenbereinigung, Feature-Engineering und Modellbereitstellung verantwortlich. Es ist die am leichtesten zu übersehende, aber kritischste Ebene im KI-Produktionssystem. Airflow wurde 2014 von Airbnb entwickelt. Es wurde 2016 in den Apache Incubator aufgenommen und 2019 als Top-Level-Projekt abgeschlossen. Es ist nach wie vor die am weitesten verbreitete Workflow-Orchestrierungsplattform im Bereich Data Engineering.
| Projekte | Öffentliche Informationen |
|---|---|
| Offizielle Positionierung | Open-Source-Workflow-Orchestrierungsplattform |
| Kernparadigma | Directed Graph (DAG), definiert im Python-Code |
| Planungs-Engine | Verteilter Scheduler + Executor (Celery, Kubernetes, CeleryKubernetes, Local, Sequential) |
| Bereitstellungsformular | Selbstgehosteter (einzelner Rechner/Cluster), verwalteter Cloud-Dienst (Amazon MWAA, Google Cloud Composer, Astronomer) |
| Open-Source-Lizenz | Apache 2.0 |
| Gemeinschaftsgröße | GitHub über 39.000+ Sterne, 2.000+ Forks, über 800 Mitwirkende |
| Anbieter-Ökosystem | Über 100 offizielle Anbieter + Hunderte von Community-Anbietern, die AWS/GCP/Azure/Snowflake/Databricks/Spark usw. abdecken. |
| Kernsprache | Python |
| Datenbank-Backend | PostgreSQL, MySQL, SQLite (für Entwicklung) |
| Nachrichtenwarteschlange | Redis/RabbitMQ |
Branchenstatus: Das DAG-as-Code-Paradigma von Airflow ist zum De-facto-Standard für die Workflow-Orchestrierung geworden. Die drei Mainstream-Cloud-Anbieter AWS, GCP und Azure bieten alle verwaltete Airflow-Dienste an, und Astronomer bietet eine mandantenfähige Verwaltungsplattform auf Unternehmensebene. Im CNCF-Cloud-Native-Panorama wird Airflow als Benchmark-Projekt im Bereich Workflow und Terminplanung aufgeführt.
Benutzer- und Markterkennung
Die Marktposition von Airflow lässt sich anhand von drei Dimensionen beobachten: Community-Aktivität, Unternehmensakzeptanz und Investitionen von Cloud-Anbietern.
Community-Aktivität: Airflow hat etwa 39.000 Sterne, über 2.000 Forks und über 800 aktive Mitwirkende auf GitHub. Dies ist die größte Community unter den Open-Source-Workflow-Orchestrierungstools. Jede Veröffentlichung einer Hauptversion (z. B. 2.0, 2.9, 2.10) löst einen Höhepunkt bei den Community-Beiträgen aus. Der Slack-Kanal von Airflow hat Zehntausende registrierte Benutzer, und jeden Monat gibt es Hunderte von Diskussionsthreads rund um die Nutzung von DAG-Writing-Anbietern und die Leistungsoptimierung.
Einführung in Unternehmen: Airflow wird in Produktionsumgebungen von Tausenden von Unternehmen auf der ganzen Welt eingesetzt und deckt vertikale Branchen wie Finanzen, E-Commerce, Technologie, medizinische Versorgung und Fertigung ab. Zu den bekannten Nutzern zählen Airbnb (Urheber), Twitter/Lyft/Slack (Early Adopters), Walmart, JPMorgan Chase, Adobe, Intuit und andere. Auf dem chinesischen Markt haben erstklassige Internetunternehmen wie ByteDance, Alibaba und Meituan Airflow oder seine eigenen Derivate in großem Umfang eingesetzt. Von Airbnb im Jahr 2021 veröffentlichte Daten zeigen, dass sein Airflow-Cluster täglich mehr als 500.000 Aufgaben ausführt.
Investitionen von Cloud-Anbietern: Amazon MWAA (Managed Workflows for Apache Airflow) hat seit GA im Jahr 2021 die verfügbaren Bereiche und Funktionen weiter ausgebaut. Google Cloud Composer ist das Hauptprodukt für die native Datenpipeline-Orchestrierung von GCP, und die Data Factory von Azure bietet auch eine integrierte Airflow-Integration. Die Hosting-Investition der drei großen Cloud-Anbieter bestätigt die unersetzliche Position von Airflow im Bereich der Orchestrierung von Datenpipelines.
Kostenvorteil
Die Kostenstruktur von Airflow unterscheidet sich von kommerziellen SaaS-Tools und muss in drei Ebenen unterteilt werden: „Open-Source-Lizenzkosten + selbstgehostete Betriebs- und Wartungskosten + Beschaffungskosten für verwaltete Dienste“.
Einzelne/C-seitige Benutzer: Keine Lizenzkosten, aber der Hardware-Grenzwert ist vorhanden. Die Airflow Community Edition ist völlig kostenlos, ohne Funktionseinschränkungen oder Kontoschließungen. Einzelpersonen können Airflow über Docker Compose oder den virtuellen Python-Kontext auf ihren Laptops zum Lernen oder für kleine Datenpipelines starten. Wenn es sich jedoch um umfangreiche DAGs oder Planungen mit hoher Parallelität handelt, kommt es schnell zu Leistungsengpässen, wenn das SQLite-Backend und der Sequential Executor auf einem einzelnen Computer bereitgestellt werden.
Entwickler/Teams: Selbstgehostet ohne Lizenz, Betriebs- und Wartungskosten summieren sich Schritt für Schritt. Selbsthosting in Produktionsqualität erfordert die Bereitstellung:
- Metadatenbank (PostgreSQL/MySQL) – Die jährlichen Kosten für die Cloud-Datenbank betragen etwa 1.200–6.000 Yuan (abhängig von den Spezifikationen).
- Nachrichtenwarteschlange (Redis/RabbitMQ) – etwa 600–3.000 Yuan/Jahr
- Scheduler- und Worker-Knoten (Kubernetes Pod oder EC2) – 5–50 Einheiten, monatliche Gebühr 3.000–30.000 Yuan
- Protokollspeicherung und -überwachung (S3/GCS + CloudWatch/Prometheus) - Floating je nach Datenvolumen
Unternehmen/Privatunternehmen: Kosten für Managed Services im Vergleich zu den vollen Kosten für Eigenbetrieb und Wartung. Der Vergleich der Mainstream-Hosting-Dienste ist wie folgt (im Folgenden sind öffentliche Referenzpreise aufgeführt, die den Echtzeitseiten jedes Dienstanbieters unterliegen):
| Vergleich | Amazon MWAA | Google Cloud Composer | Astronom | Selbstgehostet (K8s-Cluster) |
|---|---|---|---|---|
| Preismodell | Peripheriegebühr + Worker-vCPU-Stunde | Peripheriegebühr + Worker-vCPU-Stunde | Abonnement (nach Knoten oder pro Benutzer) | Nach tatsächlicher Infrastrukturnutzung + Betriebs- und Wartungspersonal |
| Kleine angeschlossene monatliche Gebühr (Schätzung) | ~3.000-8.000 Yuan | ~2.500-7.000 Yuan | Nicht bekannt gegeben | ~2.000–5.000 Yuan (nur Cloud-Ressourcen) |
| Mittlere, begrenzte monatliche Gebühr (Schätzung) | ~10.000-30.000 Yuan | ~8.000-25.000 Yuan | Geschäftsbestätigung erforderlich | ~8.000-20.000 Yuan (einschließlich Betrieb und Wartung) |
| Betriebs- und Wartungspersonal | Anteil des Cloud-Anbieters | Anteil des Cloud-Anbieters | Plattformanbieter vollständig verwaltet | Mindestens 0,5–1 Vollzeitstellen |
| Anwendbare Szenarien | Tiefe AWS-Integration | Tiefe GCP-Integration | Multi-Cloud-/Multi-Tenant-/Enterprise-Governance | Compliance-Isolierung/hoher Grad an Individualisierung |
Tipp zu versteckten Kosten:
- DAG-Debugging und -Inspektion zeitaufwändig: Der Debugging-Link von Airflow (Analysefehler → erneutes Parsen des Planers → Worker-Ausführung → Protokoll-Traceback) kann in einem umfangreichen DAG-Szenario für jedes Debugging 15 bis 60 Minuten dauern. Dies sind die versteckten Kosten, die Teams am leichtesten unterschätzen können.
- Migrationskosten: Bei der Migration von einem selbst gehosteten zu einem gehosteten Dienst oder umgekehrt ist der DAG-Code selbst portierbar, aber die Migration von Connector-Anmeldeinformationen, Kontextvariablen, historischen Metadaten und Protokollen erfordert zusätzlichen Aufwand.
Hauptfunktionen
Das Funktionssystem von Airflow dreht sich um die vier Abschnitte „Definition → Scheduling → Monitoring → Extension“. Sein Kernwert ist nicht eine einzelne Funktion, sondern die Synergie zwischen diesen Funktionen.
-
DAG-Definition (Python-as-Code): Verwenden Sie Standard-Python-Code, um Aufgaben (Operatoren), Abhängigkeiten (
>>/<</set_upstream) und Ausführungsstrategien (Anzahl der Wiederholungsversuche, Zeitüberschreitungen, Warteschlangen) zu deklarieren. Synergieeffekt: DAG-Code ist von Natur aus versioniert (Git), testbar (pytest-airflow) und wiederverwendbar (angepasste Operator-Paketverwaltung), was den Hauptschmerzpunkt herkömmlicher grafischer Orchestrierungstools löst: „nicht zu wissen, wer was geändert hat, und nicht in der Lage zu sein, CR nach der Durchführung von Änderungen durchzuführen“. -
Planungs-Engine (zeitgesteuert + Ereignis + Sensor): Unterstützt die Timing-Auslösung von Cron-Ausdrücken und unterstützt außerdem den Datensensor, der darauf wartet, dass Upstream-Daten bereit sind, den externen Aufgabensensor im DAG, der darauf wartet, dass der Dateisensor die Dateilandung überwacht usw. Synergieeffekt: Der Sensor und der Planer können kontinuierlich externe Bedingungen erkennen, ohne Worker-Ressourcen zu verbrauchen. Wenn die Bedingungen erfüllt sind, werden nachgelagerte Aufgaben automatisch ausgelöst – dadurch entfällt die manuelle Prüfung im vollautomatischen Zusammenhang „Warten auf Dateneingang → Starten der Pipeline → Berichterstellung ist abgeschlossen“.
-
Web-Benutzeroberfläche und Beobachtbarkeit: Visualisieren Sie den DAG-Laufstatus, das Aufgaben-Gantt-Diagramm, den Aufgabendauertrend, die Rasteransicht und die Herkunft auf Aufgabenebene. Synergie: Das Gantt-Diagramm macht Engpassaufgaben intuitiv sichtbar, die Herkunftsanalyse hilft dabei, die Quelle von Datenqualitätsproblemen zu lokalisieren, und die Rasteransicht zeigt die Statusverteilung jedes DAG-Laufs nach Ausführungsdatum an – die Kombination dieser drei ermöglicht es dem Betriebs- und Wartungspersonal, zu lokalisieren, „welcher Schritt einer Aufgabe in welchem Zeitfenster langsamer wird“, ohne die Protokolle einzeln lesen zu müssen.
-
Anbieter-Ökosystem (über 100 Konnektoren): Der offizielle Anbieter deckt AWS (S3, EMR, Lambda, Redshift, SageMaker), GCP (BigQuery, Cloud Storage, Dataflow, Vertex AI), Azure (Blob, Data Lake, Synapse), Snowflake, Databricks, Spark, Kubernetes, Docker, Slack, PagerDuty usw. ab. Synergie: Mehrere Anbieter können gleichzeitig in Reihe geschaltet werden DAG – zum Beispiel Daten aus Snowflake lesen → Spark-Cluster führt Konvertierung durch → in GCS schreiben → Dataflow für nachfolgende Analyse auslösen. Im gesamten Prozess muss kein API-Aufrufcode geschrieben werden. Deklarieren Sie einfach den entsprechenden Operator im DAG.
-
Erweiterbare Architektur (Operator + Hook + Executor):
- Operator: definiert „was zu tun ist“ (z. B. „PythonOperator“ führt Python-Funktionen aus, „BashOperator“ führt Shell-Befehle aus)
- Hook: Kapselt Verbindungsdetails externer Dienste (z. B. „S3Hook“, um AWS-Anmeldeinformationen und Wiederholungsversuche automatisch zu verwalten)
- Executor: Entscheiden Sie, wie ausgeführt werden soll (sequentiell → lokal seriell, lokal → lokal parallel, Celery → verteilte Warteschlange, KubernetesExecutor → unabhängiger Pod pro Aufgabe)
- Synergie: Die hierarchische Entkopplung der drei ermöglicht es Airflow, „SequentialExecutor“ im Entwicklungskontext zu verwenden und in der Produktion nahtlos zu „CeleryExecutor“ oder „KubernetesExecutor“ zu wechseln, ohne den DAG-Code zu ändern – dies ist Airflows „Null-Code-Änderung“-Erweiterungsfähigkeit vom eigenständigen Aufgabenexperiment zur Planung mit hoher Parallelität auf Produktionsebene.
Modell- und Versionsentwicklung
Als Open-Source-Projekt spiegeln die Versionsiterationen von Airflow die Entwicklung der Anforderungen an die Orchestrierung von Data-Engineering-Workflows von „Skriptplanung“ zu „Cloud Native + AI-Pipeline“ wider.
1.x-Ära (2015–2020): Etablierung des DAG-Paradigmas
- Airflow 1.0 (2015): Entwickelt innerhalb von Airbnb von Maxime Beauchemin, sind die Kernkonzepte von DAG, Operator und Scheduler alle etabliert.
- Airflow 1.8 (2018): Einführung von „SubDAG“ und „BranchOperator“ zur Verbesserung der DAG-Wiederverwendungsfunktionen. Dies ist eine der von der Community am häufigsten verwendeten 1.x-Versionen.
- Airflow 1.10 (2019-2020): Einführung der ersten Hauptversion nach der Apache-Stufe, Hinzufügen von „KubernetesPodOperator“, Stabilisierung der REST-API, Verbesserung der Protokollspeicherung und der Benutzeroberfläche. Die 1.10-Serie wird weiterhin bis 1.10.15 iteriert.
2.x-Ära (2020 bis heute): Architekturrekonstruktion und Cloud-nativ
- Airflow 2.0 (2020-12): Meilenstein-Veröffentlichung. Neuschreiben des Schedulers (Unterstützung der HA-Hochverfügbarkeit), Einführung der „TaskFlow API“ (vereinfachtes DAG-Schreiben) und native Kubernetes Executor-Unterstützung. Der Migrationspfad von 1.10 auf 2.0 erfordert eine manuelle Anpassung.
- Airflow 2.1-2.2 (2021): Einführung der Rasteransicht (anstelle der alten Baumansicht), automatischer DAG-Registrierung und Unterstützung für Aufgabengruppen. Wichtige Änderung: Grid View löst den Engpass bei der Visualisierungsleistung in Tausenden von DAG-Run-Szenarien.
- Airflow 2.3-2.4 (2022): Unterstützung für dynamische DAG-Generierung, verbesserte Scheduler-Leistung (50 %+ Reduzierung der Parsing-Zeit), Trennung von Provider-Paketen von Kernpaketen. Wichtige Änderungen: Die Anbieterentkopplung reduziert Abhängigkeitskonflikte in Kernpaketen und jeder Anbieter kann unabhängig iterieren.
- Airflow 2.5-2.6 (2023): DAG-Versionierung, Prüfprotokolle, verbesserte Unterstützung für parallele Tasks der „@task“-Decorator-Matrix.
- Airflow 2.7-2.8 (2024): Verbesserter Scheduler-Heartbeat-Mechanismus, Optimierung des Datenbankverbindungspools, Web-UI-Dark-Mode-Python-3.12-Unterstützung.
- Airflow 2.9 (2025-12): Datensatzgesteuerte DAG-Planung – abhängige Planung basierend auf der Datenausgabe ersetzt reine Zeitplanung, was ein wichtiger Schritt zur Erreichung „echter ereignisgesteuerter Datenpipelines“ ist. Das Protokoll-Streaming auf Aufgabenebene wurde ebenfalls verbessert.
- Airflow 2.10 (2026-05): Die neueste stabile Version (noch kein offizielles genaues Datum). Konzentrieren Sie sich auf die Optimierung des Metadaten-Datenbankdrucks des Schedulers in großen DAG-Szenarien (10.000+ DAG), eine verbesserte Asset-/Dataset-Verwaltungsschnittstelle und eine verbesserte Startgeschwindigkeit des KubernetesExecutor-Pods.
Schneller Überblick über den Versionsverlauf
| Versionsreihe | Zeit | Wichtige Änderungen | Notizen |
|---|---|---|---|
| 1,0-1,10 | 2015-2020 | DAG-Paradigma etabliert, Community-Akkumulation | 1.10.15 ist die endgültige Version von 1.x |
| 2,0 | 2020-12 | Scheduler HA, TaskFlow API, K8s Executor native Unterstützung | Meilenstein der Architekturrekonstruktion |
| 2.1-2.4 | 2021-2022 | Grid-Ansicht, Provider-Entkopplung, dynamische DAG, Optimierung der Scheduler-Leistung | Beobachtbarkeit und ökologische Expansion |
| 2,5-2,8 | 2023-2024 | DAG-Versionskontrolle, Audit-Protokoll Python 3.12, UI-Verbesserungen | Abschluss der Enterprise-Governance-Funktion |
| 2,9 | 2025-12 | Datensatzgesteuerte Planung, Protokoll-Streaming | Wichtige Ergänzungen zur ereignisgesteuerten Orchestrierung |
| 2.10 | 2026-05 | Umfangreiche DAG-Leistungsoptimierung und Asset-Management-Verbesserung | Neueste stabile Version |
Technische Vorteile
Airflow konnte seine Dominanz im Bereich der Workflow-Orchestrierung seit zehn Jahren behaupten. Sein technischer Vorteil liegt nicht in der „Einzelpunkt-Funktionsführung“, sondern in der langfristigen Rationalität von Entscheidungen auf Systemebene wie Architekturschichtung und Scheduler-Design, DAG-Analyse und Ausführungstrennung.
Vollständige Trennung von DAG-Analyse und -Ausführung: Dies ist die zentrale Architekturentscheidung von Airflow. Der Scheduler ist für die regelmäßige Analyse von Python-Dateien verantwortlich, um DAG-Objekte zu generieren (statische Analyse), und der Executor ist für die Verteilung von Aufgaben in DAG an Worker zur Ausführung verantwortlich. Die beiden kommunizieren über die Metabasis, und der Scheduler enthält nicht den Ausführungskontext des Workers. Das bedeutet: – Auch wenn der Worker-Knoten ausfällt, kann der Scheduler Aufgaben auf dem neuen Worker neu planen – Nachdem der DAG-Code aktualisiert wurde, führt der Scheduler eine automatische erneute Analyse durch und wird wirksam, ohne dass der Dienst neu gestartet werden muss. – Verschiedene Aufgaben derselben DAG können in unterschiedlichen Worker-Kontexten ausgeführt werden (Kubernetes-Pod, Celery-Container-Remote-EMR usw.)
Scheduler HA und intelligenter Parser: Der Scheduler von Airflow 2.0+ unterstützt die Bereitstellung mit mehreren Kopien und hoher Verfügbarkeit, und der Datenbanksperrmechanismus stellt sicher, dass es jeweils nur einen aktiven Scheduler gibt. Sein DAG-Parser führt ab 2.4 das Caching der Dateiänderungszeit und das inkrementelle Parsen ein – er analysiert nur DAG-Dateien, die sich seit dem letzten Parsen geändert haben, und komprimiert so die Parsing-Zeit von mehr als 10.000 DAGs von Minuten auf mehrere zehn Sekunden.
Feinkörnige Schichtung des Executors:
- SequentialExecutor: Für Entwicklung und Debugging, serielle Ausführung, mit SQLite-Backend.
- LocalExecutor: Eine einzelne Maschine führt Aufgaben parallel aus und nutzt dabei einen Multiprozesspool, der für die Produktion in kleinem Maßstab geeignet ist.
- CeleryExecutor: Implementiert einen verteilten Worker-Pool über Celery + Redis/RabbitMQ, geeignet für mittlere Größenordnungen (Hunderte bis Tausende von Aufgaben/Tag).
- CeleryKubernetesExecutor: Ein Hybrid-Executor, der Celery Worker als Rückgrat verwendet und einige Aufgaben zur besseren Isolierung an Kubernetes-Pods weiterleitet.
- KubernetesExecutor: Jede Task-Instanz startet einen unabhängigen Pod und wird nach der Ausführung automatisch zerstört. Es verfügt über die stärkste Ressourcenisolation und eignet sich für ML-Trainingsaufgaben, die eine feinkörnige Ressourcensteuerung (CPU/Speicher/GPU) erfordern.
Provider-Paketverwaltung und Versionsentkopplung: Airflow wird Provider in 2.3 vom Kernpaket trennen. Jeder Anbieter verfügt über eine unabhängige Versionsnummer und einen unabhängigen Veröffentlichungszyklus. Das bedeutet: – Benutzer müssen nur die Anbieter installieren, die sie benötigen („Apache-Airflow-Providers-Aws“ usw.), um eine Explosion der Abhängigkeiten zu vermeiden – Anbieteraktualisierungen blockieren nicht die Iteration der Airflow-Kernversion
- Community Provider kann unabhängig freigegeben werden, ohne in den Stamm einzubinden
Datensatzgesteuerte Planung: Der in 2.9+ eingeführte Datensatzmechanismus basiert nicht auf der Zeit, sondern darauf, „ob die Daten bereit sind“, um nachgelagerte Aufgaben auszulösen. Wenn eine Aufgabe einen Datensatz erzeugt (über „Outlets“ deklariert), löst Airflow automatisch alle Downstream-DAGs aus, die von diesem Datensatz abhängen. Dies ist die Schlüsselfunktion, die Airflow von einem „Zeitplaner“ zu einem „Datenplaner“ aufwertet – die Datenpipeline realisiert tatsächlich eine „ausgabegesteuerte“ Streaming-Automatisierung.
Wie man es benutzt
Der Nutzungspfad von Airflow ist in drei Phasen unterteilt: DAG einrichten, schreiben, bereitstellen und betreiben. In jeder Phase gibt es klare Schlüsseltechnologieoptionen.
Kontextuelle Konstruktion (drei typische Lösungen)
| So verwenden Sie | Anwendbare Stufe | Befehl/Operation | Beschreibung |
|---|---|---|---|
| Docker Compose (offizielles Beispiel) | Lokale Entwicklung/Lernen | curl -LfO 'https://airflow.apache.org/docs/apache-airflow/2.10.0/docker-compose.yaml' && mkdir -p ./dags ./logs ./plugins && docker-compose up |
Ein-Klick-Start, einschließlich Scheduler, Worker, Webserver, Datenbank |
| Pip-Installation | Verfügt bereits über eine Python-Umgebung | „pip install apache-airflow“ und dann „airflow db init && airflow webserver && airflow schemer“ ausführen | Flexibel, Abhängigkeiten müssen jedoch selbst verwaltet werden |
| Helmkarte (K8s-Produktion) | Produktionsbereitstellung | helm repo add apache-airflow https://airflow.apache.org && helm install airflow apache-airflow/airflow |
Offizielles Helm-Chart, unterstützt K8sExecutor, CeleryExecutor |
DAG-Schreibbeispiel
Das Folgende ist ein typischer KI-Datenpipeline-DAG, der Datenextraktion, -transformation, -laden und -training umfasst:
„Python from datetime import datetime aus dem Luftstrom-Import-DAG aus airflow.operators.python PythonOperator importieren aus airflow.providers.amazon.aws.hooks.s3 S3Hook importieren aus airflow.providers.snowflake.operators.snowflake import SnowflakeOperator
default_args = { „owner“: „data_team“, „depends_on_past“: Falsch, „Wiederholungen“: 2, „retry_delay“: timedelta(Minuten=5), }
mit DAG( dag_id="ai_training_pipeline", start_date=datetime(2026, 1, 1), Schedule_interval="@daily", Catchup=Falsch, tags=["ai", "training"], default_args=default_args, ) als dag:
extract_raw_data = SnowflakeOperator(
task_id="extract_raw_data",
sql="SELECT * FROM raw_events WHERE dt = '{{ ds }}'",
snowflake_conn_id="snowflake_prod",
)
def transform_data(**context):
# Datenbereinigung und Feature-Engineering-Logik
df = context["task_instance"].xcom_pull(task_ids="extract_raw_data")
transformiert = df.dropna().pipe(engineer_features)
Rückkehr transformiert.to_json()
transform_task = PythonOperator(
task_id="transform_data",
python_callable=transform_data,
)
upload_to_s3 = PythonOperator(
task_id="upload_to_s3",
python_callable=lambda: S3Hook(aws_conn_id="aws_prod")
.load_string(
string_data="{{ ti.xcom_pull(task_ids='transform_data') }}",
key="training/{{ ds }}/features.json",
Bucket_name="ml-features",
),
)
trigger_training = BashOperator(
task_id="trigger_training_job",
bash_command="aws sagemaker create-training-job --region us-east-1 ...",
)
extract_raw_data >> transform_task >> upload_to_s3 >> trigger_training
„
Wichtige Hinweise:
- „xcom_pull“ / „xcom_push“ wird verwendet, um kleine Datenmengen zwischen Aufgaben zu übertragen (empfohlen <100 KB)
- „schedule_interval“ unterstützt „@daily“, „@hourly“, Cron-Ausdrücke und Dataset-Objekte – Für große Dateiübertragungen sollte externer Speicher wie S3/GCS verwendet werden und die Übertragung über die Airflow-Metabasis vermieden werden
Schlüsselkonfiguration für die Produktionsbereitstellung
„yaml
docker-compose.yaml-Schlüsselkonfiguration
x-airflow-common: &airflow-common Bild: Apache/Airflow:2.10.0 Umgebung: AIRFLOWCOREEXECUTOR: CeleryExecutor AIRFLOWCORESQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres/airflow AIRFLOWCELERYRESULT_BACKEND: db+postgresql://airflow:airflow@postgres/airflow AIRFLOWCELERYBROKER_URL: redis://:@redis:6379/0 AIRFLOWSCHEDULERDAG_DIR_LIST_INTERVAL: 30 LUFTSTROMKERNPARALLELITÄT: 128 AIRFLOWCOREDAG_CONCURRENCY: 16 „
Produktpreise
Die Preisgestaltung von Airflow ist in zwei orthogonale Dimensionen unterteilt: vollständig Open Source und Managed Services. Die beiden sind kein Ersatz, sondern eine Wahl zwischen „eigenem Betrieb und Wartung vs. ausgelagertem Betrieb und Wartung“.
Community Edition (völlig kostenlos): Apache 2.0-Lizenz, keine Funktionsbeschränkung, keine Benutzerbeschränkung, keine kommerzielle Nutzungsbeschränkung. Jede Organisation kann es kostenlos herunterladen, ändern, bereitstellen und kommerziell nutzen. Dies ist der größte Preisvorteil von Airflow – keine Lizenzkosten.
Tatsächliche Kosten für Selbsthosting (in Jahren):
- Kleiner Maßstab (Einzelperson/kleines Team, <50 DAG/Tag): Die monatliche Cloud-Servergebühr beträgt etwa 200–800 Yuan, und die jährlichen Gesamtkosten betragen etwa 2.400–10.000 Yuan. – Mittlerer Umfang (Team, 200–500 DAG/Tag): 3–5 Worker-Knoten + verwaltete Datenbank + Nachrichtenwarteschlange, monatliche Gebühr beträgt etwa 5.000–15.000 Yuan, jährliche Kosten betragen etwa 60.000–180.000 Yuan.
- Großer Maßstab (Unternehmensebene, 1000+ DAG/Tag, hohe Verfügbarkeit): K8s-Cluster (10–30 Pod) + Hochverfügbarkeitsdatenbank + Redis Sentinel, monatliche Gebühr beträgt etwa 20.000–60.000 Yuan, jährliche Kosten betragen etwa 240.000–720.000 Yuan, und mindestens 0,5–1 Betriebs- und Wartungs-VZÄ sind erforderlich.
Referenz zu den gehosteten Servicekosten:
- Amazon MWAA: Es gibt eine Grenzgebühr (ab ca. 1.400 Yuan/Monat) + Worker vCPU-Stundengebühr. Geeignet für Unternehmen, die bereits im AWS-Ökosystem sind.
- Google Cloud Composer: Es gibt eine Grenzgebühr (ab ca. 1.200 Yuan/Monat) + Arbeitergebühr. Geeignet für Unternehmen, die bereits im GCP-Ökosystem sind.
- Astronomer: Abonnementbasiert, Abrechnung nach Anzahl der Knoten oder Benutzer, Bereitstellung von Mandantenfähigkeit, Zugriffskontrolle auf Teamebene und zusätzlichen Sicherheitsüberprüfungsfunktionen. Spezifische Preise erfordern eine geschäftliche Bestätigung.
Der Kernwert von Hosting-Diensten besteht darin, Betriebs- und Wartungsarbeiten wie die Hochverfügbarkeitskonfiguration des Schedulers, Datenbankwartung, Versionsaktualisierungen sowie Überwachung und Alarmierung an Cloud-Anbieter oder Plattformanbieter auszulagern. Für kleine und mittlere Unternehmen, die kein eigenes Airflow-Betriebsteam haben, sind Managed Services oft wirtschaftlicher als Selbsthosting.
Anwendungsszenarien
Die Anwendungsszenarien von Airflow gehen weit über das traditionelle ETL hinaus und es spielt eine immer zentralere Rolle in KI-gesteuerten Datenpipelines.
-
AI Training Pipeline Orchestration: Dies ist das am schnellsten wachsende Szenario für Airflow im Zeitraum 2024–2026. Typische Links: Rohdatenerfassung → Datenbereinigung und Annotation → Feature Engineering → Modellschulung (SageMaker/Kubernetes/Kubeflow) → Modellbewertung → Modellregistrierung → Modellbereitstellung (A/B-Tests). Der „KubernetesPodOperator“ oder „SageMakerOperator“ von Airflow kann GPU-Trainingsaufgaben direkt im DAG starten und Ressourcen nach Abschluss des Trainings automatisch recyceln. Kostenreduzierung und Effizienzsteigerung: Bei herkömmlichen Methoden arrangieren ML-Ingenieure manuell Trainingsschritte, überprüfen Zwischenergebnisse und lösen die nächste Stufe aus. Der manuelle Vorgang dauert etwa 30–60 Minuten, um eine einzelne Trainingspipeline zu starten. Nach der Verbindung mit Airflow wird die Pipeline vollautomatisch ausgelöst und ausgeführt. Ein manueller Eingriff ist nur erforderlich, wenn die Ergebnisse der Modellbewertung abnormal sind. Die Zeit einer einzelnen Pipeline wird auf 5–10 Minuten komprimiert, wodurch etwa 70–80 % der Orchestrierungszeit eingespart werden.
-
Data Lake/Warehouse ETL Pipeline: Extrahieren Sie Daten aus mehreren Quellsystemen (OLTP-Datenbank, Log-Stream-SaaS-API), aggregieren und bereinigen Sie sie und schreiben Sie sie dann in den Data Lake (S3/GCS/ADLS) oder das Data Warehouse (Snowflake/BigQuery/Redshift). Synergie: Die Sensor + Provider-Kombination von Airflow kann eine Echtzeit-Pipeline implementieren, die „die Extraktion auslöst, wenn Daten eintreffen“ – S3KeySensor überwacht die Dateilandung → S3ToSnowflakeOperator löst das Laden aus → SnowflakeOperator führt die Konvertierung durch → SlackWebhookOperator benachrichtigt das Datenteam. Implementierungstipps: In Cross-Cloud-Szenarien müssen Sie auf die Kompatibilität der Anbieterversion und jedes Cloud-SDK achten. Es wird empfohlen, anbieterübergreifende Integrationstests in CI hinzuzufügen.
-
Cloud-Infrastruktur und DevOps-Automatisierung: Orchestrieren Sie die Erstellung von Multi-Cloud-Ressourcen, die Erstellung von AMI-Images, die Datenbankmigration, die Zertifikatsrotation, die Compliance-Inspektion und andere Betriebs- und Wartungsprozesse. Grenze der Zusammenarbeit zwischen Mensch und Maschine: Infrastrukturerstellung, Konfigurationsprüfung, Statusbestätigung und andere Schritte können zu 100 % automatisiert werden; Für Vorgänge mit produktionsgebundenem Rollback, Datenbankschemaänderungen, Berechtigungsgenehmigung und anderen Vorgängen müssen jedoch manuelle Bestätigungspunkte eingerichtet werden („BranchPythonOperator“ oder „trigger_rule="none_failed"“ auf Aufgabenebene in Verbindung mit der manuellen Genehmigungsaufgabe). Airflow bietet „AirflowSkipException“ und „DagRunState.FAILED“ sowie andere Mechanismen zur Handhabung von Genehmigungs- und Ablehnungspfaden.
-
BI-Bericht und Datenproduktbetrieb: Geschäftsdaten automatisch täglich/wöchentlich extrahieren → Vorberechnung und Aggregation durchführen → Push an BI-Tools (Tableau/Power BI/Metabase) oder Datenprodukt-API. Der „BranchPythonOperator“ von Airflow kann automatisch die Alarmpipeline auslösen, wenn die Datenqualität nicht dem Standard entspricht, anstatt schmutzige Daten direkt zu übertragen, um die Meldung von Unfällen zu vermeiden.
Nicht für Szenarien geeignet: Echtzeit-Stream-Verarbeitung (Verzögerung auf Millisekundenebene), einmalige Skripte (Betriebs- und Wartungsaufwand übersteigt den Nutzen), Logik außerhalb der reinen DAG-Definition (z. B. die direkte Durchführung der Datentransformation in Airflow erschöpft den Worker-Speicher).
Anwendbare Personen
Die anwendbare Zielgruppe von Airflow konzentriert sich auf „mehrstufige, abhängige und geplante“ Datenverarbeitungsaufgaben und ist nicht für einstufige Skripts oder Echtzeit-Stream-Verarbeitungsszenarien geeignet.
-
Data Engineering Team (Core Users): Das Team besteht in der Regel aus mehr als 3 Data Engineers und ist für den Aufbau, die Wartung und die Überwachung von Datenpipelines auf Unternehmensebene verantwortlich. Das DAG-as-Code-Paradigma von Airflow ermöglicht die Codeüberprüfung, Versionierung und Unit-Tests von Datenpipelines genau wie Anwendungscode. Nicht für die Grenze geeignet: Wenn das Team keine Python-Grundlage hat oder nur eine Person Teilzeit an der Datenpipeline arbeitet, können die Lern-, Betriebs- und Wartungskosten von Airflow den Nutzen übersteigen. In diesem Fall empfiehlt es sich, zunächst Prefect (die Lernkurve ist flacher) oder das integrierte Planungstool des Cloud-Anbieters zu evaluieren.
-
MLOps/AI Engineer: Es ist notwendig, die mehreren Schritte des Modelltrainings, der Bewertung und der Bereitstellung in einer automatisierten Pipeline zu organisieren und diese mit CI/CD zu kombinieren, um die automatische Freigabe des Modells von der Codeübermittlung an Online-Dienste zu realisieren. „KubernetesPodOperator“ und „SageMakerOperator“ von Airflow können GPU-Jobs direkt auf dem Trainingscluster starten, erfordern jedoch, dass das Team über grundlegende Betriebs- und Wartungskenntnisse von K8s oder SageMaker verfügt. Implementierungstipps: In ML-Szenarien wird empfohlen, die Modelltrainingslogik in ein Docker-Image zu kapseln. Der DAG ist nur für die Orchestrierung und Auslösung verantwortlich und nicht für die Ausführung des kontextbezogenen Abhängigkeitsmanagements. Auf diese Weise erfordern Schulungscode-Upgrades keine Änderung des DAG.
-
Plattformbetrieb und -wartung/Plattformteam: Stellen Sie eine einheitliche Aufgabenplanungsplattform für mehrere Teams bereit (Daten-ML, Analyse, Geschäft) und müssen die mandantenfähige DAG-Isolation, Ressourcenkontingente, Protokollprüfung und Alarme verwalten. Die RBAC (rollenbasierte Zugriffskontrolle) von Airflow ist in 2.0+ ausgereift und kann mit LDAP/SSO mit der einheitlichen Unternehmensauthentifizierung verbunden werden. Nicht für Grenzen geeignet: Wenn die Organisation bereits über ein vollständiges K8s CronJob + Argo Workflows-System verfügt und keine mehrstufigen Orchestrierungsanforderungen hat, erhöht die Einführung von Airflow die Redundanz der Toolkette.
-
Datenanalyst (begrenzte Anpassung): Sehen Sie sich den Betriebsstatus des vorhandenen DAG-Frameworks an und führen Sie einfache Trigger aus (z. B. das Auffüllen historischer Daten). Die tägliche Analysearbeit basiert immer noch auf SQL und Notebook, und DAG wird nicht direkt geschrieben. Es wird empfohlen, dass das Data-Engineering-Team eine Standard-DAG-Vorlage kapselt und Analysten nur die Parameter eingeben müssen, um die Ausführung auszulösen.
Zusammenfassung und Ausblick
Apache Airflow hat sich mit seinem DAG-as-Code-Paradigma und seinem riesigen Provider-Ökosystem eine nahezu standardisierte Wettbewerbsposition im Bereich der Workflow-Orchestrierung aufgebaut. Die Hauptbarriere ist nicht eine einzelne Funktion, sondern eine Kombination der folgenden drei: Versionierbare DAG-Definition + Anbieter-Ökosystem, das gängige Cloud- und Datendienste abdeckt + Nahtlose Erweiterungsmöglichkeiten von Standalone zu Kubernetes. Diese Kombination macht Airflow zu einer wesentlichen „Basisschicht“ für Datentechnik und KI-Infrastruktur.
Aktuelle Kernvorteile:
- Die Community-Größe und die Anbieterabdeckung übertreffen ähnliche Konkurrenzprodukte (Prefect, Dagster, Argo Workflows) bei weitem. Neue Datendienste unterstützen Airflow Provider in der Regel zunächst nach ihrer Einführung.
- Flexible Sekundärentwicklungs- und Anpassungsfunktionen – vom benutzerdefinierten Operator bis zum benutzerdefinierten Executor können Unternehmen eine umfassende Kontrolle über das Planungsverhalten haben. – Die Verbesserung der Hosting-Dienste von Cloud-Anbietern hat die Hemmschwelle für kleine und mittlere Unternehmen gesenkt, Airflow zu nutzen.
Wichtige aktuelle Einschränkungen:
- Leistungsengpässe sind offensichtlich, wenn der Scheduler auf einen sehr großen Umfang erweitert wird (10.000+ DAG) - Die DAG-Analysezeit für den Metabasis-Verbindungspool und die Planungs-Heartbeat-Konkurrenz müssen durch Datenbank-Sharding und benutzerdefinierte Planungskonfiguration in großen Bereitstellungen verringert werden.
- Es gibt immer noch Probleme beim Schreiben und Debuggen von DAGs – Das lokale Debuggen basiert auf dem „Airflow-Dags-Test“, um die Ausführung zu simulieren, und Python-Syntaxfehler werden nur beim Parsen durch den Scheduler aufgedeckt, was eine Ebene langsamer ist als der REPL-Entwicklungsmodus herkömmlicher Python-Skripte. Erfordert die Unterstützung von Tools wie Pytest-Airflow oder Community Dag-Factory.
- Echtzeit- und Stream-Verarbeitung sind nicht seine Designziele - Das minimale Planungsintervall von Airflow ist auf „min_file_process_interval“ (normalerweise 30 Sekunden) begrenzt und kann nicht in Echtzeitszenarien unter einer Minute verwendet werden. Für Stream-Verarbeitungsaufgaben wird die Zusammenarbeit mit Kafka/Flink empfohlen. Airflow dient lediglich als Batch-Orchestrierungsschicht.
- Dataset-gesteuerte Planung befindet sich noch im Reifeprozess - Der in 2.9+ eingeführte Dataset-Mechanismus löst DAG-übergreifende Datenabhängigkeiten, aber die Konsistenzgarantie des Planungsdiagramms unter dem großen Dataset-Netzwerk und die Zuverlässigkeitsüberprüfung des Produktionskontexts erfordern noch mehr Community-Feedback.
Konkurrenzproduktvergleich auf einen Blick:
| Abmessungen vergleichen | Luftstrom | Präfekt | Dolch | Argo-Workflows |
|---|---|---|---|---|
| Definitionssprache | Python-DAG | Python-Dekorateur | Python + Asset-Definition | YAML |
| Planungsgranularität | Minutenebene | Zweite Ebene | Minutenebene | Minutenebene |
| UI-Beobachtbarkeit | Raster + Gantt + Herkunft | Moderne Benutzeroberfläche + Zeitleiste | Asset-Herkunftsdiagramm | Grundlegende Pod-Ansicht |
| Cloud-Native-Abschluss | K8sExecutor + Helm | K8s nativ + serverlos | Dagit + K8s | Kubernetes nativ |
| Unternehmensführung | RBAC + Audit-Protokolle | RBAC + SSO | RBAC + Teamisolation | K8s RBAC-Vererbung |
| Community und Anbieter | Über 100 Anbieter | Weniger native Anbieter | Weniger native Anbieter | Keine eigenständigen Anbieter |
| Lernkurve | Mittel-Hoch (erfordert Kenntnisse der Airflow-Architektur) | Mittel-niedrig | Mittel (Anpassbarkeit an Asset-Konzepte erforderlich) | Niedrig (YAML-Definition) |
| Anwendbarer Maßstab | Kleiner bis großer Allzweck | Mittlerer bis großer Maßstab | Mittlerer bis großer Maßstab | Kleiner bis mittlerer Maßstab |
Beschaffungs- und Einführungsrisikobewertung:
Für individuelles Lernen und Pilotversuche in kleinen Teams ist Airflow dank der Null-Lizenzkosten und der Ein-Klick-Startfunktion von Docker Compose eine nahezu risikofreie Wahl – ein Wochenende damit zu verbringen, die Umgebung einzurichten und das offizielle Tutorial auszuführen, reicht aus, um zu beurteilen, ob es den Anforderungen entspricht.
Für mittlere und große Unternehmen sollten die folgenden drei Punkte vor einer Investition sorgfältig geprüft werden:
- Investitionen in Betrieb und Wartung vs. Wahl der Hosting-Dienste: Im Selbsthosting-Modus sind mindestens 0,5 FTE für den Vollzeitbetrieb und die Wartung erforderlich (Planeroptimierung, Datenbankwartung, DAG-Debugging-Unterstützung für Versionsaktualisierungen). Wenn die Organisation nicht über Erfahrung im Airflow-Betrieb verfügt, wird dringend empfohlen, mit einem verwalteten Dienst (MWAA / Cloud Composer / Astronomer) zu beginnen – die Hosting-Gebühr ist in der Regel niedriger als die versteckten Arbeitskosten beim Selbsthosting, und der Cloud-Anbieter ist für Versions-Upgrades und die Behandlung von Infrastrukturfehlern verantwortlich.
- Der Lock-in-Effekt des DAG-Technologie-Stacks: Der DAG-Code selbst ist portierbar, aber die Migration der Anbieterkonfiguration (Verbindungszeichenfolge, Anmeldeinformationsverwaltung) und begrenzter Abhängigkeiten (Python-Pakete, Systembibliotheken) zwischen verschiedenen Bereitstellungsmethoden erfordert Tests und Überprüfung. Es wird empfohlen, Containerisierung zu verwenden, um alle DAG-Aufgaben in den frühen Phasen des Projekts auszuführen, und kontextbezogene Abhängigkeiten in Docker-Images zu kapseln, um Reibungsverluste bei zukünftigen Migrationen zu reduzieren.
- GPU-Orchestrierungseinschränkungen in AI/ML-Szenarien: Beim Anordnen von GPU-Trainingsaufgaben in Airflow müssen Sie sicherstellen, dass der Pod von KubernetesExecutor GPU-Ressourcen anfordern kann, und auf den Timeout-Wiederholungsmechanismus des Schedulers achten, der durch langfristige Trainingsaufgaben (>12 Stunden) ausgelöst werden kann. Es wird empfohlen, „execution_timeout“ und „retries=0“ für langfristige Trainingsaufgaben festzulegen, um zu verhindern, dass der Planer wiederholt neue Instanzen abruft, wenn das Training nicht abgeschlossen ist.
Versionsinfo
- Luftstrom 2.10 :Einen offiziellen genauen Termin gibt es noch nicht.
- Luftstrom 2.9 :Einen offiziellen genauen Termin gibt es noch nicht.
Benutzerbewertungen