如何在Cloud Function中识别可用的Composer 2环境以调用DAG?
检测Composer 2环境可用性并实现高可用DAG调用
一、用现有函数实现可用性检测
你提供的make_composer2_web_server_request函数完全可以用来判断Composer环境是否可用,核心思路是请求Airflow的健康检查端点:
- 健康检查端点:Composer 2的Airflow服务提供
/health端点,返回环境的健康状态。 - 扩展检测逻辑:基于现有函数封装检测方法,通过响应状态和内容判断可用性。
示例代码
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
- 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
相关产品推荐
相关产品推荐

