Memicu DAG Managed Service for Apache Airflow dengan fungsi Cloud Run dan Airflow REST API

Managed Airflow (Generasi ke-3) | Managed Airflow (Generasi ke-2) | Managed Airflow (Generasi ke-1 Lama)

Halaman ini menjelaskan cara menggunakan Cloud Run Functions untuk memicu DAG Managed Service for Apache Airflow sebagai respons terhadap peristiwa.

Apache Airflow dirancang untuk menjalankan DAG sesuai jadwal reguler, tetapi Anda juga dapat memicu DAG sebagai respons terhadap peristiwa. Salah satu cara untuk melakukannya adalah menggunakan Cloud Run Functions untuk memicu DAG Managed Airflow saat peristiwa tertentu terjadi.

Anda juga dapat:

Contoh dalam panduan ini menunjukkan fungsi yang memicu DAG sebagai respons terhadap peristiwa:

  1. Anda mengonfigurasi pemicu untuk fungsi di Cloud Run Functions.
  2. Saat dipicu, fungsi akan membuat permintaan untuk memicu DAG melalui Airflow REST API dari lingkungan Managed Airflow Anda. Permintaan tersebut berisi ID dan jenis peristiwa, serta payload peristiwa.
  3. Airflow memproses permintaan ini dan menjalankan DAG yang ditentukan dalam permintaan. DAG menampilkan data yang diteruskan ke DAG dari fungsi.

Sebelum memulai

Bagian ini mencantumkan langkah-langkah persiapan.

Memeriksa konfigurasi jaringan lingkungan Anda

Solusi ini tidak berfungsi dalam konfigurasi IP Pribadi dan Kontrol Layanan VPC karena konektivitas dari Cloud Run Functions ke server web Airflow tidak dapat dikonfigurasi dalam konfigurasi ini.

Di Managed Airflow (Generasi ke-3), Anda dapat menggunakan pendekatan lain: Memicu DAG menggunakan Cloud Run Functions dan Pesan Pub/Sub.

Mengaktifkan API untuk project Anda

Konsol

Aktifkan Managed Airflow dan Cloud Run Functions API.

Peran yang diperlukan untuk mengaktifkan API

Untuk mengaktifkan API, Anda memerlukan izin serviceusage.services.enable. Jika Anda membuat project, kemungkinan besar Anda sudah memiliki izin ini melalui peran Pemilik (roles/owner). Jika tidak, Anda bisa mendapatkan izin ini melalui peran Admin Penggunaan Layanan (roles/serviceusage.serviceUsageAdmin). Pelajari cara memberikan peran.

Aktifkan API

gcloud

Aktifkan Managed Airflow dan Cloud Run Functions API:

Peran yang diperlukan untuk mengaktifkan API

Untuk mengaktifkan API, Anda memerlukan izin serviceusage.services.enable. Jika Anda membuat project, kemungkinan besar Anda sudah memiliki izin ini melalui peran Pemilik (roles/owner). Jika tidak, Anda bisa mendapatkan izin ini melalui peran Admin Penggunaan Layanan (roles/serviceusage.serviceUsageAdmin). Pelajari cara memberikan peran.

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

Mengizinkan panggilan API ke Airflow REST API menggunakan kontrol akses jaringan server web

Cloud Run Functions dapat menjangkau Airflow REST API melalui alamat IPv4 atau IPv6.

Jika Anda tidak yakin dengan rentang IP panggilan, gunakan opsi konfigurasi default di Kontrol Akses Server Web yang merupakan All IP addresses have access (default) agar tidak memblokir Cloud Run Functions Anda secara tidak sengaja. Anda selalu dapat mengonfigurasi akses jaringan server web nanti.

Mendapatkan URL server web Airflow

Contoh ini membuat permintaan REST API ke endpoint server web Airflow. Anda menggunakan URL server web Airflow dalam kode Cloud Function.

Konsol

  1. Di Google Cloud konsol, buka halaman Environments.

    Buka Environments

  2. Klik nama lingkungan Anda.

  3. Di halaman Environment details, buka tab Environment configuration.

  4. URL server web Airflow tercantum dalam item Airflow web UI.

gcloud

Jalankan perintah berikut:

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

Ganti:

  • ENVIRONMENT_NAME dengan nama lingkungan.
  • LOCATION dengan region tempat lingkungan berada.

Mengupload DAG ke lingkungan Anda

Upload DAG ke lingkungan Anda. Contoh DAG berikut menampilkan konfigurasi eksekusi DAG yang diterima. Anda memicu DAG ini dari fungsi, yang akan Anda buat nanti dalam panduan ini.

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

Men-deploy fungsi yang memicu DAG

Anda dapat men-deploy fungsi menggunakan bahasa pilihan yang didukung oleh Cloud Run Functions atau Cloud Run. Tutorial ini menunjukkan a Cloud Function yang diimplementasikan di Python dan Java.

Menentukan parameter konfigurasi fungsi

  • Trigger: Pilih pemicu Eventarc atau beberapa pemicu untuk fungsi Anda.

    Untuk mengetahui informasi selengkapnya tentang cara membuat pemicu, lihat Membuat pemicu dengan Eventarc. Misalnya, Anda dapat memicu fungsi dari Cloud Storage menggunakan Eventarc.

  • Akun layanan: Akun layanan yang Anda tentukan untuk pemicu harus memiliki izin yang cukup untuk memicu DAG di lingkungan Managed Airflow.

    Sebaiknya ikuti prinsip hak istimewa minimum dan berikan hanya peran Pengguna Composer (composer.user). Untuk mengetahui informasi selengkapnya tentang cara mengonfigurasi izin, lihat Peran dan izin untuk target Cloud Run.

  • Function entry point:

    • (Python) Saat menambahkan kode untuk contoh ini, pilih runtime Python 3.10 atau yang lebih baru dan tentukan trigger_dag_with_gcf sebagai titik entri.

    • (Java) Saat menambahkan kode untuk contoh ini, pilih runtime Java 17 atau yang lebih baru dan tentukan functions.TriggerDagExample sebagai titik entri.

Menambahkan persyaratan

Python

Tentukan dependensi dalam file requirements.txt:

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

Java

Tambahkan dependensi berikut ke bagian dependencies di 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>

Menambahkan kode fungsi

Python

Masukkan kode berikut ke file main.py:

  • Ganti nilai variabel web_server_url dengan alamat server web Airflow yang Anda dapatkan sebelumnya.

  • (Airflow 3) Ganti nilai variabel airflow_major_version dengan 3, yang merupakan versi utama Airflow di lingkungan Anda.

  • Jika Anda memicu DAG yang berbeda, ganti nilai variabel 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

Masukkan kode berikut ke file TriggerDagExample.java (tempatkan file ini ke direktori src/main/java/gcfv2/):

  • Ganti nilai variabel webServerUrl dengan alamat server web Airflow yang Anda dapatkan sebelumnya.

  • (Airflow 3) Ganti nilai variabel majorAirflowVersion dengan 3, yang merupakan versi utama Airflow di lingkungan Anda.

  • Jika Anda memicu DAG yang berbeda, ganti nilai variabel 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());
        }
      }
    }
  }
}

Menguji fungsi

Untuk memeriksa apakah fungsi dan DAG Anda berfungsi sesuai harapan:

  1. Tunggu hingga fungsi Anda di-deploy.
  2. Picu fungsi sesuai dengan pemicu yang ditentukan. Anda juga dapat memicu fungsi secara manual dengan memilih tindakan Test the function untuk fungsi tersebut di Google Cloud konsol.
  3. Periksa halaman DAG di antarmuka web Airflow. DAG harus memiliki satu eksekusi DAG yang aktif atau sudah selesai.
  4. Di UI Airflow, periksa log tugas untuk eksekusi ini. Anda akan melihat bahwa tugas print_gcs_info menampilkan data yang diterima dari fungsi ke log:

Contoh perintah untuk menguji fungsi:

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

Contoh output:

[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

Pemecahan masalah:

  • Jika fungsi Anda gagal dengan error NullPointerException: Null data dan stack trace mengarah ke fungsi BackgroundFunctionExecutor.parseLegacyEvent, artinya peristiwa yang diterima oleh fungsi tidak memiliki header metadata CloudEvent standar. Fungsi ini mengasumsikan bahwa Anda mengirim peristiwa latar belakang lama, mencoba mengurai kolom data dari peristiwa tersebut, dan gagal. Hal ini dapat terjadi, misalnya, jika Anda mengirim payload peristiwa arbitrer saat menguji fungsi.
  • Jika fungsi Anda gagal dengan 500 Internal Server Error: The server encountered an internal error and was unable to complete your request., periksa kembali nilai variabel airflow_major_version. Variabel ini menentukan endpoint Airflow REST API, yang berbeda di Airflow 2 dan Airflow 3.

Langkah berikutnya