Airflow gerenciado (Geração 3) | Airflow gerenciado (Geração 2) | Airflow gerenciado (Geração 1 legada)
Nesta página, descrevemos como usar funções do Cloud Run para acionar DAGs do Serviço gerenciado para Apache Airflow em resposta a eventos.
O Apache Airflow foi projetado para executar DAGs em uma programação regular, mas também é possível acionar DAGs em resposta a eventos. Uma maneira de fazer isso é usar as funções do Cloud Run para acionar DAGs gerenciados do Airflow quando um evento especificado ocorrer.
Você também pode:
- Acione DAGs usando apenas a API REST do Airflow.
- Crie uma função que aciona um DAG quando uma mensagem é enviada por push para um tópico do Pub/Sub.
O exemplo neste guia demonstra uma função que aciona um DAG em resposta a um evento:
- Você configura acionadores para sua função no Cloud Run functions.
- Quando a função é acionada, ela faz uma solicitação para acionar um DAG pela API REST do Airflow do seu ambiente do Airflow Gerenciado. A solicitação contém o identificador e o tipo do evento, além da carga útil dele.
- O Airflow processa essa solicitação e executa o DAG especificado nela. O DAG gera os dados que foram transmitidos a ele pela função.
Antes de começar
Esta seção lista as etapas preparatórias.
Verificar a configuração de rede do ambiente
Essa solução não funciona em configurações de IP privado e VPC Service Controls porque não é possível configurar a conectividade das funções do Cloud Run com o servidor da Web do Airflow nessas configurações.
No Airflow Gerenciado (Geração 2), você pode usar outra abordagem: acionar DAGs usando funções do Cloud Run e mensagens do Pub/Sub.
Ativar as APIs do projeto
Console
Ative as APIs do Airflow gerenciado e do Cloud Run functions, se alguma delas ainda não estiver ativada.
Funções necessárias para ativar APIs
Para ativar APIs, você precisa da permissão serviceusage.services.enable. Se você
criou o projeto, provavelmente já tem essa permissão com o papel de
Proprietário (roles/owner). Caso contrário, é possível receber essa permissão com o papel de
Administrador do Service Usage (roles/serviceusage.serviceUsageAdmin).
Saiba como conceder papéis.
gcloud
Ative as APIs do Airflow gerenciado e do Cloud Run functions, se alguma delas ainda não estiver ativada:
Funções necessárias para ativar APIs
Para ativar APIs, você precisa da permissão serviceusage.services.enable. Se você
criou o projeto, provavelmente já tem essa permissão com o papel de
Proprietário (roles/owner). Caso contrário, é possível receber essa permissão com o papel de
Administrador do Service Usage (roles/serviceusage.serviceUsageAdmin).
Saiba como conceder papéis.
gcloud services enable cloudfunctions.googleapis.comcomposer.googleapis.com
Ativar a API REST do Airflow
Para o Airflow 2, a API REST estável já está ativada por padrão. Se a API estável estiver desativada no ambiente, ative a API REST estável.
Permitir chamadas de API para a API REST do Airflow usando o controle de acesso à rede do servidor da Web
O Cloud Run functions pode acessar a API REST do Airflow usando um endereço IPv4 ou IPv6.
Se você não tiver certeza de qual será o intervalo de IP de chamada, use uma opção de configuração padrão em Controle de acesso ao servidor da Web, que é All IP addresses have access (default), para não bloquear acidentalmente as funções do Cloud Run. Você pode
configurar o acesso à rede do servidor da Web mais tarde.
Ver o URL do servidor da Web do Airflow
Este exemplo faz solicitações da API REST para o endpoint do servidor da Web do Airflow. Use o URL do servidor da Web do Airflow no código da função do Cloud.
Console
No console Google Cloud , acesse a página Ambientes.
Clique no nome do seu ambiente.
Na página Detalhes do ambiente, acesse a guia Configuração do ambiente.
O URL do servidor da Web do Airflow está listado no item da IU da Web do Airflow.
gcloud
Execute este comando:
gcloud composer environments describe ENVIRONMENT_NAME \
--location LOCATION \
--format='value(config.airflowUri)'
Substitua:
ENVIRONMENT_NAMEpelo nome do ambienteLOCATIONpela região em que o ambiente está localizado;
Fazer upload de um DAG para o ambiente
Faça o upload de um DAG para seu ambiente. O exemplo a seguir mostra a configuração de execução do DAG recebida. Você acionará esse DAG a partir de uma função que vai criar neste guia depois.
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 }}}}')
Implantar uma função que aciona o DAG
É possível implantar uma função usando sua linguagem preferida compatível com o Cloud Run Functions ou o Cloud Run. Neste tutorial, demonstramos uma função do Cloud implementada em Python e Java.
Especificar parâmetros de configuração da função
Gatilho: selecione um ou vários gatilhos do Eventarc para sua função.
Para mais informações sobre como criar gatilhos, consulte Criar gatilhos com o Eventarc. Por exemplo, é possível acionar funções do Cloud Storage usando o Eventarc.
Conta de serviço: a conta de serviço especificada para o gatilho precisa ter permissões suficientes para acionar DAGs em ambientes do Airflow Gerenciado.
Recomendamos seguir o princípio de menor privilégio e conceder apenas o papel de Usuário do Composer (
composer.user). Para mais informações sobre como configurar permissões, consulte Papéis e permissões para destinos do Cloud Run.Ponto de entrada da função:
(Python) Ao adicionar código para este exemplo, selecione o ambiente de execução Python 3.10 ou mais recente e especifique
trigger_dag_with_gcfcomo o ponto de entrada.(Java) Ao adicionar código para este exemplo, selecione o tempo de execução Java 17 ou e especifique
functions.TriggerDagExamplecomo o ponto de entrada.
Adicionar requisitos
Python
Especifique as dependências no arquivo requirements.txt:
google-auth>=2.38.0
requests>=2.34.2
functions-framework==3.*
Java
Adicione as seguintes dependências à seção dependencies no 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>
Adicionar código da função
Python
Coloque o seguinte código no arquivo main.py:
Substitua o valor da variável
web_server_urlpelo endereço do servidor da Web do Airflow que você recebeu anteriormente.Se você estiver acionando um DAG diferente, substitua o valor da variável
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
Coloque o seguinte código no arquivo TriggerDagExample.java (coloque esse arquivo no diretório src/main/java/gcfv2/):
Substitua o valor da variável
webServerUrlpelo endereço do servidor da Web do Airflow que você recebeu anteriormente.Se você estiver acionando um DAG diferente, substitua o valor da variável
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());
}
}
}
}
}
Testar a função
Para verificar se a função e o DAG funcionam conforme o esperado:
- Aguarde até que a função seja implantada.
- Acione a função de acordo com o gatilho especificado. Também é possível acionar a função manualmente selecionando a ação Testar a função em Google Cloud console.
- Verifique a página do DAG na interface da Web do Airflow. O DAG precisa ter uma execução ativa ou já concluída.
- Na IU do Airflow, verifique os registros de tarefas desta execução. Você verá que a tarefa
print_gcs_infogera os dados recebidos da função para os registros:
Exemplo de comando para testar a função:
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"
}'
Exemplo de saída:
[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
Solução de problemas:
- Se a função falhar com um erro
NullPointerException: Null datae o stack trace apontar para a funçãoBackgroundFunctionExecutor.parseLegacyEvent, significa que o evento recebido pela função não tem cabeçalhos de metadadosCloudEventpadrão. A função pressupõe que você está enviando um evento em segundo plano legado, tenta analisar o campodatadele e falha. Isso pode acontecer, por exemplo, se você enviar um payload de evento arbitrário ao testar a função. - Se a função falhar com
500 Internal Server Error: The server encountered an internal error and was unable to complete your request., verifique o valor da variávelairflow_major_version. Essa variável determina o endpoint de API REST do Airflow, que é diferente no Airflow 2 e no Airflow 3.
A seguir
- Acessar a IU do Airflow
- Acessar a API REST do Airflow
- Gravar DAGs
- Gravar funções do Cloud Run
- Gatilhos do Cloud Storage