Managed Airflow (Gen 3) | Managed Airflow (Gen 2) | Managed Airflow (Legacy Gen 1)
Auf dieser Seite wird beschrieben, wie Sie Cloud Run-Funktionen verwenden, um Managed Service for Apache Airflow-DAGs als Reaktion auf Ereignisse auszulösen.
Apache Airflow ist für die regelmäßige Ausführung von DAGs konzipiert. Sie können DAGs aber auch als Reaktion auf Ereignisse auslösen. Eine Möglichkeit hierfür ist die Verwendung von Cloud Run Functions, um Managed Airflow-DAGs auszulösen, wenn ein bestimmtes Ereignis eintritt.
Weitere Möglichkeiten:
- DAGs nur mit der Airflow REST API auslösen:
- Erstellen Sie eine Funktion, die einen DAG auslöst, wenn eine Nachricht in ein Pub/Sub-Thema übertragen wird.
Das Beispiel in diesem Leitfaden zeigt eine Funktion, die als Reaktion auf ein Ereignis einen DAG auslöst:
- Sie konfigurieren Trigger für Ihre Funktion in Cloud Run Functions.
- Wenn die Funktion ausgelöst wird, wird eine Anfrage gesendet, um einen DAG über die Airflow REST API Ihrer Managed Airflow-Umgebung auszulösen. Die Anfrage enthält die Kennung 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 die Vorbereitungsschritte aufgeführt.
Netzwerkkonfiguration Ihrer Umgebung prüfen
Diese Lösung funktioniert nicht in Konfigurationen mit privaten IP-Adressen 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- und Cloud Run Functions APIs, falls sie noch nicht aktiviert sind.
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 „Service Usage-Administrator“ (roles/serviceusage.serviceUsageAdmin) erhalten. Informationen zum Zuweisen von Rollen
gcloud
Aktivieren Sie die Managed Airflow und Cloud Run Functions APIs, falls sie noch nicht aktiviert sind:
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 „Service Usage-Administrator“ (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 standardmäßig aktiviert. Wenn die stabile API in Ihrer Umgebung deaktiviert ist, 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 Netzwerkzugriffssteuerung des Webservers zulassen
Cloud Run Functions können die Airflow REST API über eine IPv4- oder IPv6-Adresse erreichen.
Wenn Sie sich nicht sicher sind, welcher IP-Bereich für Aufrufe verwendet wird, verwenden Sie in der Webserver-Zugriffssteuerung die Standardkonfigurationsoption All IP addresses have access (default), um Ihre Cloud Run-Funktionen nicht versehentlich zu blockieren. Sie können den Netzwerkzugriff des Webservers später jederzeit konfigurieren.
URL des Airflow-Webservers abrufen
In diesem Beispiel werden REST API-Anfragen an den Airflow-Webserver-Endpunkt gesendet.
Sie verwenden den Teil der Airflow-Weboberflächen-URL 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 wird unter 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. Stattdessen können Sie eine nicht authentifizierte Anfrage an den Airflow-Webserver senden und die Client-ID aus der Weiterleitungs-URL erfassen:
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 die Werte für project_id, location und composer_environment ein und führen Sie den Code in Cloud Shell oder in 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 der von Ihnen bevorzugten Sprache bereitstellen, die von Cloud Run-Funktionen oder Cloud Run unterstützt wird. In dieser Anleitung wird eine Cloud Functions-Funktion in Python und Java implementiert.
Konfigurationsparameter für Funktionen 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 Dienstkonto, das Sie für den Trigger angeben, muss über genügend Berechtigungen verfügen, um DAGs in Managed Airflow-Umgebungen auszulösen.
Wir empfehlen, dem Prinzip der geringsten Berechtigung zu folgen 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 Änderungen vor:
Ersetzen Sie den Wert der Variablen
client_iddurch den Wertclient_id, den Sie zuvor abgerufen haben.Ersetzen Sie den Wert der Variable
webserver_iddurch die Mandantenprojekt-ID, die Teil der URL der Airflow-Weboberfläche vor.appspot.comist. Sie haben die URL der Airflow-Weboberfläche bereits 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 vorgesehen funktionieren:
- Warten Sie, bis die 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 für sie auswählen.
- Sehen Sie sich die DAG-Seite in der Airflow-Weboberfläche an. Der DAG sollte einen aktiven oder bereits abgeschlossenen DAG-Lauf haben.
- Sehen Sie sich in der Airflow-UI die Aufgabenlogs für diesen Lauf an. Sie sollten sehen, dass die Daten, die von der Funktion empfangen wurden, von der Aufgabe
print_gcs_infoin die Logs ausgegeben werden:
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