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

求助:如何实现支持Cloud Composer的GKE Autopilot集群夜间暂停?

GKE Autopilot集群(Cloud Composer)夜间暂停方案修复

现有代码核心问题

  1. 变量名冲突:代码中client变量被同时赋值给container_v1.ClusterManagerClient和kubernetes客户端,导致后续调用K8s API时使用错误的客户端实例,引发异常。
  2. Autopilot集群缩容逻辑错误:GKE Autopilot集群无法通过K8s API直接修改节点数,必须通过GCP Container API调整节点池的自动扩缩容参数(min_node_count和max_node_count)。
  3. 未处理Cloud Composer自动扩容机制:Cloud Composer的Airflow组件自带自动扩缩容逻辑(如HPA、Airflow自身调度器),直接缩容节点会触发组件自动扩容Pod,进而重新拉起节点。
  4. 副本数无持久化存储:代码中硬编码恢复副本数为3,未存储缩容前的实际副本数,恢复时会破坏原有部署配置。

修复步骤及代码示例

1. 配置必要IAM权限

给Cloud Function关联的服务账号添加以下角色:

  • roles/container.clusterAdmin:管理GKE集群及节点池
  • roles/container.developer:操作K8s集群内资源
  • roles/cloudfunctions.invoker:允许Pub/Sub触发Cloud Function

2. 修正后的暂停/恢复代码

以下代码实现暂停集群(缩容节点池到0,暂停Airflow自动扩缩容)和恢复集群(恢复节点池配置及Airflow副本数)的逻辑,通过Pub/Sub消息中的action字段区分操作:

import os
from google.cloud import container_v1
from google.auth import default
from kubernetes import client as k8s_client, config
from kubernetes.client.rest import ApiException

# 集群配置(替换为实际值)
CLUSTER_NAME = "your-cluster-name"
LOCATION = "your-cluster-location"
PROJECT_ID = "your-project-id"
NODE_POOL_NAME = "default-pool"  # Autopilot集群默认节点池名称
# 存储原节点池配置的临时路径(生产环境建议改用Cloud Storage/Firestore持久化)
ORIGINAL_CONFIG_PATH = "/tmp/original_cluster_config.json"

def handle_cluster_action(event, context):
    # 解析Pub/Sub消息中的操作指令("pause"或"resume")
    action = event.get('attributes', {}).get('action', 'resume')
    
    # 初始化GCP Container客户端
    credentials, _ = default()
    container_client = container_v1.ClusterManagerClient(credentials=credentials)
    cluster_path = f"projects/{PROJECT_ID}/locations/{LOCATION}/clusters/{CLUSTER_NAME}"

    if action == "pause":
        # 1. 获取当前节点池配置并保存
        node_pool = container_client.get_node_pool(
            name=f"{cluster_path}/nodePools/{NODE_POOL_NAME}"
        )
        original_min = node_pool.autoscaling.min_node_count
        original_max = node_pool.autoscaling.max_node_count
        
        # 保存原配置
        with open(ORIGINAL_CONFIG_PATH, 'w') as f:
            f.write(f'{{"min": {original_min}, "max": {original_max}}}')

        # 2. 调整节点池自动扩缩容参数为0
        update_mask = "autoscaling.min_node_count,autoscaling.max_node_count"
        container_client.update_node_pool(
            name=f"{cluster_path}/nodePools/{NODE_POOL_NAME}",
            node_pool=container_v1.NodePool(
                autoscaling=container_v1.NodePoolAutoscaling(
                    min_node_count=0,
                    max_node_count=0
                )
            ),
            update_mask=update_mask
        )
        print(f"节点池 {NODE_POOL_NAME} 已缩容至0")

        # 3. 暂停Cloud Composer的Airflow自动扩缩容
        # 配置K8s客户端
        config.load_kube_config_from_dict(_build_kube_config(container_client, cluster_path))
        apps_v1 = k8s_client.AppsV1Api()
        autoscaling_v2 = k8s_client.AutoscalingV2Api()

        # 缩容Airflow核心组件deployments到0
        namespace = "composer-<your-composer-environment>-namespace"  # 替换为实际Composer命名空间
        deployments = ["airflow-webserver", "airflow-scheduler", "airflow-worker"]
        for dep_name in deployments:
            try:
                apps_v1.patch_namespaced_deployment(
                    name=dep_name,
                    namespace=namespace,
                    body={"spec": {"replicas": 0}}
                )
                print(f"Deployment {dep_name} 已缩容至0")
            except ApiException as e:
                print(f"缩容Deployment {dep_name} 失败: {e}")

        # 暂停HPA(如果存在)
        for hpa_name in deployments:
            try:
                autoscaling_v2.patch_namespaced_horizontal_pod_autoscaler(
                    name=hpa_name,
                    namespace=namespace,
                    body={"spec": {"minReplicas": 0, "maxReplicas": 0}}
                )
                print(f"HPA {hpa_name} 已暂停")
            except ApiException as e:
                print(f"HPA {hpa_name} 不存在或操作失败: {e}")

    elif action == "resume":
        # 1. 读取原节点池配置并恢复
        with open(ORIGINAL_CONFIG_PATH, 'r') as f:
            original_config = eval(f.read())
        original_min = original_config["min"]
        original_max = original_config["max"]

        update_mask = "autoscaling.min_node_count,autoscaling.max_node_count"
        container_client.update_node_pool(
            name=f"{cluster_path}/nodePools/{NODE_POOL_NAME}",
            node_pool=container_v1.NodePool(
                autoscaling=container_v1.NodePoolAutoscaling(
                    min_node_count=original_min,
                    max_node_count=original_max
                )
            ),
            update_mask=update_mask
        )
        print(f"节点池 {NODE_POOL_NAME} 已恢复原配置: min={original_min}, max={original_max}")

        # 2. 恢复Airflow组件副本数(生产环境建议从持久化存储读取原副本数)
        namespace = "composer-<your-composer-environment>-namespace"
        deployments = {
            "airflow-webserver": 2,
            "airflow-scheduler": 1,
            "airflow-worker": 3
        }
        config.load_kube_config_from_dict(_build_kube_config(container_client, cluster_path))
        apps_v1 = k8s_client.AppsV1Api()
        autoscaling_v2 = k8s_client.AutoscalingV2Api()

        for dep_name, replicas in deployments.items():
            try:
                apps_v1.patch_namespaced_deployment(
                    name=dep_name,
                    namespace=namespace,
                    body={"spec": {"replicas": replicas}}
                )
                print(f"Deployment {dep_name} 已恢复至 {replicas} 副本")
            except ApiException as e:
                print(f"恢复Deployment {dep_name} 失败: {e}")

        # 恢复HPA配置
        hpa_configs = {
            "airflow-worker": {"min": 2, "max": 10}
        }
        for hpa_name, config in hpa_configs.items():
            try:
                autoscaling_v2.patch_namespaced_horizontal_pod_autoscaler(
                    name=hpa_name,
                    namespace=namespace,
                    body={"spec": {"minReplicas": config["min"], "maxReplicas": config["max"]}}
                )
                print(f"HPA {hpa_name} 已恢复原配置")
            except ApiException as e:
                print(f"恢复HPA {hpa_name} 失败: {e}")

    print(f"集群{action}操作完成")

def _build_kube_config(container_client, cluster_path):
    """构建K8s客户端配置"""
    cluster_info = container_client.get_cluster(name=cluster_path)
    credentials, _ = default()
    return {
        "apiVersion": "v1",
        "clusters": [{
            "cluster": {
                "server": f"https://{cluster_info.endpoint}",
                "certificate-authority-data": cluster_info.master_auth.cluster_ca_certificate,
            },
            "name": cluster_info.name,
        }],
        "contexts": [{
            "context": {
                "cluster": cluster_info.name,
                "user": "default",
            },
            "name": "default",
        }],
        "current-context": "default",
        "kind": "Config",
        "preferences": {},
        "users": [{
            "name": "default",
            "user": {
                "auth-provider": {
                    "config": {
                        "access-token": credentials.token,
                        "cmd-args": "config config-helper --format=json",
                        "cmd-path": "gcloud",
                        "expiry-key": "{.credential.token_expiry}",
                        "token-key": "{.credential.access_token}",
                    },
                    "name": "gcp",
                },
            },
        }],
    }

3. Cloud Scheduler配置

创建两个调度任务:

  • 夜间暂停:触发Pub/Sub消息,attributes设置action=pause,调度时间设为夜间非业务时段
  • 日间恢复:触发Pub/Sub消息,attributes设置action=resume,调度时间设为业务开始前

关键说明

  • 生产环境中,原节点池配置和Airflow副本数建议存储在Cloud Storage或Firestore中,避免Cloud Function实例重启后丢失数据。
  • 需替换代码中的集群名称、地域、项目ID及Composer命名空间为实际值。
  • 若Composer环境有自定义组件(如额外的workers),需在代码中添加对应的缩容/恢复逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 11:34:55