Déclencher des DAG Managed Service pour Apache Airflow avec des fonctions Cloud Run et l'API REST Airflow

Airflow géré (3e génération) | Airflow géré (2e génération) | Airflow géré (1re génération héritée)

Cette page explique comment utiliser Cloud Run Functions pour déclencher des DAG Managed Service for Apache Airflow en réponse à des événements.

Apache Airflow est conçu pour exécuter des DAG de façon régulière, mais vous pouvez également déclencher des DAG en réponse à des événements. Pour ce faire, vous pouvez utiliser Cloud Run Functions afin de déclencher des DAG Managed Airflow lorsqu'un événement spécifique se produit.

Vous pouvez également :

L'exemple de ce guide présente une fonction qui déclenche un DAG en réponse à un événement :

  1. Vous configurez des déclencheurs pour votre fonction dans Cloud Run Functions.
  2. Lorsque la fonction est déclenchée, elle envoie une requête pour déclencher un DAG via l'API REST Airflow de votre environnement Managed Airflow. La requête contient l'identifiant et le type de l'événement, ainsi que sa charge utile.
  3. Airflow traite cette requête et exécute le DAG spécifié dans la requête. Le DAG génère les données qui lui ont été transmises par la fonction.

Avant de commencer

Cette section répertorie les étapes préparatoires.

Vérifier la configuration réseau de votre environnement

Cette solution ne fonctionne pas dans les configurations d'adresse IP privée et de VPC Service Controls, car il n'est pas possible de configurer la connectivité entre Cloud Run Functions et le serveur Web Airflow dans ces configurations.

Dans Managed Airflow (3e génération), vous pouvez utiliser une autre approche : déclencher des DAG à l'aide de Cloud Run Functions et de messages Pub/Sub.

Activer les API pour votre projet.

Console

Activez les API Managed Airflow et Cloud Run Functions.

Rôles requis pour activer les API

Pour activer les API, vous devez disposer de l'autorisation serviceusage.services.enable. Si vous avez créé le projet, vous disposez probablement déjà de cette autorisation via le rôle Propriétaire (roles/owner). Sinon, vous pouvez l'obtenir via le rôle Administrateur d'utilisation du service (roles/serviceusage.serviceUsageAdmin). Découvrez comment attribuer des rôles.

Activer les API

gcloud

Activez les API Managed Airflow et Cloud Run Functions :

Rôles requis pour activer les API

Pour activer les API, vous devez disposer de l'autorisation serviceusage.services.enable. Si vous avez créé le projet, vous disposez probablement déjà de cette autorisation via le rôle Propriétaire (roles/owner). Sinon, vous pouvez l'obtenir via le rôle Administrateur d'utilisation du service (roles/serviceusage.serviceUsageAdmin). Découvrez comment attribuer des rôles.

gcloud services enable cloudfunctions.googleapis.com composer.googleapis.com

Autoriser les appels d'API à l'API REST Airflow à l'aide du contrôle des accès réseau du serveur Web

Cloud Run Functions peut accéder à l'API REST Airflow via une adresse IPv4 ou IPv6.

Si vous n'êtes pas sûr de la plage d'adresses IP appelantes, utilisez une option de configuration par défaut dans Contrôle d'accès au serveur Web , à savoir All IP addresses have access (default) (Toutes les adresses IP ont accès (par défaut)), afin de ne pas bloquer accidentellement vos fonctions Cloud Run. Vous pourrez toujours configurer l'accès réseau au serveur Web ultérieurement.

Obtenir l'URL du serveur Web Airflow

Cet exemple envoie des requêtes API REST au point de terminaison du serveur Web Airflow. Vous utilisez l'URL du serveur Web Airflow dans le code de votre fonction Cloud.

Console

  1. Dans la Google Cloud console, accédez à la page Environnements.

    Accéder à la page Environnements

  2. Cliquez sur le nom de votre environnement.

  3. Sur la page Détails de l'environnement, accédez à l'onglet Configuration de l'environnement.

  4. L'URL du serveur Web Airflow est répertoriée dans l'élément UI Web Airflow.

gcloud

Exécutez la commande suivante :

gcloud composer environments describe ENVIRONMENT_NAME \
    --location LOCATION \
    --format='value(config.airflowUri)'

Remplacez :

  • ENVIRONMENT_NAME par le nom de l'environnement.
  • LOCATION par la région dans laquelle se trouve l'environnement.

Importer un DAG dans votre environnement

Importer un DAG dans votre environnement. L'exemple de DAG suivant génère la configuration d'exécution de DAG reçue. Vous allez déclencher ce DAG à partir d'une fonction que vous créerez plus loin dans ce guide.

Airflow 3

import datetime

import airflow
from airflow.providers.standard.operators.bash 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 }}}}')

Airflow 2

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 }}}}')

Déployer une fonction qui déclenche le DAG

Vous pouvez déployer une fonction à l'aide du langage de votre choix compatible avec Cloud Run Functions ou Cloud Run. Ce tutoriel présente une fonction Cloud implémentée en Python et Java.

Spécifier les paramètres de configuration de la fonction

  • Déclencheur : sélectionnez un ou plusieurs déclencheurs Eventarc pour votre fonction.

    Pour en savoir plus sur la création de déclencheurs, consultez la section Créer des déclencheurs avec Eventarc. Par exemple, vous pouvez déclencher des fonctions depuis Cloud Storage à l'aide d'Eventarc.

  • Compte de service : le compte de service que vous spécifiez pour le déclencheur doit disposer des autorisations suffisantes pour déclencher des DAG dans les environnements Managed Airflow.

    Nous vous recommandons de suivre le principe de moindre privilège et de n'accorder que le rôle Utilisateur Composer (composer.user). Pour en savoir plus sur la configuration des autorisations, consultez la page Rôles et autorisations pour les cibles Cloud Run.

  • Point d'entrée de la fonction :

    • (Python) Lorsque vous ajoutez du code pour cet exemple, sélectionnez l'environnement d'exécution Python 3.10 ou une version ultérieure, puis spécifiez trigger_dag_with_gcf comme point d'entrée.

    • (Java) Lorsque vous ajoutez du code pour cet exemple, sélectionnez l'environnement d'exécution Java 17 ou une version ultérieure, puis spécifiez functions.TriggerDagExample comme point d'entrée.

Ajouter des conditions

Python

Spécifiez les dépendances dans le fichier requirements.txt :

google-auth>=2.38.0
requests>=2.34.2
functions-framework==3.*

Java

Ajoutez les dépendances suivantes à la section dependencies du fichier pom.xml :

    <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>

Ajouter un code de fonction

Python

Placez le code suivant dans le fichier main.py :

  • Remplacez la valeur de la variable web_server_url par l'adresse du serveur Web Airflow obtenue précédemment.

  • (Airflow 3) Remplacez la valeur de la variable airflow_major_version par 3, qui correspond à la version majeure d'Airflow dans votre environnement.

  • Si vous déclenchez un autre DAG, remplacez la valeur de la variable 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

Placez le code suivant dans le fichier TriggerDagExample.java (placez ce fichier dans le répertoire src/main/java/gcfv2/) :

  • Remplacez la valeur de la variable webServerUrl par l'adresse du serveur Web Airflow obtenue précédemment.

  • (Airflow 3) Remplacez la valeur de la variable majorAirflowVersion par 3, qui correspond à la version majeure d'Airflow dans votre environnement.

  • Si vous déclenchez un autre DAG, remplacez la valeur de la variable 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());
        }
      }
    }
  }
}

Tester votre fonction

Pour vérifier que votre fonction et votre DAG fonctionnent comme prévu, procédez comme suit :

  1. Attendez que votre fonction soit déployée.
  2. Déclenchez la fonction en fonction du déclencheur spécifié. Vous pouvez également déclencher la fonction manuellement en sélectionnant l'action Tester la fonction dans la Google Cloud console.
  3. Consultez la page du DAG dans l'interface Web Airflow. Le DAG doit comporter une exécution de DAG active ou déjà terminée.
  4. Dans l'UI d'Airflow, consultez les journaux des tâches associées à cette exécution. Vous devriez voir que la tâche print_gcs_info écrit les données reçues de la fonction dans les journaux :

Exemple de commande pour tester la fonction :

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"
  }'

Exemple de résultat :

[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

Dépannage :

  • Si votre fonction échoue avec une erreur NullPointerException: Null data et que la trace de la pile pointe vers la fonction BackgroundFunctionExecutor.parseLegacyEvent, cela signifie que l'événement reçu par la fonction ne comporte pas d'en-têtes de métadonnées CloudEvent standards. La fonction suppose que vous envoyez un événement d'arrière-plan hérité, tente d'analyser le champ data et échoue. Cela peut se produire, par exemple, si vous envoyez une charge utile d'événement arbitraire lorsque vous testez la fonction.
  • Si votre fonction échoue avec l'erreur 500 Internal Server Error: The server encountered an internal error and was unable to complete your request., vérifiez la valeur de la variable airflow_major_version. Cette variable détermine le point de terminaison de l'API REST Airflow, qui est différent dans Airflow 2 et Airflow 3.

Étape suivante