Managed Service for Apache Airflow-DAGs mit Cloud Run-Funktionen und der Airflow REST API auslösen

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:

Das Beispiel in diesem Leitfaden zeigt eine Funktion, die als Reaktion auf ein Ereignis einen DAG auslöst:

  1. Sie konfigurieren Trigger für Ihre Funktion in Cloud Run Functions.
  2. 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.
  3. 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.

In Managed Airflow (2. Generation) können Sie einen anderen Ansatz verwenden: DAGs mit Cloud Run Functions und Pub/Sub-Nachrichten auslösen.

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

APIs aktivieren

gcloud

Aktivieren Sie die Managed Airflow API und die Cloud Run Functions API, 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.com composer.googleapis.com

Airflow REST API aktivieren

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.

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 die URL des Airflow-Webservers in Ihrem Cloud Function-Code.

Console

  1. Rufen Sie in der Google Cloud Console die Seite Umgebungen auf.

    Zur Seite Umgebungen

  2. Klicken Sie auf den Namen Ihrer Umgebung.

  3. Rufen Sie auf der Seite Umgebungsdetails den Tab Umgebungskonfiguration auf.

  4. 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_NAME durch den Namen der Umgebung.
  • LOCATION durch die Region, in der sich die Umgebung befindet.

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_dag_with_gcf als Einstiegspunkt an.

    • (Java) Wenn Sie Code für dieses Beispiel hinzufügen, wählen Sie die Java 17-Laufzeit aus und geben Sie functions.TriggerDagExample als 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 dependencies-Abschnitt im 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 den folgenden Code in die main.py-Datei ein:

  • Ersetzen Sie den Wert der Variablen web_server_url durch die Airflow-Webserveradresse, 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 den folgenden Code in die Datei TriggerDagExample.java ein (legen Sie diese Datei im Verzeichnis src/main/java/gcfv2/ ab):

  • Ersetzen Sie den Wert der Variablen webServerUrl durch die Airflow-Webserveradresse, 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 vorgesehen funktionieren:

  1. Warten Sie, bis die Funktion bereitgestellt wurde.
  2. 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.
  3. Sehen Sie sich die DAG-Seite in der Airflow-Weboberfläche an. Der DAG sollte einen aktiven oder bereits abgeschlossenen DAG-Lauf haben.
  4. 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_info in die Logs ausgegeben werden:

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 einem NullPointerException: Null data-Fehler fehlschlägt und der Stacktrace auf die Funktion BackgroundFunctionExecutor.parseLegacyEvent verweist, bedeutet das, dass das von der Funktion empfangene Ereignis keine Standard-CloudEvent-Metadatenheader hat. Die Funktion geht davon aus, dass Sie ein altes Hintergrundereignis senden, versucht, das Feld data daraus 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 Variablen airflow_major_version. Diese Variable bestimmt den Airflow REST API-Endpunkt, der sich in Airflow 2 und Airflow 3 unterscheidet.

Nächste Schritte