Managed Airflow (第 3 代) | Managed Airflow (第 2 代) | Managed Airflow (舊版第 1 代)
本頁說明如何使用 Cloud Run 函式,根據事件觸發 Managed Service for Apache Airflow DAG。
Apache Airflow 的設計是定期執行 DAG,但您也可以在發生事件時觸發 DAG。其中一種做法是使用 Cloud Run 函式,在發生指定事件時觸發 Managed Airflow DAG。
您也能執行下列操作:
- 只使用 Airflow REST API 觸發 DAG。
- 建立函式,在訊息推送至 Pub/Sub 主題時觸發 DAG。
本指南中的範例會示範函式如何觸發 DAG 來回應事件:
- 您可以在 Cloud Run 函式中設定函式的觸發條件。
- 函式觸發後,會透過 Managed Airflow 環境的 Airflow REST API 提出要求,觸發 DAG。要求包含事件的 ID 和類型,以及事件的酬載。
- Airflow 會處理這項要求,並執行要求中指定的 DAG。 DAG 會輸出從函式傳遞給 DAG 的資料。
事前準備
本節列出準備步驟。
檢查環境的網路設定
這個解決方案不適用於私人 IP 和 VPC Service Controls 設定,因為在這些設定中,無法設定從 Cloud Run 函式到 Airflow 網路伺服器的連線。
為專案啟用 API
控制台
啟用 Managed Airflow 和 Cloud Run functions API。
啟用 API 時所需的角色
如要啟用 API,您必須具備 serviceusage.services.enable 權限。如果您建立了專案,可能已透過「擁有者」角色 (roles/owner) 取得這項權限。否則,您可以透過「服務使用情形管理員」角色 (roles/serviceusage.serviceUsageAdmin) 取得這項權限。瞭解如何授予角色。
gcloud
啟用 Managed Airflow 和 Cloud Run functions API:
啟用 API 時所需的角色
如要啟用 API,您必須具備 serviceusage.services.enable 權限。如果您建立了專案,可能已透過「擁有者」角色 (roles/owner) 取得這項權限。否則,您可以透過「服務使用情形管理員」角色 (roles/serviceusage.serviceUsageAdmin) 取得這項權限。瞭解如何授予角色。
gcloud services enable cloudfunctions.googleapis.comcomposer.googleapis.com
啟用 Airflow REST API
視 Airflow 版本而定:
- 對於 Airflow 2,穩定版 REST API 預設為啟用。如果您的環境停用了穩定版 API,請啟用穩定版 REST API。
- 如果是 Airflow 1,請啟用實驗性 REST API。
使用網路伺服器網路存取控管,允許對 Airflow REST API 進行 API 呼叫
Cloud Run 函式可透過 IPv4 或 IPv6 位址連線至 Airflow REST API。
如果不確定呼叫 IP 範圍,請使用「Webserver Access Control」中的預設設定選項 All IP addresses have access (default),以免不慎封鎖 Cloud Run functions。您隨時可以稍後設定網路伺服器網路存取權。
取得 Airflow 網路伺服器網址
這個範例會向 Airflow 網路伺服器端點發出 REST API 要求。在 Cloud Functions 程式碼中,您會使用 Airflow 網頁介面網址中 .appspot.com 前的部分。
控制台
前往 Google Cloud 控制台的「Environments」(環境) 頁面。
按一下環境名稱。
在「環境詳細資料」頁面中,前往「環境設定」分頁。
Airflow 網路伺服器的網址會列在「Airflow web UI」(Airflow 網頁版 UI) 項目中。
gcloud
執行下列指令:
gcloud composer environments describe ENVIRONMENT_NAME \
--location LOCATION \
--format='value(config.airflowUri)'
更改項目:
- 將
ENVIRONMENT_NAME替換為環境的名稱。 - 將
LOCATION替換為環境所在的區域。
取得 IAM Proxy 的 client_id
如要向 Airflow REST API 端點提出要求,您必須為函式提供保護 Airflow 網路伺服器的 Identity and Access Management Proxy 用戶端 ID。
Managed Airflow 不會直接提供這項資訊,而會向 Airflow 網路伺服器發出未經授權的要求,並從重新導向網址中擷取用戶端 ID:
cURL
curl -v AIRFLOW_URL 2>&1 >/dev/null | grep -o "client_id\=[A-Za-z0-9-]*\.apps\.googleusercontent\.com"
將 AIRFLOW_URL 替換為 Airflow 網頁介面的網址。
在輸出內容中,搜尋 client_id 後方的字串。例如:
client_id=836436932391-16q2c5f5dcsfnel77va9bvf4j280t35c.apps.googleusercontent.com
Python
將下列程式碼儲存為 get_client_id.py 檔案。填入 project_id、location 和 composer_environment 的值,然後在 Cloud Shell 或本機環境中執行程式碼。
將 DAG 上傳至環境
將 DAG 上傳至環境。 下列 DAG 範例會輸出收到的 DAG 執行設定。您會從函式觸發這個 DAG,該函式會在稍後建立。
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 }}}}')
部署會觸發 DAG 的函式
您可以使用 Cloud Run functions 或 Cloud Run 支援的偏好語言部署函式。本教學課程將示範以 Python 和 Java 實作的 Cloud 函式。
指定函式設定參數
觸發條件:為函式選取一個或多個 Eventarc 觸發條件。
如要進一步瞭解如何建立觸發條件,請參閱使用 Eventarc 建立觸發條件。 舉例來說,您可以使用 Eventarc,透過 Cloud Storage 觸發函式。
服務帳戶:您為觸發條件指定的服務帳戶必須具備足夠的權限,才能在 Managed Airflow 環境中觸發 DAG。
建議您遵循最低權限原則,只授予該帳戶「Composer 使用者」 (
composer.user) 角色。如要進一步瞭解如何設定權限,請參閱「Cloud Run 目標的角色和權限」。函式進入點:
(Python) 新增這個範例的程式碼時,請選取 Python 3.10 以上版本的執行階段,並將
trigger_dag指定為進入點。
新增規定
在 requirements.txt 檔案中指定依附元件:
新增函式程式碼
將下列程式碼放入 main.py 檔案,並進行下列替換:
將
client_id變數的值替換為先前取得的client_id值。將
webserver_id變數的值替換為您的租戶專案 ID,這是 Airflow 網頁介面網址中.appspot.com前的部分。您先前已取得 Airflow 網頁介面網址。指定您使用的 Airflow REST API 版本:
- 如果您使用穩定版 Airflow REST API,請將
USE_EXPERIMENTAL_API變數設為False。 - 如果您使用實驗性 Airflow REST API,則無需進行任何變更。
USE_EXPERIMENTAL_API變數已設為True。
- 如果您使用穩定版 Airflow REST API,請將
測試函式
如要確認函式和 DAG 是否正常運作,請按照下列步驟操作:
- 等待函式部署完成。
- 根據指定的觸發條件觸發函式。您也可以在 Google Cloud 控制台中選取函式的「測試函式」動作,手動觸發函式。
- 在 Airflow 網頁介面中查看 DAG 頁面。DAG 應有一個有效或已完成的 DAG 執行作業。
- 在 Airflow UI 中,查看這項執行的工作記錄。您應該會看到
print_gcs_info工作將從函式收到的資料輸出至記錄:
輸出內容範例:
[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