Managed Airflow (Gen 3) | Managed Airflow (Gen 2) | Managed Airflow (Legacy Gen 1)
Halaman ini menjelaskan cara menggunakan fungsi Cloud Run untuk memicu DAG Managed Service untuk Apache Airflow sebagai respons terhadap peristiwa.
Apache Airflow dirancang untuk menjalankan DAG sesuai jadwal rutin, tetapi Anda juga dapat memicu DAG sebagai respons terhadap peristiwa. Salah satu cara untuk melakukannya adalah dengan menggunakan Cloud Run Functions untuk memicu DAG Managed Airflow saat peristiwa tertentu terjadi.
Anda juga dapat:
- Memicu DAG hanya menggunakan Airflow REST API.
- Buat fungsi yang memicu DAG saat pesan dikirim ke topik Pub/Sub.
Contoh dalam panduan ini menunjukkan fungsi yang memicu DAG sebagai respons terhadap suatu peristiwa:
- Anda mengonfigurasi pemicu untuk fungsi di Cloud Run Functions.
- Saat dipicu, fungsi akan membuat permintaan untuk memicu DAG melalui Airflow REST API dari lingkungan Managed Airflow Anda. Permintaan berisi ID dan jenis peristiwa, serta payload peristiwa.
- Airflow memproses permintaan ini dan menjalankan DAG yang ditentukan dalam permintaan. DAG menampilkan data yang diteruskan kepadanya dari fungsi.
Sebelum memulai
Bagian ini mencantumkan langkah-langkah persiapan.
Periksa konfigurasi jaringan lingkungan Anda
Solusi ini tidak berfungsi dalam konfigurasi Kontrol Layanan VPC dan IP Pribadi karena tidak mungkin mengonfigurasi konektivitas dari fungsi Cloud Run ke server web Airflow dalam konfigurasi ini.
Di Managed Airflow (Gen 2), Anda dapat menggunakan pendekatan lain: Memicu DAG menggunakan fungsi Cloud Run 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 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
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 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.comcomposer.googleapis.com
Mengaktifkan Airflow REST API
Untuk Airflow 2, REST API yang stabil sudah diaktifkan secara default. Jika API stabil dinonaktifkan di lingkungan Anda, aktifkan REST API stabil.
Mengizinkan panggilan API ke Airflow REST API menggunakan kontrol akses jaringan server web
Fungsi Cloud Run dapat menghubungi Airflow REST API melalui alamat IPv4 atau IPv6.
Jika Anda tidak yakin dengan rentang IP panggilan, gunakan opsi konfigurasi default di
Kontrol Akses Webserver, yaitu 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 Anda.
Konsol
Di konsol Google Cloud , buka halaman Environments.
Klik nama lingkungan Anda.
Di halaman Detail lingkungan, buka tab Konfigurasi lingkungan.
URL server web Airflow tercantum di item Airflow web UI.
gcloud
Jalankan perintah berikut:
gcloud composer environments describe ENVIRONMENT_NAME \
--location LOCATION \
--format='value(config.airflowUri)'
Ganti:
ENVIRONMENT_NAMEdengan nama lingkungan.LOCATIONdengan region tempat lingkungan berada.
Mengupload DAG ke lingkungan Anda
Mengupload DAG ke lingkungan Anda. DAG contoh berikut menampilkan output konfigurasi eksekusi DAG yang diterima. Anda memicu DAG ini dari fungsi, yang akan Anda buat nanti dalam panduan ini.
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 fungsi Cloud Run atau Cloud Run. Tutorial ini menunjukkan Cloud Function yang diimplementasikan di Python dan Java.
Menentukan parameter konfigurasi fungsi
Pemicu: 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 memadai untuk memicu DAG di lingkungan Managed Airflow.
Sebaiknya ikuti prinsip hak istimewa minimum dan berikan hanya peran Pengguna Composer (
composer.user) kepadanya. Untuk mengetahui informasi selengkapnya tentang mengonfigurasi izin, lihat Peran dan izin untuk target Cloud Run.Titik entri fungsi:
(Python) Saat menambahkan kode untuk contoh ini, pilih runtime Python 3.10 atau yang lebih baru dan tentukan
trigger_dag_with_gcfsebagai titik entri.(Java) Saat menambahkan kode untuk contoh ini, pilih runtime Java 17 atau yang lebih baru dan tentukan
functions.TriggerDagExamplesebagai 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_urldengan alamat server web Airflow yang Anda dapatkan sebelumnya.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
(masukkan file ini ke direktori src/main/java/gcfv2/):
Ganti nilai variabel
webServerUrldengan alamat server web Airflow yang Anda dapatkan sebelumnya.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 sebagaimana mestinya:
- Tunggu hingga fungsi Anda di-deploy.
- Memicu fungsi sesuai dengan pemicu yang ditentukan. Anda juga dapat memicu fungsi secara manual dengan memilih tindakan Uji fungsi untuk fungsi tersebut di konsol Google Cloud .
- Periksa halaman DAG di antarmuka web Airflow. DAG harus memiliki satu proses DAG yang aktif atau sudah selesai.
- Di UI Airflow, periksa log tugas untuk proses ini. Anda akan melihat
bahwa tugas
print_gcs_infomenampilkan 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 datadan stack trace mengarah ke fungsiBackgroundFunctionExecutor.parseLegacyEvent, artinya peristiwa yang diterima oleh fungsi tidak memiliki header metadataCloudEventstandar. Fungsi ini mengasumsikan Anda mengirim peristiwa latar belakang lama, mencoba mengurai kolomdatadari 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 variabelairflow_major_version. Variabel ini menentukan endpoint Airflow REST API, yang berbeda di Airflow 2 dan Airflow 3.
Langkah berikutnya
- Mengakses UI Airflow
- Mengakses Airflow REST API
- Menulis DAG
- Menulis fungsi Cloud Run
- Pemicu Cloud Storage