Managed Airflow (Gen 3) | Managed Airflow (Gen 2) | Managed Airflow (Legacy Gen 1)
Auf dieser Seite wird beschrieben, wie Sie mit Cloud Run-Funktionen Managed Service for Apache Airflow-DAGs als Reaktion auf Ereignisse auslösen.
Apache Airflow ist für das regelmäßige Ausführen von DAGs nach Plan konzipiert. Sie können DAGs aber auch als Reaktion auf Ereignisse auslösen lassen. Eine Möglichkeit dazu ist, Cloud Run-Funktionen zu verwenden, um Managed Airflow-DAGs auszulösen, wenn ein bestimmtes Ereignis eintritt.
Weitere Möglichkeiten:
- DAGs nur mit der Airflow REST API auslösen.
- Eine Funktion erstellen, die einen DAG auslöst, wenn eine Nachricht an ein Pub/Sub-Thema gesendet wird.
Das Beispiel in dieser Anleitung zeigt eine Funktion, die einen DAG als Reaktion auf ein Ereignis auslöst:
- Sie konfigurieren Trigger für Ihre Funktion in Cloud Run-Funktionen.
- Wenn die Funktion ausgelöst wird, sendet sie eine Anfrage zum Auslösen eines DAG über die Airflow REST API Ihrer Managed Airflow-Umgebung. Die Anfrage enthält die ID und den Typ des Ereignisses sowie die Nutzlast des Ereignisses.
- Airflow verarbeitet diese Anfrage und führt den in der Anfrage angegebenen DAG aus. Der DAG gibt die Daten aus, die von der Funktion an ihn übergeben wurden.
Hinweis
In diesem Abschnitt werden vorbereitende Schritte aufgeführt.
Netzwerkkonfiguration Ihrer Umgebung prüfen
Diese Lösung funktioniert nicht in Konfigurationen mit privater IP-Adresse und VPC Service Controls, da es in diesen Konfigurationen nicht möglich ist, die Verbindung von Cloud Run-Funktionen zum Airflow-Webserver zu konfigurieren.
Die APIs für Ihr Projekt aktivieren
Console
Aktivieren Sie die Managed Airflow API und die Cloud Run Functions API.
Rollen, die zum Aktivieren von APIs erforderlich sind
Zum Aktivieren von APIs benötigen Sie die Berechtigung serviceusage.services.enable. Wenn Sie das Projekt erstellt haben, haben Sie diese Berechtigung wahrscheinlich bereits über die Rolle „Inhaber“ (roles/owner). Andernfalls können Sie diese Berechtigung über die Rolle „Administrator der Dienstnutzung“ (roles/serviceusage.serviceUsageAdmin) erhalten. Informationen zum Zuweisen von Rollen.
gcloud
Aktivieren Sie die Managed Airflow API und die Cloud Run Functions API:
Rollen, die zum Aktivieren von APIs erforderlich sind
Zum Aktivieren von APIs benötigen Sie die Berechtigung serviceusage.services.enable. Wenn Sie das Projekt erstellt haben, haben Sie diese Berechtigung wahrscheinlich bereits über die Rolle „Inhaber“ (roles/owner). Andernfalls können Sie diese Berechtigung über die Rolle „Administrator der Dienstnutzung“ (roles/serviceusage.serviceUsageAdmin) erhalten. Informationen zum Zuweisen von Rollen.
gcloud services enable cloudfunctions.googleapis.comcomposer.googleapis.com
Airflow REST API aktivieren
Je nach Airflow-Version:
- In Airflow 2 ist die stabile REST API bereits standardmäßig aktiviert. Wenn die stabile API in Ihrer Umgebung deaktiviert ist, dann aktivieren Sie die stabile REST API.
- Aktivieren Sie für Airflow 1, die experimentelle REST API.
API-Aufrufe an die Airflow REST API über die Netzwerkkontrolle des Webservers zulassen
Cloud Run-Funktionen können die Airflow REST API über eine IPv4- oder IPv6-Adresse erreichen.
Wenn Sie nicht sicher sind, welcher aufrufende IP-Bereich verwendet wird, verwenden Sie die Standardkonfigurationsoption in der Webserver-Zugriffssteuerung , nämlich All IP addresses have access (default) (Alle IP-Adressen haben Zugriff (Standard)), um zu verhindern, dass Ihre Cloud Run-Funktionen versehentlich blockiert werden. Sie können
den Netzwerkzugriff auf den Webserver später konfigurieren.
URL des Airflow-Webservers abrufen
In diesem Beispiel werden REST API-Anfragen an den Airflow-Webserver-Endpunkt gesendet.
Verwenden Sie den Teil der URL der Airflow-Weboberfläche vor .appspot.com in Ihrem Cloud Function-Code.
Console
Rufen Sie in der Google Cloud Console die Seite Umgebungen auf.
Klicken Sie auf den Namen Ihrer Umgebung.
Rufen Sie auf der Seite Umgebungsdetails den Tab Umgebungskonfiguration auf.
Die URL des Airflow-Webservers ist im Element Airflow-Web-UI aufgeführt.
gcloud
Führen Sie dazu diesen Befehl aus:
gcloud composer environments describe ENVIRONMENT_NAME \
--location LOCATION \
--format='value(config.airflowUri)'
Ersetzen Sie:
ENVIRONMENT_NAMEdurch den Namen der Umgebung.LOCATIONdurch die Region, in der sich die Umgebung befindet.
client_id des IAM-Proxys abrufen
Um eine Anfrage an den Airflow REST API-Endpunkt zu senden, benötigt die Funktion die Client-ID des Identity and Access Management-Proxys, der den Airflow-Webserver schützt.
Managed Airflow stellt diese Informationen nicht direkt zur Verfügung. Senden Sie stattdessen eine nicht authentifizierte Anfrage an den Airflow-Webserver und erfassen Sie die Client-ID aus der Weiterleitungs-URL:
cURL
curl -v AIRFLOW_URL 2>&1 >/dev/null | grep -o "client_id\=[A-Za-z0-9-]*\.apps\.googleusercontent\.com"
Ersetzen Sie AIRFLOW_URL durch die URL der Airflow-Weboberfläche.
Suchen Sie in der Ausgabe nach dem String nach client_id. Beispiel:
client_id=836436932391-16q2c5f5dcsfnel77va9bvf4j280t35c.apps.googleusercontent.com
Python
Speichern Sie den folgenden Code in einer Datei mit dem Namen get_client_id.py. Geben Sie Ihre Werte für project_id, location und composer_environment ein und führen Sie den Code in Cloud Shell oder Ihrer lokalen Umgebung aus.
DAG in Ihre Umgebung hochladen
Laden Sie einen DAG in Ihre Umgebung hoch. Der folgende Beispiel-DAG gibt die empfangene DAG-Ausführungskonfiguration aus. Sie lösen diesen DAG über eine Funktion aus, die Sie später in dieser Anleitung erstellen.
import datetime
import airflow
from airflow.operators.bash_operator import BashOperator
with airflow.DAG(
'composer_sample_trigger_response_dag',
start_date=datetime.datetime(2026, 1, 1),
# Not scheduled, trigger only
schedule=None) as dag:
# Print the dag_run's configuration, which includes information about the
# Cloud Storage object change.
print_gcs_info = BashOperator(
task_id='print_gcs_info', bash_command='echo {{ dag_run.conf }}}}')
Funktion bereitstellen, die den DAG auslöst
Sie können eine Funktion in Ihrer bevorzugten Sprache bereitstellen, die von Cloud Run-Funktionen oder Cloud Run unterstützt wird. In dieser Anleitung wird eine Cloud Functions-Funktion gezeigt, die in Python und Java implementiert ist.
Konfigurationsparameter für die Funktion angeben
Trigger: Wählen Sie einen oder mehrere Eventarc Trigger für Ihre Funktion aus.
Weitere Informationen zum Erstellen von Triggern finden Sie unter Trigger mit Eventarc erstellen. Sie können beispielsweise Funktionen aus Cloud Storage mit Eventarc auslösen.
Dienstkonto: Das für den Trigger angegebene Dienstkonto muss über ausreichende Berechtigungen verfügen, um DAGs in Managed Airflow Umgebungen auszulösen.
Wir empfehlen, das Prinzip der geringsten Berechtigung zu befolgen und dem Dienstkonto nur die Rolle Composer-Nutzer (
composer.user) zuzuweisen. Weitere Informationen zum Konfigurieren von Berechtigungen finden Sie unter Rollen und Berechtigungen für Cloud Run-Ziele.Funktionseinstiegspunkt:
(Python) Wenn Sie Code für dieses Beispiel hinzufügen, wählen Sie die Laufzeit Python 3.10 oder höher aus und geben Sie
trigger_dagals Einstiegspunkt an.
Anforderungen hinzufügen
Geben Sie die Abhängigkeiten in der Datei requirements.txt an:
Funktionscode hinzufügen
Fügen Sie den folgenden Code in die Datei main.py ein und nehmen Sie die folgenden Ersetzungen vor:
Ersetzen Sie den Wert der Variablen
client_iddurch den Wert vonclient_id, den Sie zuvor abgerufen haben.Ersetzen Sie den Wert der Variablen
webserver_iddurch die Mandantenprojekt-ID, die Teil der URL der Airflow-Weboberfläche vor.appspot.comist. Sie haben die URL der Airflow-Weboberfläche zuvor abgerufen.Geben Sie die von Ihnen verwendete Airflow REST API-Version an:
- Wenn Sie die stabile Airflow API verwenden, legen Sie die Variable
USE_EXPERIMENTAL_APIaufFalsefest. - Wenn Sie die experimentelle Airflow REST API verwenden, sind keine Änderungen erforderlich. Die Variable
USE_EXPERIMENTAL_APIist bereits aufTruefestgelegt.
- Wenn Sie die stabile Airflow API verwenden, legen Sie die Variable
Funktion testen
So prüfen Sie, ob Ihre Funktion und Ihr DAG wie erwartet funktionieren:
- Warten Sie, bis Ihre Funktion bereitgestellt wurde.
- Lösen Sie die Funktion gemäß dem angegebenen Trigger aus. Sie können die Funktion auch manuell auslösen, indem Sie in der Google Cloud Console die Aktion Funktion testen auswählen.
- Prüfen Sie die DAG-Seite in der Airflow-Weboberfläche. Der DAG sollte eine aktive oder bereits abgeschlossene DAG-Ausführung haben.
- Prüfen Sie in der Airflow-UI die Aufgabenlogs für diese Ausführung. Sie sollten sehen, dass die Aufgabe
print_gcs_infodie von der Funktion empfangenen Daten in die Logs ausgibt:
Beispielausgabe:
[2021-04-04 18:25:44,778] {bash_operator.py:154} INFO - Output:
[2021-04-04 18:25:44,781] {bash_operator.py:158} INFO - Triggered from GCF:
{bucket: example-storage-for-gcf-triggers, contentType: text/plain,
crc32c: dldNmg==, etag: COW+26Sb5e8CEAE=, generation: 1617560727904101,
... }
[2021-04-04 18:25:44,781] {bash_operator.py:162} INFO - Command exited with
return code 0h
Nächste Schritte
- Auf die Airflow-UI zugreifen
- Auf die Airflow REST API zugreifen
- DAGs schreiben
- Cloud Run-Funktionen schreiben
- Google Cloud Storage-Trigger