You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Cloud Function中识别可用的Composer 2环境以调用DAG?

检测Composer 2环境可用性并实现高可用DAG调用

一、用现有函数实现可用性检测

你提供的make_composer2_web_server_request函数完全可以用来判断Composer环境是否可用,核心思路是请求Airflow的健康检查端点:

  1. 健康检查端点:Composer 2的Airflow服务提供/health端点,返回环境的健康状态。
  2. 扩展检测逻辑:基于现有函数封装检测方法,通过响应状态和内容判断可用性。

示例代码

import google.auth
from google.auth.transport.requests import AuthorizedSession

# 全局获取GCP凭证
CREDENTIALS, _ = google.auth.default(scopes=["https://www.googleapis.com/auth/cloud-platform"])

def make_composer2_web_server_request(url: str, method: str = "GET", **kwargs) -> google.auth.transport.Response:
    """向Composer 2环境Web服务器发起认证请求"""
    authed_session = AuthorizedSession(CREDENTIALS)
    kwargs.setdefault("timeout", 10)  # 健康检查用较短超时
    return authed_session.request(method, url, **kwargs)

def is_composer_available(composer_web_url: str) -> bool:
    """检测Composer环境是否可用"""
    health_url = f"{composer_web_url.rstrip('/')}/health"
    try:
        resp = make_composer2_web_server_request(health_url)
        return resp.status_code == 200 and "healthy" in resp.text.lower()
    except (ConnectionError, TimeoutError, Exception) as e:
        print(f"环境不可用: {str(e)}")
        return False
  1. Cloud Function中的故障转移逻辑:维护两个环境的URL列表,依次检测找到可用环境后触发DAG:
def on_bucket_file_upload(event, context):
    # 配置双Composer环境的Web URL
    COMPOSER_ENVS = [
        "https://your-primary-composer-web-url",
        "https://your-secondary-composer-web-url"
    ]
    TARGET_DAG = "your-dag-id"
    available_env = None

    # 遍历找第一个可用环境
    for env_url in COMPOSER_ENVS:
        if is_composer_available(env_url):
            available_env = env_url
            break

    if not available_env:
        raise RuntimeError("所有Composer环境均不可用")
    
    # 触发目标DAG
    trigger_url = f"{available_env.rstrip('/')}/api/v1/dags/{TARGET_DAG}/dagRuns"
    payload = {"conf": {"uploaded_file": event["name"]}}
    make_composer2_web_server_request(trigger_url, method="POST", json=payload)
    print(f"通过环境 {available_env} 成功触发DAG {TARGET_DAG}")

二、其他检测方案

1. 调用Composer API查询环境状态

直接通过Composer API获取环境的官方状态,准确性更高:

from google.cloud import composer_v1

def get_composer_state(project_id: str, region: str, env_name: str) -> str:
    """获取Composer环境的运行状态"""
    client = composer_v1.EnvironmentsClient()
    env_path = client.environment_path(project_id, region, env_name)
    return client.get_environment(name=env_path).state.name

# 使用示例:判断环境是否处于运行状态
if get_composer_state("your-gcp-project", "us-central1", "composer-primary") == "RUNNING":
    # 调用该环境的DAG
    pass

权限要求:给Cloud Function的服务账号添加roles/composer.environmentViewer权限。

2. Cloud Monitoring告警驱动(进阶)

提前配置Cloud Monitoring告警,当Composer环境状态变为非RUNNING时,自动更新Cloud Function的环境变量(比如切换可用环境URL)。这种方式适合大规模场景,避免每次触发都做主动检测。

三、关键注意事项

  • 权限配置:确保Cloud Function服务账号拥有roles/composer.user权限(用于触发DAG),以及对应检测方案所需的权限。
  • 重试机制:给健康检查添加重试逻辑,排除临时网络波动导致的误判。
  • DAG同步:双环境的DAG代码、依赖必须保持一致,避免执行结果差异。

内容的提问来源于stack exchange,提问作者PROTOCOL

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.19 04:40:17