Acionar DAGs do Serviço gerenciado para Apache Airflow com funções do Cloud Run e a API REST do Airflow

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:

O exemplo neste guia demonstra uma função que aciona um DAG em resposta a um evento:

  1. Você configura acionadores para sua função no Cloud Run functions.
  2. 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.
  3. 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.

Ativar as APIs

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.com composer.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

  1. No console Google Cloud , acesse a página Ambientes.

    Acessar "Ambientes"

  2. Clique no nome do seu ambiente.

  3. Na página Detalhes do ambiente, acesse a guia Configuração do ambiente.

  4. 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_NAME pelo nome do ambiente
  • LOCATION pela 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_gcf como o ponto de entrada.

    • (Java) Ao adicionar código para este exemplo, selecione o tempo de execução Java 17 ou e especifique functions.TriggerDagExample como 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_url pelo 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 webServerUrl pelo 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:

  1. Aguarde até que a função seja implantada.
  2. 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.
  3. Verifique a página do DAG na interface da Web do Airflow. O DAG precisa ter uma execução ativa ou já concluída.
  4. Na IU do Airflow, verifique os registros de tarefas desta execução. Você verá que a tarefa print_gcs_info gera 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 data e o stack trace apontar para a função BackgroundFunctionExecutor.parseLegacyEvent, significa que o evento recebido pela função não tem cabeçalhos de metadados CloudEvent padrão. A função pressupõe que você está enviando um evento em segundo plano legado, tenta analisar o campo data dele 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ável airflow_major_version. Essa variável determina o endpoint de API REST do Airflow, que é diferente no Airflow 2 e no Airflow 3.

A seguir