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 :
- déclencher des DAG à l'aide de l'API REST Airflow.
- créer une fonction qui déclenche un DAG lorsqu'un message est envoyé à un sujet Pub/Sub.
L'exemple de ce guide illustre une fonction qui déclenche un DAG en réponse à un événement :
- Vous configurez des déclencheurs pour votre fonction dans Cloud Run Functions.
- 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.
- 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 (2e 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.
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.comcomposer.googleapis.com
Activer l'API REST Airflow
Pour Airflow 2, l'API REST stable est déjà activée par défaut. Si l'API stable est désactivée dans votre environnement, alors activez l'API REST stable.
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
Dans la Google Cloud console, accédez à la page Environnements.
Cliquez sur le nom de votre environnement.
Sur la page Détails de l'environnement, accédez à l'onglet Configuration de l'environnement.
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_NAMEpar le nom de l'environnement.LOCATIONpar 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.
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 section 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_gcfcomme 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.TriggerDagExamplecomme 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_urlpar l'adresse du serveur Web Airflow obtenue précédemment.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
webServerUrlpar l'adresse du serveur Web Airflow obtenue précédemment.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 :
- Attendez que votre fonction soit déployée.
- 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.
- 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.
- 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 dataet que la trace de la pile pointe vers la fonctionBackgroundFunctionExecutor.parseLegacyEvent, cela signifie que l'événement reçu par la fonction ne comporte pas d'en-têtes de métadonnéesCloudEventstandards. La fonction suppose que vous envoyez un événement d'arrière-plan hérité, tente d'analyser le champdataet é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 variableairflow_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
- Accéder à l'UI d'Airflow
- Accéder à l'API REST Airflow
- Écrire des DAG
- Écrire des fonctions Cloud Run
- Déclencheurs Cloud Storage