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.yamlpod_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权限:
- 创建一个ClusterRole:
apiVersion: rbac.authorization.k8s.io/v1 kind: ClusterRole metadata: name: airflow-pod-annotator rules: - apiGroups: [""] resources: ["pods"] verbs: ["patch", "get"] - 将该角色绑定到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

