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

Cloud Composer 2如何避免Worker Pod被驱逐并兼顾集群缩容?

解决Cloud Composer 2中Worker Pod的自动扩缩容驱逐问题

我完全理解你升级到Cloud Composer 2时的顾虑——GKE Autopilot的自动扩缩容确实会在节点需要缩容时驱逐可调度的Pod,而任务重试容忍度低的话很容易出故障。下面分两种方案帮你解决这个问题,重点推荐动态管理注释的方案,兼顾任务稳定性和集群缩容效率:


方案一:静态添加safe-to-evict注释(全局生效,需注意缩容影响)

如果想让所有Composer Worker Pod默认都不被驱逐,可以通过自定义Airflow的Pod模板来实现:

  • 方法1:通过gcloud命令更新Composer环境
    在更新环境时指定自定义Pod模板文件:
    gcloud composer environments update YOUR_COMPOSER_ENV_NAME \
      --location YOUR_REGION \
      --update-airflow-config kubernetes.pod_template_file=gs://YOUR_BUCKET/pod_template.yaml
    
    其中pod_template.yaml内容如下,添加需要的注释:
    apiVersion: v1
    kind: Pod
    metadata:
      annotations:
        cluster-autoscaler.kubernetes.io/safe-to-evict: "false"
    spec:
      containers:
        - name: airflow-worker
          image: YOUR_COMPOSER_WORKER_IMAGE
          # 其他默认配置保持和Composer一致
    
  • 方法2:通过Cloud Console配置
    在Composer环境的「Airflow配置」页面,找到kubernetes.pod_template_file配置项,输入GCS存储桶中自定义模板的路径即可。

⚠️ 注意:这种静态配置会让所有Worker Pod(包括空闲状态的)都无法被驱逐,可能导致集群节点无法正常缩容,长期下来会增加成本,因此更推荐下面的动态方案。


方案二:动态管理注释(结合任务生命周期,最优解)

这个方案会在任务运行时给对应的Worker Pod添加safe-to-evict: "false"注释,任务**完成(成功/失败)**后自动移除注释,既保证运行中任务不被打断,又不影响空闲Pod被驱逐和集群缩容。

步骤1:编写回调函数(操作Kubernetes API)

首先需要在Airflow中编写两个回调函数,分别用于添加和移除注释。Composer的Airflow Pod默认有权限访问GKE集群的Kubernetes API,所以不需要额外配置认证:

from kubernetes import client, config
from airflow.models import TaskInstance

def add_safe_to_evict(task_instance: TaskInstance):
    # 获取当前任务所在的Worker Pod名称(Celery Worker的hostname就是Pod名称)
    pod_name = task_instance.hostname
    namespace = "composer-<YOUR_ENV_NAME>-<SUFFIX>"  # 替换成你的Composer命名空间,可通过kubectl get namespaces查看
    
    # 加载集群内部K8s配置
    config.load_incluster_config()
    v1 = client.CoreV1Api()
    
    # 更新Pod注释
    body = {
        "metadata": {
            "annotations": {
                "cluster-autoscaler.kubernetes.io/safe-to-evict": "false"
            }
        }
    }
    v1.patch_namespaced_pod(name=pod_name, namespace=namespace, body=body)

def remove_safe_to_evict(task_instance: TaskInstance):
    pod_name = task_instance.hostname
    namespace = "composer-<YOUR_ENV_NAME>-<SUFFIX>"
    
    config.load_incluster_config()
    v1 = client.CoreV1Api()
    
    # 删除指定注释
    body = {
        "metadata": {
            "annotations": {
                "cluster-autoscaler.kubernetes.io/safe-to-evict": None
            }
        }
    }
    v1.patch_namespaced_pod(name=pod_name, namespace=namespace, body=body)

步骤2:在DAG任务中绑定回调

在你的DAG定义里,给需要保护的任务添加回调函数:

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime

with DAG(
    dag_id="protected_task_dag",
    start_date=datetime(2024, 1, 1),
    schedule_interval="@daily",
    catchup=False
) as dag:
    protected_task = PythonOperator(
        task_id="my_protected_task",
        python_callable=your_task_function,  # 替换成你的任务函数
        on_execute_callback=add_safe_to_evict,
        on_success_callback=remove_safe_to_evict,
        on_failure_callback=remove_safe_to_evict
    )

步骤3:验证权限(可选)

如果遇到Kubernetes API权限问题,需要给Composer的Airflow ServiceAccount添加pods/annotate权限:

  1. 创建一个ClusterRole:
    apiVersion: rbac.authorization.k8s.io/v1
    kind: ClusterRole
    metadata:
      name: airflow-pod-annotator
    rules:
    - apiGroups: [""]
      resources: ["pods"]
      verbs: ["patch", "get"]
    
  2. 将该角色绑定到Airflow的ServiceAccount:
    apiVersion: rbac.authorization.k8s.io/v1
    kind: ClusterRoleBinding
    metadata:
      name: airflow-pod-annotator-binding
    subjects:
    - kind: ServiceAccount
      name: composer-<YOUR_ENV_NAME>-airflow-sa
      namespace: composer-<YOUR_ENV_NAME>-<SUFFIX>
    roleRef:
      kind: ClusterRole
      name: airflow-pod-annotator
      apiGroup: rbac.authorization.k8s.io
    

验证方法

运行任务后,用以下命令查看Pod注释:

kubectl describe pod <worker-pod-name> -n <composer-namespace>

在Annotations部分应该能看到cluster-autoscaler.kubernetes.io/safe-to-evict: false;任务完成后,该注释会被移除。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 10:40:39