Managed Airflow(第 3 代) | Managed Airflow(第 2 代) | Managed Airflow(旧版第 1 代)
本页面介绍如何使用 Cloud Run 函数触发 Managed Service for Apache Airflow DAG 来响应事件。
Apache Airflow 设计为定期运行 DAG,但您也可以触发 DAG 来响应事件。方法之一是使用 Cloud Run functions 在发生指定事件时触发 托管式 Airflow DAG。
您还可以:
本指南中的示例演示了一个函数,该函数会触发 DAG 来响应事件:
- 您可以在 Cloud Run functions 中为函数配置触发器。
- 触发函数后,它会发出请求,以通过 Managed Airflow 环境的 Airflow REST API 触发 DAG。该请求包含事件的标识符和类型,以及事件的载荷。
- Airflow 会处理此请求并运行请求中指定的 DAG。 DAG 会输出从函数传递给它的数据。
准备工作
本部分列出了准备步骤。
检查环境的网络配置
此解决方案不适用于专用 IP 和 VPC Service Controls 配置,因为在这些配置中,无法配置从 Cloud Run functions 到 Airflow Web 服务器的连接。
为您的项目启用 API
控制台
启用托管式 Airflow 和 Cloud Run functions API。
启用 API 所需的角色
如需启用 API,您需要拥有 serviceusage.services.enable 权限。如果您
创建了项目,则可能已通过
所有者角色 (roles/owner) 拥有此权限。否则,您可以通过
Service Usage Admin 角色 (roles/serviceusage.serviceUsageAdmin) 获得此权限。
了解如何授予角色。
gcloud
启用 Managed Airflow 和 Cloud Run functions API:
启用 API 所需的角色
如需启用 API,您需要拥有 serviceusage.services.enable 权限。如果您创建了项目,则可能已通过
所有者角色 (roles/owner) 拥有此权限。否则,您可以通过
Service Usage Admin 角色 (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。
使用 Web 服务器网络访问权限控制允许对 Airflow REST API 进行 API 调用
Cloud Run functions 可以通过 IPv4 或 IPv6 地址 访问 Airflow REST API。
如果您不确定调用 IP 范围,请使用网络服务器访问权限控制 中的默认配置选项 All IP addresses have access (default),以免意外屏蔽 Cloud Run functions。您可以随时
配置 Web 服务器网络访问权限。
获取 Airflow Web 服务器网址
此示例向 Airflow 网络服务器端点发出 REST API 请求。您需要在 Cloud Functions 函数代码中使用 .appspot.com 之前的 Airflow 网页界面网址部分。
控制台
在 Google Cloud 控制台中,前往 环境 页面。
点击您的环境的名称。
在环境详情 页面上,前往环境配置 标签页。
Airflow 网页界面 项中列出了 Airflow Web 服务器的网址。
gcloud
运行以下命令:
gcloud composer environments describe ENVIRONMENT_NAME \
--location LOCATION \
--format='value(config.airflowUri)'
替换:
ENVIRONMENT_NAME替换为环境的名称。LOCATION替换为环境所在的区域。
获取 IAM 代理的 client_id
要向 Airflow REST API 端点发出请求,函数需要获得保护 Airflow Web 服务器的 Identity and Access Management 代理的客户端 ID。
Managed Airflow 不直接提供此信息。而是向 Airflow Web 服务器发出未经身份验证的请求,并从重定向网址捕获此客户端 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 Function。
指定函数配置参数
触发器:为函数选择一个或多个 Eventarc 触发器。
如需详细了解如何创建触发器,请参阅 使用 Eventarc 创建触发器。 例如,您可以使用 Eventarc 从 Cloud Storage 触发函数。
服务账号:您为触发器指定的服务账号必须 具有足够的权限,才能在 Managed Airflow 环境中触发 DAG。
我们建议遵循最小权限原则,仅向其授予 Composer User (
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 界面中,检查此运行的任务日志。您应该会看到
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