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.
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
Bergantung pada versi Airflow Anda:
- Untuk Airflow 2, REST API yang stabil sudah diaktifkan secara default. Jika API stabil dinonaktifkan di lingkungan Anda, aktifkan REST API stabil.
- Untuk Airflow 1, aktifkan REST API eksperimental.
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 bagian URL antarmuka web Airflow sebelum .appspot.com 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.
Mendapatkan client_id proxy IAM
Untuk membuat permintaan ke endpoint Airflow REST API, fungsi ini memerlukan ID klien proxy Identity and Access Management yang melindungi server web Airflow.
Managed Airflow tidak memberikan informasi ini secara langsung. Sebagai gantinya, buat permintaan yang tidak diautentikasi ke server web Airflow dan ambil ID klien dari URL pengalihan:
cURL
curl -v AIRFLOW_URL 2>&1 >/dev/null | grep -o "client_id\=[A-Za-z0-9-]*\.apps\.googleusercontent\.com"
Ganti AIRFLOW_URL dengan URL antarmuka web Airflow.
Pada output, telusuri string setelah client_id. Contoh:
client_id=836436932391-16q2c5f5dcsfnel77va9bvf4j280t35c.apps.googleusercontent.com
Python
Simpan kode berikut dalam file bernama get_client_id.py. Isi nilai
untuk project_id, location, dan composer_environment, lalu jalankan
kode di Cloud Shell atau lingkungan lokal Anda.
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_dagsebagai titik entri.
Menambahkan persyaratan
Tentukan dependensi dalam file requirements.txt:
Menambahkan kode fungsi
Masukkan kode berikut ke file main.py dan lakukan penggantian berikut:
Ganti nilai variabel
client_iddengan nilaiclient_idyang Anda peroleh sebelumnya.Ganti nilai variabel
webserver_iddengan ID project tenant Anda, yang merupakan bagian dari URL antarmuka web Airflow sebelum.appspot.com. Anda telah mendapatkan URL antarmuka web Airflow sebelumnya.Tentukan versi Airflow REST API yang Anda gunakan:
- Jika Anda menggunakan Airflow REST API yang stabil, tetapkan variabel
USE_EXPERIMENTAL_APIkeFalse. - Jika Anda menggunakan Airflow REST API eksperimental, tidak ada perubahan yang diperlukan. Variabel
USE_EXPERIMENTAL_APIsudah disetel keTrue.
- Jika Anda menggunakan Airflow REST API yang stabil, tetapkan variabel
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 output:
[2021-04-04 18:25:44,778] {bash_operator.py:154} INFO - Output:
[2021-04-04 18:25:44,781] {bash_operator.py:158} INFO - Triggered from GCF:
{bucket: example-storage-for-gcf-triggers, contentType: text/plain,
crc32c: dldNmg==, etag: COW+26Sb5e8CEAE=, generation: 1617560727904101,
... }
[2021-04-04 18:25:44,781] {bash_operator.py:162} INFO - Command exited with
return code 0h
Langkah berikutnya
- Mengakses UI Airflow
- Mengakses Airflow REST API
- Menulis DAG
- Menulis fungsi Cloud Run
- Pemicu Cloud Storage