如何通过Azure Function重启AKS Pod?连接集群失败求助
问题描述
我正在做一个项目,需要实现每次调用Azure Function时,让这个Function重启Kubernetes里的Pod。我尝试用SDK在Function里写脚本,但连不上集群:已经成功获取Azure凭证,却遇到连接被拒绝、401未授权的错误,相关错误日志和Function代码如下,希望能解决连接问题,成功重启Pod。
错误日志
DefaultAzureCredential acquired a token from EnvironmentCredential [2024-09-20T03:41:52.511Z] C:\Users\Andres_Sanchez1\AppData\Roaming\Python\Python311\site-packages\urllib3\connectionpool.py:1061: InsecureRequestWarning: Unverified HTTPS request is being made to host 'abccall-demo-lgq05iz2.hcp.eastus.azmk8s.io'. Adding certificate verification is strongly advised. See: https://urllib3.readthedocs.io/en/1.26.x/advanced-usage.html#ssl-warnings [2024-09-20T03:41:52.514Z] warnings.warn( [2024-09-20T03:42:04.837Z] Retrying (Retry(total=0, connect=None, read=None, redirect=None, status=None)) after connection broken by 'NewConnectionError('<urllib3.connection.HTTPConnection object at 0x000002835461E050>: Failed to establish a new connection: [WinError 10061] No connection could be made because the target machine actively refused it')': /api/v1/namespaces/eks_demo/pods [2024-09-20T03:42:08.923Z] 401 [2024-09-20T03:42:08.923Z] Unexpected error: HTTPConnectionPool(host='localhost', port=80): Max retries exceeded with url: /api/v1/namespaces/eks_demo/pods (Caused by NewConnectionError('<urllib3.connection.HTTPConnection object at 0x000002835461E9D0>: Failed to establish a new connection: [WinError 10061] No connection could be made because the target machine actively refused it')) [2024-09-20T03:42:08.928Z] {'kind': 'Status', 'apiVersion': 'v1', 'metadata': {}, 'status': 'Failure', 'message': 'Unauthorized', 'reason': 'Unauthorized', 'code': 401}
Function代码
import logging import os import azure.functions as func from kubernetes import client, config from kubernetes.client.rest import ApiException from azure.identity import DefaultAzureCredential import requests app = func.FunctionApp(http_auth_level=func.AuthLevel.FUNCTION) @app.route(route="http_trigger") def http_trigger(req: func.HttpRequest) -> func.HttpResponse: logging.info('Processing a request to restart a Kubernetes pod using managed identity.') # Get pod name and namespace from query parameters pod_name = req.params.get('pod_name') namespace = req.params.get('namespace', 'default') if not pod_name: return func.HttpResponse( "Please pass a pod_name in the query string.", status_code=400 ) try: # Use ManagedIdentityCredential for managed identity authentication credential = DefaultAzureCredential() # Get the AKS API server endpoint from environment variables aks_api_server = 'https://abccall-demo-lgq05iz2.hcp.eastus.azmk8s.io' # Set this in your Function App settings if not aks_api_server: return func.HttpResponse("AKS_API_SERVER environment variable is not set.", status_code=500) # Get the access token token = credential.get_token("https://management.azure.com/.default").token # Create a Kubernetes API client configuration configuration = client.Configuration() configuration.host = aks_api_server configuration.verify_ssl = False # Consider enabling SSL verification in production configuration.api_key = {"authorization": f"Bearer {token}"} response = requests.get(f"{aks_api_server}/api/v1/namespaces/default/pods", verify=False) # Change verify=True in production print(response.status_code) print(response.json()) # Create the Kubernetes API client k8s_client = client.CoreV1Api(client.ApiClient(configuration)) v1 = client.CoreV1Api() pods = v1.list_namespaced_pod(namespace) for pod in pods.items: print(f"Pod Name: {pod.metadata.name}") # Delete the pod to trigger a restart logging.info(f"Attempting to restart pod {pod_name} in namespace {namespace}.") k8s_client.delete_namespaced_pod(pod_name, namespace, body=client.V1DeleteOptions()) return func.HttpResponse( f"Pod {pod_name} in namespace {namespace} has been restarted.", status_code=200 ) except ApiException as e: logging.error(f"Exception when calling CoreV1Api->delete_namespaced_pod: {e}") return func.HttpResponse( f"Error: {str(e)}", status_code=e.status ) except Exception as e: logging.error(f"Unexpected error: {e}") return func.HttpResponse( f"Error: {str(e)}", status_code=500 )
解决方案
问题根源
- Token受众不匹配:获取的token是给Azure管理API用的,AKS集群需要针对Kubernetes API的token,受众应为
https://kubernetes.azure.com/.default,这直接导致401未授权。 - 客户端实例未复用配置:代码创建了带AKS配置的
k8s_client,但后续又新建了v1 = client.CoreV1Api(),这个新实例默认连接本地K8s集群(localhost:80),引发连接被拒绝错误。 - SSL验证关闭:生产环境存在安全风险,必须开启。
修复步骤
- 修正Token受众
替换获取token的代码:token = credential.get_token("https://kubernetes.azure.com/.default").token - 复用配置好的客户端
删除v1 = client.CoreV1Api(),改用已配置的k8s_client调用API:pods = k8s_client.list_namespaced_pod(namespace) - 开启SSL验证(生产环境)
移除configuration.verify_ssl = False和requests.get中的verify=False,若有证书问题可指定AKS的CA证书路径。 - 给托管身份授权AKS权限
在Azure门户进入AKS集群的「访问控制(IAM)」,添加角色分配,将「Azure Kubernetes Service RBAC Cluster Admin」或自定义Pod删除权限角色分配给Function的托管身份。 - 环境变量管理
将AKS地址改为从环境变量读取,避免硬编码:aks_api_server = os.getenv("AKS_API_SERVER")
修复后完整代码
import logging import os import azure.functions as func from kubernetes import client, config from kubernetes.client.rest import ApiException from azure.identity import DefaultAzureCredential app = func.FunctionApp(http_auth_level=func.AuthLevel.FUNCTION) @app.route(route="http_trigger") def http_trigger(req: func.HttpRequest) -> func.HttpResponse: logging.info('Processing a request to restart a Kubernetes pod using managed identity.') pod_name = req.params.get('pod_name') namespace = req.params.get('namespace', 'default') if not pod_name: return func.HttpResponse( "请在查询字符串中传入pod_name参数。", status_code=400 ) try: credential = DefaultAzureCredential() aks_api_server = os.getenv("AKS_API_SERVER") if not aks_api_server: return func.HttpResponse("未设置AKS_API_SERVER环境变量。", status_code=500) # 获取针对K8s API的token token = credential.get_token("https://kubernetes.azure.com/.default").token configuration = client.Configuration() configuration.host = aks_api_server # 生产环境请移除下面这行,开启SSL验证 # configuration.verify_ssl = False configuration.api_key = {"authorization": f"Bearer {token}"} k8s_client = client.CoreV1Api(client.ApiClient(configuration)) # 使用配置好的客户端查询Pod pods = k8s_client.list_namespaced_pod(namespace) for pod in pods.items: logging.info(f"Pod名称: {pod.metadata.name}") logging.info(f"尝试重启命名空间{namespace}中的Pod {pod_name}。") k8s_client.delete_namespaced_pod(pod_name, namespace, body=client.V1DeleteOptions()) return func.HttpResponse( f"命名空间{namespace}中的Pod {pod_name}已重启。", status_code=200 ) except ApiException as e: logging.error(f"调用CoreV1Api->delete_namespaced_pod时出错: {e}") return func.HttpResponse( f"错误: {str(e)}", status_code=e.status ) except Exception as e: logging.error(f"意外错误: {e}") return func.HttpResponse( f"错误: {str(e)}", status_code=500 )
内容的提问来源于stack exchange,提问作者Andres Sanchez
相关产品推荐
相关产品推荐

