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

如何通过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
    )
解决方案

问题根源

  1. Token受众不匹配:获取的token是给Azure管理API用的,AKS集群需要针对Kubernetes API的token,受众应为https://kubernetes.azure.com/.default,这直接导致401未授权。
  2. 客户端实例未复用配置:代码创建了带AKS配置的k8s_client,但后续又新建了v1 = client.CoreV1Api(),这个新实例默认连接本地K8s集群(localhost:80),引发连接被拒绝错误。
  3. SSL验证关闭:生产环境存在安全风险,必须开启。

修复步骤

  1. 修正Token受众
    替换获取token的代码:
    token = credential.get_token("https://kubernetes.azure.com/.default").token
    
  2. 复用配置好的客户端
    删除v1 = client.CoreV1Api(),改用已配置的k8s_client调用API:
    pods = k8s_client.list_namespaced_pod(namespace)
    
  3. 开启SSL验证(生产环境)
    移除configuration.verify_ssl = False和requests.get中的verify=False,若有证书问题可指定AKS的CA证书路径。
  4. 给托管身份授权AKS权限
    在Azure门户进入AKS集群的「访问控制(IAM)」,添加角色分配,将「Azure Kubernetes Service RBAC Cluster Admin」或自定义Pod删除权限角色分配给Function的托管身份。
  5. 环境变量管理
    将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 04:57:32