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 Functions.
- 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.
In Managed Airflow (Gen 2) können Sie einen anderen Ansatz verwenden: DAGs mit Cloud Run-Funktionen und Pub/Sub-Nachrichten auslösen.
Die APIs für Ihr Projekt aktivieren
Console
Aktivieren Sie die Managed Airflow API und die Cloud Run Functions API.
Erforderliche Rollen zum Aktivieren von APIs
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 API und die Cloud Run Functions API:
Erforderliche Rollen zum Aktivieren von APIs
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
Für Airflow 2 ist die stabile REST API bereits standardmäßig aktiviert. Wenn die stabile API in Ihrer Umgebung deaktiviert ist, dann aktivieren Sie sie.
API-Aufrufe an die Airflow REST API über die Netzwerkzugriffssteuerung 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 All IP addresses have access (default), um Ihre Cloud Run-Funktionen nicht versehentlich zu blockieren. Sie können den Netzwerkzugriff für Webserver später jederzeit
konfigurieren.
URL des Airflow-Webservers abrufen
In diesem Beispiel werden REST API-Anfragen an den Airflow-Webserver-Endpunkt gesendet. Sie verwenden die URL des Airflow-Webservers 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.
Laden Sie einen DAG in Ihre Umgebung hoch.
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 Functions 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 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_dag_with_gcfals Einstiegspunkt an.(Java) Wenn Sie Code für dieses Beispiel hinzufügen, wählen Sie die Laufzeit Java 17 oder höher aus und geben Sie
functions.TriggerDagExampleals Einstiegspunkt an.
Anforderungen hinzufügen
Python
Geben Sie die Abhängigkeiten in der Datei requirements.txt an:
google-auth>=2.38.0
requests>=2.34.2
functions-framework==3.*
Java
Fügen Sie dem Abschnitt dependencies in der Datei pom.xml die folgenden Abhängigkeiten hinzu:
<dependency>
<groupId>com.google.apis</groupId>
<artifactId>google-api-services-docs</artifactId>
<version>v1-rev20250917-2.0.0</version>
</dependency>
<dependency>
<groupId>com.google.api-client</groupId>
<artifactId>google-api-client</artifactId>
<version>2.9.0</version>
</dependency>
<dependency>
<groupId>com.google.auth</groupId>
<artifactId>google-auth-library-credentials</artifactId>
<version>1.49.0</version>
</dependency>
<dependency>
<groupId>com.google.auth</groupId>
<artifactId>google-auth-library-oauth2-http</artifactId>
<version>1.49.0</version>
</dependency>
Funktionscode hinzufügen
Python
Fügen Sie der Datei main.py den folgenden Code hinzu:
Ersetzen Sie den Wert der Variablen
web_server_urldurch die Adresse des Airflow-Webservers, die Sie zuvor abgerufen haben.Wenn Sie einen anderen DAG auslösen, ersetzen Sie den Wert der Variablen
dag_id.
from __future__ import annotations
from typing import Any
from datetime import datetime, timezone
import google.auth
from google.auth.transport.requests import AuthorizedSession
import requests
import functions_framework
# Following Google Cloud best practices, these credentials should be
# constructed at start-up time and used throughout
# https://cloud.google.com/apis/docs/client-libraries-best-practices
AUTH_SCOPE = "https://www.googleapis.com/auth/cloud-platform"
CREDENTIALS, _ = google.auth.default(scopes=[AUTH_SCOPE])
def make_managed_airflow_web_server_request(
url: str, method: str = "GET", **kwargs: Any
) -> google.auth.transport.Response:
"""
Make a request to environment's web server.
Args:
url: The URL to fetch.
method: The request method to use ('GET', 'OPTIONS', 'HEAD', 'POST',
'PUT', 'PATCH', 'DELETE')
**kwargs: Any of the parameters defined for the request function:
https://github.com/requests/requests/blob/master/requests/api.py
If no timeout is provided, it is set to 90 by default.
"""
authed_session = AuthorizedSession(CREDENTIALS)
# Set the default timeout, if missing
if "timeout" not in kwargs:
kwargs["timeout"] = 90
return authed_session.request(method, url, **kwargs)
def trigger_dag_request(web_server_url: str, airflow_version: str, dag_id: str, data: dict, logical_date: str) -> str:
"""
Make a request to trigger a dag using the Airflow REST API.
https://airflow.apache.org/docs/apache-airflow/stable/stable-rest-api-ref.html
Args:
web_server_url: The URL of the Airflow web server.
airflow_version: Major version of Airflow. Determines the API endpoint.
dag_id: The DAG ID.
data: Additional configuration parameters for the DAG run (json).
logical_date: Data interval for which to run the DAG.
"""
if airflow_version == "2":
endpoint = f"api/v1/dags/{dag_id}/dagRuns"
elif airflow_version == "3":
endpoint = f"api/v2/dags/{dag_id}/dagRuns"
else:
raise ValueError(
f"Invalid Airflow version: {airflow_version}. Expected: 2 or 3.")
request_url = f"{web_server_url}/{endpoint}"
json_data = {
"conf": data,
"logical_date": logical_date,
}
response = make_managed_airflow_web_server_request(
request_url, method="POST", json=json_data
)
if response.status_code == 403:
raise requests.HTTPError(
"You do not have a permission to perform this operation. "
"Check Airflow RBAC roles for your account."
f"{response.headers} / {response.text}"
)
elif response.status_code != 200:
response.raise_for_status()
else:
return response.text
@functions_framework.cloud_event
def trigger_dag_with_gcf(cloud_event: CloudEvent) -> None:
"""
Entry point for the Cloud Function. Triggers a DAG and passes event data.
"""
# cloud_event.data contains the resource payload (e.g., storage object
# details or pub/sub body)
event_data = {
"id": cloud_event["id"],
"subject": cloud_event["subject"],
"type": cloud_event["type"],
"data": cloud_event.data
}
# TODO(developer): replace with your values
# Replace web_server_url with the Airflow web server address. To obtain this
# URL, run the following command for your environment:
# gcloud composer environments describe example-environment \
# --location=your-composer-region \
# --format="value(config.airflowUri)"
web_server_url = (
"https://example-airflow-ui-url-dot-us-central1.composer.googleusercontent.com"
)
# TODO(developer): If your environment uses Airflow 3, replace with "3"
airflow_major_version = "2"
# Replace with the ID of the DAG that you want to run.
dag_id = "composer_sample_trigger_response_dag"
# The data interval for which to run the DAG
# Format example: "2026-07-15T15:00:00Z"
now = datetime.now(timezone.utc)
logical_date = now.strftime("%Y-%m-%dT%H:%M:%SZ")
trigger_dag_request(web_server_url, airflow_major_version, dag_id, event_data, logical_date)
Java
Fügen Sie der Datei TriggerDagExample.java den folgenden Code hinzu (fügen Sie diese Datei in das Verzeichnis src/main/java/gcfv2/ ein):
Ersetzen Sie den Wert der Variablen
webServerUrldurch die Adresse des Airflow-Webservers, die Sie zuvor abgerufen haben.Wenn Sie einen anderen DAG auslösen, ersetzen Sie den Wert der Variablen
dagName.
package gcfv2;
import com.google.api.client.http.GenericUrl;
import com.google.api.client.http.HttpContent;
import com.google.api.client.http.HttpRequest;
import com.google.api.client.http.HttpRequestFactory;
import com.google.api.client.http.HttpResponse;
import com.google.api.client.http.HttpResponseException;
import com.google.api.client.http.javanet.NetHttpTransport;
import com.google.api.client.http.json.JsonHttpContent;
import com.google.api.client.json.gson.GsonFactory;
import com.google.auth.http.HttpCredentialsAdapter;
import com.google.auth.oauth2.GoogleCredentials;
import com.google.cloud.functions.CloudEventsFunction;
import com.google.gson.Gson;
import io.cloudevents.CloudEvent;
import java.nio.charset.StandardCharsets;
import java.time.Instant;
import java.util.logging.Logger;
import java.util.HashMap;
import java.util.Map;
/**
* Function that triggers an Airflow DAG in response to an event ad passes data.
*/
public class TriggerDagExample implements CloudEventsFunction {
private static final Logger logger = Logger.getLogger(TriggerDagExample.class.getName());
@Override
public void accept(CloudEvent event) throws Exception{
// TODO(developer): replace with your values
// Replace webServerUrl with the Airflow web server address. To obtain this
// URL, run the following command for your environment:
// gcloud composer environments describe example-environment \
// --location=your-composer-region \
// --format="value(config.airflowUri)"
String webServerUrl = "https://example-airflow-ui-url-dot-us-central1.composer.googleusercontent.com";
// TODO(developer): If your environment uses Airflow 3, replace with "3"
String majorAirflowVersion = "2";
String apiVersion = switch (majorAirflowVersion) {
case "2" -> "v1";
case "3" -> "v2";
default -> throw new IllegalArgumentException("Invalid Airflow version: " + majorAirflowVersion);
};
String dagName = "composer_sample_trigger_response_dag";
String url = String.format("%s/api/%s/dags/%s/dagRuns", webServerUrl, apiVersion, dagName);
logger.info(String.format("Triggering DAG %s as a result of an event on the object %s.",
dagName, event.getSubject()));
logger.info(String.format("Triggering DAG through the following URL: %s", url));
GoogleCredentials googleCredentials = GoogleCredentials.getApplicationDefault()
.createScoped("https://www.googleapis.com/auth/cloud-platform");
HttpCredentialsAdapter credentialsAdapter = new HttpCredentialsAdapter(googleCredentials);
HttpRequestFactory requestFactory =
new NetHttpTransport().createRequestFactory(credentialsAdapter);
Map<String, Object> conf = new HashMap<>();
conf.put("id", event.getId());
conf.put("subject", event.getSubject());
conf.put("type", event.getType());
if (event.getData() != null) {
String dataJson = new String(event.getData().toBytes(), StandardCharsets.UTF_8);
Gson gson = new Gson();
Map<String, Object> dataMap = gson.fromJson(dataJson, Map.class);
conf.put("data", dataMap);
}
String currentUtcTime = Instant.now().toString();
Map<String, Object> json = new HashMap<>();
json.put("conf", conf);
json.put("logical_date", currentUtcTime);
HttpContent content = new JsonHttpContent(new GsonFactory(), json);
HttpRequest request = requestFactory.buildPostRequest(new GenericUrl(url), content);
request.getHeaders().setContentType("application/json");
HttpResponse response = null;
try {
response = request.execute();
int statusCode = response.getStatusCode();
logger.info("Response code: " + statusCode);
logger.info(response.parseAsString());
} catch (HttpResponseException e) {
logger.info("Received HTTP exception");
logger.info(e.getLocalizedMessage());
logger.info("- 400 error: wrong arguments passed to Airflow API");
logger.info("- 401 error: check if service account has Composer User role");
logger.info("- 403 error: check Airflow RBAC roles assigned to service account");
logger.info("- 404 error: check Web Server URL");
} catch (Exception e) {
logger.info("Received exception");
logger.info(e.getLocalizedMessage());
} finally {
// Safely close and release the HTTP connection pool resource
if (response != null) {
try {
response.disconnect();
} catch (Exception e) {
logger.warning("Failed to disconnect response: " + e.getMessage());
}
}
}
}
}
Funktion testen
So prüfen Sie, ob Ihre Funktion und Ihr DAG wie erwartet 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.
- 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:
Beispielbefehl zum Testen der Funktion:
curl -X POST "https://service-id.region.run.app" \
-H "Authorization: bearer $(gcloud auth print-identity-token)" \
-X POST \
-H "Content-Type: application/json" \
-H "ce-id: 1234567890" \
-H "ce-specversion: 1.0" \
-H "ce-type: google.cloud.storage.object.v1.finalized" \
-H "ce-source: //storage.googleapis.com/projects/_/buckets/example-bucket" \
-d '{
"name": "example-file.csv",
"bucket": "example-bucket"
}'
Beispielausgabe:
[2026-07-14, 15:10:12 UTC] {subprocess.py:88} INFO - Running command: ['/usr/bin/bash', '-c', "echo {'data': {'name': 'example-file.csv', 'bucket': 'example-bucket'}, 'id': '1234567890', 'type': 'google.cloud.storage.object.v1.finalized'}"]
[2026-07-14, 15:10:12 UTC] {subprocess.py:99} INFO - Output:
[2026-07-14, 15:10:12 UTC] {subprocess.py:106} INFO - {data: {name: example-file.csv, bucket: my-bucket}, id: 1234567890, type: google.cloud.storage.object.v1.finalized}
[2026-07-14, 15:10:12 UTC] {subprocess.py:110} INFO - Command exited with return code 0
[2026-07-15, 10:06:32 UTC] {subprocess.py:88} INFO - Running command: ['/usr/bin/bash', '-c', "echo {'id': '1234567890', 'subject': 'objects/example-file.csv', 'type': 'google.cloud.storage.object.v1.finalized', 'data': {'name': 'example-file.csv', 'bucket': 'example-bucket'}}"]
[2026-07-15, 10:06:32 UTC] {subprocess.py:99} INFO - Output:
[2026-07-15, 10:06:32 UTC] {subprocess.py:106} INFO - {id: 1234567890, subject: objects/example-file.csv, type: google.cloud.storage.object.v1.finalized, data: {name: example-file.csv, bucket: example-bucket}}
[2026-07-15, 10:06:32 UTC] {subprocess.py:110} INFO - Command exited with return code 0
Fehlerbehebung:
- Wenn Ihre Funktion mit dem Fehler
NullPointerException: Null datafehlschlägt und der Stacktrace auf die FunktionBackgroundFunctionExecutor.parseLegacyEventverweist, bedeutet das, dass die von der Funktion empfangene Ereignis keine Standard-Metadatenheader vom TypCloudEventhat. Die Funktion geht davon aus, dass Sie ein altes Hintergrundereignis senden, versucht, das Felddatadaraus zu parsen, und schlägt fehl. Das kann beispielsweise passieren, wenn Sie beim Testen der Funktion eine beliebige Ereignisnutzlast senden. - Wenn Ihre Funktion mit
500 Internal Server Error: The server encountered an internal error and was unable to complete your request.fehlschlägt, prüfen Sie den Wert der Variablenairflow_major_version. Diese Variable bestimmt den Airflow REST API-Endpunkt, der sich in Airflow 2 und Airflow 3 unterscheidet.
Nächste Schritte
- Auf die Airflow-UI zugreifen
- Auf die Airflow REST API zugreifen
- DAGs schreiben
- Cloud Run-Funktionen schreiben
- Google Cloud Storage-Trigger