client-go 入门

2025-11-05T14:11:02+08:00 | 19分钟阅读 | 更新于 2025-11-05T14:11:02+08:00

@

学习目标

学完本章,你应该能够:

  1. 讲清 client-go 的几大客户端怎么选:Clientset、RESTClient、DynamicClient、DiscoveryClient 分别解决什么场景,90% 的控制器为什么是 “Clientset + Informer + Workqueue” 这套组合。
  2. 写出一个带超时、限流、dry-run 的运维 CLI:用 cobra 搭骨架,用 clientcmd 加载 kubeconfig,用 FieldSelector 在 API 端过滤,而不是把几万个对象拉回本地再筛。
  3. 写出带冲突重试的 CRUD:理解 K8s 的乐观锁(resourceVersion)和 errors.IsConflict 重试、Limit 分页、FieldManager 幂等这些生产必踩点。
  4. 写出一个能跑的生产级 Controller:List-Watch → 入队 → Worker 消费 → Reconcile 的完整闭环,理解 Informer 本地缓存、WaitForCacheSync、Workqueue 重试与幂等。
  5. 配好 Leader Election 避免脑裂:讲清 LeaseDuration > RenewDeadline > RetryPeriod 的硬约束,以及失主为什么必须 os.Exit

前置知识

  • Go 基础:goroutine、channelcontext、接口、defer
  • 一点 Kubernetes 常识:Pod / ConfigMap / Deployment 是什么,kubectl 基本命令(get / apply / delete)。
  • 一个能连的 K8s 集群(或 kind/minikube),以及本地的 ~/.kube/config

本章你会动手做的事

  1. 把文章里的 clean-evicted CLI 跑起来,先加 --dry-run 看它会删哪些 Pod,再真删一遍。
  2. PodAnnotateController 编译运行,故意改一个 Pod 让它缺注解,观察控制器自动补上 created-by=myctl
  3. 起两个副本的控制器进程,看 Lease 选主:只有一个在 reconcile,kill 掉 leader,另一个在 LeaseDuration 内接管。

client-go 简介

client-go 是 Kubernetes 官方提供的 Go 语言客户端库,是与 K8s API Server 交互的标准方式。无论是简单的运维脚本、复杂的控制器(Controller)、Operator,还是自定义的 CLI 工具,底层都大量依赖 client-go。源码仓库为 k8s.io/client-go,通常与 k8s.io/apimachineryk8s.io/api 一起用。

我第一次用 client-go 是为了写一个批量清理 Evicted Pod 的脚本,之前用 shell + kubectl 写了一版,跑倒是能跑,但每次循环调 kubectl 启动开销巨大,1000 个 Pod 能跑两分钟。换成 client-go 之后 3 秒搞定。从此我就再没回去过。

client-go 最大的优势不是性能,是 类型安全。kubectl 是命令行工具,参数都是字符串,写错了运行时才报;client-go 是 Go 代码,参数类型编译期就给你卡住。运维工具上线生产,类型安全这点能省你一半的 review 时间。

client-go 整体架构

client-go 的核心模块包括:

模块作用
Clientset提供类型化的 REST 客户端,支持各类 K8s 资源的 CRUD
RESTClient底层 REST 调用封装,灵活但使用较繁琐
DynamicClient动态客户端,无需预先知道资源类型,适合处理 CRD
DiscoveryClient发现集群支持的 API 组、版本和资源
Informers基于 List-Watch 的本地缓存机制,高效监听资源变化
Workqueue事件队列,配合 Informer 实现可靠的事件处理
Lister只读本地缓存查询接口,性能高

理解这些模块的关系,是编写高效、稳定 K8s 程序的基础。

我按使用频率排个序,方便你心里有数:

  • 90% 场景:Clientset + Informer + Workqueue(写 Controller 的标配)
  • 5% 场景:DynamicClient(操作 CRD,比如自定义的 TenantAppConfig 资源)
  • 3% 场景:DiscoveryClient(写 kubectl-like 工具,列资源类型)
  • 2% 场景:裸 RESTClient(特殊接口 clientset 没封装时才用,极少)

下面这张图把各模块的协作关系串起来,先看一眼有个整体印象:

flowchart LR
    subgraph 用户代码
        U[运维脚本 / CLI / Controller]
    end
    U --> CS[Clientset
类型化客户端] U --> DC[DynamicClient
操作 CRD] U --> DIS[DiscoveryClient
发现 API 资源] U --> RC[RESTClient
裸 REST 调用] CS --> INF[Informers
List-Watch 本地缓存] INF --> L[Lister
只读本地查询] INF --> WQ[Workqueue
事件队列] WQ --> REC[Reconcile
调谐业务逻辑] CS --> LE[Leader Election
高可用选主]

这张图在讲什么:上层用户代码按场景选不同客户端;写控制器时 Informer 把 List-Watch 的结果缓存到本地,变更事件经 Workqueue 交给 Reconcile 处理,最后用 Leader Election 保证同一时刻只有一个实例在干活。

下面我会挨个把代码写出来,能跑的那种。

使用场景

client-go 常见的使用场景有:

  1. 运维脚本:批量查询、修改、删除资源。
  2. 自定义 CLI 工具:为团队封装内部使用的 kubectl 插件或独立工具。
  3. 控制器与 Operator:监听资源变化并调谐期望状态。
  4. 平台集成:将 K8s 能力集成到 PaaS、CICD 或 AIOps 平台。

我自己做过的实际项目:定时清理 Evicted Pod 的 CronJob、给所有命名空间打 cost-center 标签的批量工具、监听 ConfigMap 变化自动 reload nginx 的 Controller、把 K8s Event 流转发到企业微信的 sidecar。共同点是 需要稳定、可测试、能跑在生产,shell 脚本搞不定。

Golang 命令行工具开发

使用 Go 开发命令行工具时,常用 cobra 框架。基本结构如下:

package main

import (
    "fmt"
    "github.com/spf13/cobra"
)

var rootCmd = &cobra.Command{
    Use:   "myctl",
    Short: "一个自定义的 K8s 运维工具",
}

var listCmd = &cobra.Command{
    Use:   "list",
    Short: "列出 Pod",
    Run: func(cmd *cobra.Command, args []string) {
        fmt.Println("列出所有 Pod...")
    },
}

func main() {
    rootCmd.AddCommand(listCmd)
    if err := rootCmd.Execute(); err != nil {
        panic(err)
    }
}

cobra 支持子命令、参数解析、配置文件读取,是构建 kubectl 插件风格工具的首选。

上面这版太玩具了,下面给一个完整能用的 CLI 模板,带 kubeconfig flag、namespace flag、dry-run flag,结构跟 kubectl 类似:

package main

import (
    "context"
    "fmt"
    "os"
    "path/filepath"

    "github.com/spf13/cobra"
    "k8s.io/client-go/kubernetes"
    "k8s.io/client-go/tools/clientcmd"
    "k8s.io/client-go/util/homedir"
)

// 全局 flag,所有子命令共用
var (
    kubeconfig string
    namespace  string
    dryRun     bool
)

var rootCmd = &cobra.Command{
    Use:   "myctl",
    Short: "我自己的 K8s 运维 CLI",
    Long:  "myctl 是基于 client-go 的运维工具,提供 list/clean 等子命令。",
    // SilenceUsage 让出错时不打印 usage,避免刷屏
    SilenceUsage: true,
}

var listCmd = &cobra.Command{
    Use:     "pods",
    Aliases: []string{"pod", "po"},
    Short:   "列出指定 namespace 的 Pod",
    RunE: func(cmd *cobra.Command, args []string) error {
        // 1. 加载 kubeconfig,构造 clientset
        config, err := clientcmd.BuildConfigFromFlags("", kubeconfig)
        if err != nil {
            return fmt.Errorf("加载 kubeconfig 失败: %w", err)
        }
        clientset, err := kubernetes.NewForConfig(config)
        if err != nil {
            return fmt.Errorf("创建 clientset 失败: %w", err)
        }

        // 2. 拉取 Pod 列表
        // 关键:永远带 context,并设超时,别让一个慢请求卡死整个 CLI
        ctx, cancel := context.WithTimeout(context.Background(), 30)
        defer cancel()
        pods, err := clientset.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{})
        if err != nil {
            return fmt.Errorf("list pods 失败: %w", err)
        }

        // 3. 输出
        fmt.Printf("%-40s %-10s %-10s %-5s\n", "NAME", "READY", "STATUS", "RESTARTS")
        for _, p := range pods.Items {
            ready, total := 0, len(p.Spec.Containers)
            for _, cs := range p.Status.ContainerStatuses {
                if cs.Ready {
                    ready++
                }
            }
            restarts := 0
            if len(p.Status.ContainerStatuses) > 0 {
                restarts = int(p.Status.ContainerStatuses[0].RestartCount)
            }
            fmt.Printf("%-40s %d/%d        %-10s %-5d\n",
                p.Name, ready, total, p.Status.Phase, restarts)
        }
        return nil
    },
}

var cleanEvictedCmd = &cobra.Command{
    Use:   "clean-evicted",
    Short: "清理所有 Evicted 状态的 Pod",
    RunE: func(cmd *cobra.Command, args []string) error {
        config, err := clientcmd.BuildConfigFromFlags("", kubeconfig)
        if err != nil {
            return err
        }
        clientset, err := kubernetes.NewForConfig(config)
        if err != nil {
            return err
        }

        ctx, cancel := context.WithTimeout(context.Background(), 60)
        defer cancel()
        // 关键:用 FieldSelector 在 API 端过滤,省带宽
        pods, err := clientset.CoreV1().Pods("").List(ctx, metav1.ListOptions{
            FieldSelector: "status.phase=Failed",
        })
        if err != nil {
            return err
        }

        deleted := 0
        for _, p := range pods.Items {
            if p.Status.Reason != "Evicted" {
                continue
            }
            if dryRun {
                fmt.Printf("[dry-run] 将删除 %s/%s\n", p.Namespace, p.Name)
                continue
            }
            if err := clientset.CoreV1().Pods(p.Namespace).Delete(
                ctx, p.Name, metav1.DeleteOptions{}); err != nil {
                // 单个失败不影响整体,记日志继续
                fmt.Fprintf(os.Stderr, "删除 %s/%s 失败: %v\n", p.Namespace, p.Name, err)
                continue
            }
            deleted++
        }
        fmt.Printf("共清理 %d 个 Evicted Pod\n", deleted)
        return nil
    },
}

func init() {
    // kubeconfig 默认走 ~/.kube/config,跟 kubectl 一致
    if home := homedir.HomeDir(); home != "" {
        rootCmd.PersistentFlags().StringVar(&kubeconfig, "kubeconfig",
            filepath.Join(home, ".kube", "config"), "kubeconfig 文件路径")
    } else {
        rootCmd.PersistentFlags().StringVar(&kubeconfig, "kubeconfig", "", "kubeconfig 文件路径")
    }
    rootCmd.PersistentFlags().StringVarP(&namespace, "namespace", "n", "default", "命名空间")
    // 关键:写操作命令一律带 dry-run,养成习惯
    cleanEvictedCmd.Flags().BoolVar(&dryRun, "dry-run", false, "只打印不执行")

    rootCmd.AddCommand(listCmd)
    rootCmd.AddCommand(cleanEvictedCmd)
}

func main() {
    if err := rootCmd.Execute(); err != nil {
        // cobra 自己会打印错误信息,这里直接退出就行
        os.Exit(1)
    }
}

踩坑提示:

  1. RunERunRun 不返回 error,你只能 panic 或者 os.Exit(1),错误处理不优雅。
  2. SilenceUsage: true 必加。否则每次运行出错都打印一大坨 usage,烦死人。
  3. PersistentFlags vs Flags:要所有子命令都能用的放 PersistentFlags,只当前命令用的放 Flags。kubeconfig/namespace 用 Persistent,dry-run 用 Flags。
  4. kubeconfig 路径要做默认值兜底,不然用户每次都得敲一长串路径。

client-go 深度实战

创建 Clientset

import (
    "k8s.io/client-go/kubernetes"
    "k8s.io/client-go/tools/clientcmd"
)

func main() {
    config, err := clientcmd.BuildConfigFromFlags("", clientcmd.RecommendedHomeFile)
    if err != nil {
        panic(err)
    }
    clientset, err := kubernetes.NewForConfig(config)
    if err != nil {
        panic(err)
    }
    // 使用 clientset 操作资源
}

这版能用但太简陋,生产代码得加几样东西:QPS 限流、超时、TLS 配置、User-Agent。我下面给个生产可用的工厂函数:

package kube

import (
    "fmt"
    "time"

    "k8s.io/client-go/kubernetes"
    "k8s.io/client-go/rest"
    "k8s.io/client-go/tools/clientcmd"
)

// BuildClientset 从 kubeconfig 或 in-cluster config 构建 clientset
// kubeconfigPath 为空时尝试 in-cluster(Pod 里跑的情况)
func BuildClientset(kubeconfigPath string) (*kubernetes.Clientset, error) {
    var config *rest.Config
    var err error

    if kubeconfigPath == "" {
        // 1. 先试 in-cluster config(运行在 K8s 里时)
        config, err = rest.InClusterConfig()
        if err != nil {
            // 2. 不在集群里就回退到 ~/.kube/config
            config, err = clientcmd.BuildConfigFromFlags("", clientcmd.RecommendedHomeFile)
            if err != nil {
                return nil, fmt.Errorf("既无法用 in-cluster 也找不到 kubeconfig: %w", err)
            }
        }
    } else {
        config, err = clientcmd.BuildConfigFromFlags("", kubeconfigPath)
        if err != nil {
            return nil, fmt.Errorf("加载 kubeconfig %s 失败: %w", kubeconfigPath, err)
        }
    }

    // 关键:调高 QPS 和 Burst,默认 5/10 太小,写 Controller 经常被限流
    config.QPS = 50
    config.Burst = 100
    // 单次请求超时,别让它无限等
    config.Timeout = 30 * time.Second
    // 标识自己,便于 API Server 端识别和限流
    config.UserAgent = "myctl/v1.0.0"

    return kubernetes.NewForConfig(config)
}

调谐改查

操作示例方法
查询clientset.CoreV1().Pods(ns).List(ctx, metav1.ListOptions{})
创建clientset.CoreV1().Pods(ns).Create(ctx, pod, metav1.CreateOptions{})
更新clientset.CoreV1().Pods(ns).Update(ctx, pod, metav1.UpdateOptions{})
删除clientset.CoreV1().Pods(ns).Delete(ctx, name, metav1.DeleteOptions{})

光看表格没啥用,我把完整代码写出来。下面这段是带错误处理、context、retry 逻辑的 CRUD 全家桶:

package kube

import (
    "context"
    "fmt"
    "time"

    corev1 "k8s.io/api/core/v1"
    "k8s.io/apimachinery/pkg/api/errors"
    metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
    "k8s.io/client-go/kubernetes"
    "k8s.io/client-go/util/retry"
)

// ListPods 列出指定 namespace 的所有 Pod
func ListPods(ctx context.Context, clientset *kubernetes.Clientset, namespace string) ([]corev1.Pod, error) {
    pods, err := clientset.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{
        Limit: 500, // 关键:大数据集一定要分页,别一次 list 几万个 Pod 把 API Server 拖垮
    })
    if err != nil {
        return nil, fmt.Errorf("list pods in %s: %w", namespace, err)
    }
    return pods.Items, nil
}

// CreateConfigMap 创建 ConfigMap,已存在则返回原对象(幂等)
func CreateConfigMap(ctx context.Context, clientset *kubernetes.Clientset, namespace string, cm *corev1.ConfigMap) (*corev1.ConfigMap, error) {
    created, err := clientset.CoreV1().ConfigMaps(namespace).Create(ctx, cm, metav1.CreateOptions{
        FieldManager: "myctl", // 关键:server-side apply 的字段管理器,便于 kubectl diff 看是谁改的
    })
    if err != nil {
        if errors.IsAlreadyExists(err) {
            // 幂等:已存在不算错,返回当前对象
            existing, getErr := clientset.CoreV1().ConfigMaps(namespace).Get(ctx, cm.Name, metav1.GetOptions{})
            if getErr != nil {
                return nil, fmt.Errorf("get existing cm: %w", getErr)
            }
            return existing, nil
        }
        return nil, fmt.Errorf("create configmap: %w", err)
    }
    return created, nil
}

// UpdateConfigMapWithRetry 更新 ConfigMap,带冲突重试
// 关键:K8s 资源更新是乐观锁(resourceVersion),并发改同一个对象会冲突,必须重试
func UpdateConfigMapWithRetry(ctx context.Context, clientset *kubernetes.Clientset, namespace, name string, mutate func(*corev1.ConfigMap)) error {
    return retry.OnError(retry.DefaultRetry, errors.IsConflict, func() error {
        // 每次重试都要重新 Get,拿到最新的 resourceVersion
        current, err := clientset.CoreV1().ConfigMaps(namespace).Get(ctx, name, metav1.GetOptions{})
        if err != nil {
            return err
        }
        // 在最新对象上应用变更
        mutate(current)
        _, err = clientset.CoreV1().ConfigMaps(namespace).Update(ctx, current, metav1.UpdateOptions{
            FieldManager: "myctl",
        })
        return err
    })
}

// DeletePod 删除 Pod,支持优雅终止期
func DeletePod(ctx context.Context, clientset *kubernetes.Clientset, namespace, name string, gracePeriodSeconds int64) error {
    err := clientset.CoreV1().Pods(namespace).Delete(ctx, name, metav1.DeleteOptions{
        GracePeriodSeconds: &gracePeriodSeconds, // 0 = 立即强杀,>0 = 给容器 SIGTERM 时间
    })
    if err != nil {
        if errors.IsNotFound(err) {
            return nil // 已经没了,幂等返回
        }
        return fmt.Errorf("delete pod %s/%s: %w", namespace, name, err)
    }
    return nil
}

踩坑提示:

  1. errors.IsConflict 一定要处理。多控制器并发改同一个对象太常见了,不重试就间歇性失败。
  2. GracePeriodSeconds 要谨慎。强杀(=0)适合调试,生产建议 30s 以上,给应用收尾时间。StatefulSet 的 Pod 尤其不能强杀。
  3. Limit 字段别忘。我见过有人 List 全集群 Pod 把 API Server 内存打爆的,集群直接雪崩。
  4. FieldManager 名字别乱起。一旦上线就别改,否则之前管的那批字段会变成 “orphaned”,server-side apply 会出问题。

操作 CRD:DynamicClient

CRD(自定义资源)没法用 Clientset 直接操作,因为类型不是预定义的。这时候就得用 DynamicClient,它操作的是 unstructured.Unstructured,啥类型都能塞。

package kube

import (
    "context"
    "fmt"

    metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
    "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
    "k8s.io/apimachinery/pkg/runtime/schema"
    "k8s.io/client-go/dynamic"
)

// ListTenants 列出所有 Tenant CRD 实例
// 假设我们有个 CRD:tenants.example.com/v1alpha1
func ListTenants(ctx context.Context, dynClient dynamic.Interface) ([]*unstructured.Unstructured, error) {
    gvr := schema.GroupVersionResource{
        Group:    "example.com",
        Version:  "v1alpha1",
        Resource: "tenants", // 关键:这里是复数 resource 名,不是 kind
    }

    list, err := dynClient.Resource(gvr).Namespace("").List(ctx, metav1.ListOptions{})
    if err != nil {
        return nil, fmt.Errorf("list tenants: %w", err)
    }
    return list.Items, nil
}

// PatchTenantReplicas 修改 Tenant 的 spec.replicas 字段
func PatchTenantReplicas(ctx context.Context, dynClient dynamic.Interface,
    namespace, name string, replicas int64) error {
    gvr := schema.GroupVersionResource{
        Group:    "example.com",
        Version:  "v1alpha1",
        Resource: "tenants",
    }

    // 用 JSON merge patch,比 strategic patch 通用
    // 关键:DynamicClient 没有类型化方法,所有字段都得用 string 路径
    patch := map[string]interface{}{
        "spec": map[string]interface{}{
            "replicas": replicas,
        },
    }

    _, err := dynClient.Resource(gvr).Namespace(namespace).Patch(
        ctx, name,
        metav1.MergePatchType, // 或者 PatchType(fieldManager) 用 server-side apply
        mustToJSON(patch),
        metav1.PatchOptions{FieldManager: "myctl"},
    )
    return err
}

// GetTenantSpecReplicas 从 unstructured 中安全取 spec.replicas
// 关键:unstructured.NestedInt64 返回 (value, found, err) 三值,found 一定要检查
func GetTenantSpecReplicas(obj *unstructured.Unstructured) (int64, error) {
    replicas, found, err := unstructured.NestedInt64(obj.Object, "spec", "replicas")
    if err != nil {
        return 0, fmt.Errorf("spec.replicas 字段类型错误: %w", err)
    }
    if !found {
        return 0, fmt.Errorf("spec.replicas 不存在")
    }
    return replicas, nil
}

踩坑提示:

  1. GVR 的 Resource 字段是复数小写tenants),不是 Tenant 也不是 tenant。搞错直接 404。
  2. NestedInt64 三返回值必须检查 found。否则字段不存在时返回 0,你以为是默认值其实是没设。
  3. DynamicClient 没有 Lister。要用 Informer 缓存 CRD,得自己用 dynamicinformer 包,类型安全不如 typed client。
  4. CRD schema 变更要小心。新增字段 OK,删字段或改类型会导致老对象反序列化失败,集群里一片报错。

控制器与 Informer

控制器遵循 Kubernetes 的声明式控制循环:

观察(Watch) -> 分析(Diff) -> 执行(Act) -> 重试(Retry)

下面的时序图把"事件怎么从 API Server 流到你的 Reconcile"画清楚:

sequenceDiagram
    participant API as API Server
    participant INF as Informer(本地缓存)
    participant WQ as Workqueue
    participant W as Worker
    participant C as Reconcile
    API->>INF: List 全量 + Watch 增量事件
    INF->>WQ: Add(key) 事件入队
    WQ->>W: Get() 取出一个 key
    W->>C: reconcile(key)
    C->>API: Get / Update 对齐期望态
    C-->>WQ: Done(key),失败则 AddRateLimited 退避重试

这张图在讲什么:Informer 负责跟 API Server 保持 List-Watch,把变更以 namespace/name 这样的 key 丢进 Workqueue;Worker 从队列取 key 交给 reconcile,成功后 Done,失败按 rate limiter 退避重试。注意 reconcile 拿数据优先走本地缓存,不直接打 API Server。

使用 Informer 监听资源变化,配合 Workqueue 实现异步处理:

informer := factory.Core().V1().Pods().Informer()
informer.AddEventHandler(cache.ResourceEventHandlerFuncs{
    AddFunc: func(obj interface{}) {
        key, _ := cache.MetaNamespaceKeyFunc(obj)
        workqueue.Add(key)
    },
    UpdateFunc: func(old, new interface{}) {
        key, _ := cache.MetaNamespaceKeyFunc(new)
        workqueue.Add(key)
    },
    DeleteFunc: func(obj interface{}) {
        key, _ := cache.DeletionHandlingMetaNamespaceKeyFunc(obj)
        workqueue.Add(key)
    },
})

上面这版是骨架,真正能用的 Controller 还差不少。我把完整版写出来:

package controller

import (
    "context"
    "fmt"
    "time"

    corev1 "k8s.io/api/core/v1"
    metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
    "k8s.io/apimachinery/pkg/util/runtime"
    "k8s.io/apimachinery/pkg/util/wait"
    "k8s.io/client-go/informers"
    "k8s.io/client-go/kubernetes"
    "k8s.io/client-go/tools/cache"
    "k8s.io/client-go/util/workqueue"
    "k8s.io/klog/v2"
)

// PodAnnotateController 给所有缺少 owner 注解的 Pod 补上 created-by=myctl 注解
// 一个最简但完整可生产的 Controller 示例
type PodAnnotateController struct {
    clientset *kubernetes.Clientset
    queue     workqueue.RateLimitingInterface
    informer  cache.SharedIndexInformer
}

func NewPodAnnotateController(clientset *kubernetes.Clientset) *PodAnnotateController {
    // 1. 用 SharedInformerFactory 创建 informer
    factory := informers.NewSharedInformerFactory(clientset, 30*time.Minute)
    podInformer := factory.Core().V1().Pods().Informer()

    c := &PodAnnotateController{
        clientset: clientset,
        queue:     workqueue.NewNamedRateLimitingQueue(
            workqueue.DefaultControllerRateLimiter(),
            "pod-annotate", // 关键:给队列起名,metrics 里能看到
        ),
        informer: podInformer,
    }

    // 2. 注册事件处理
    podInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
        AddFunc: func(obj interface{}) {
            key, err := cache.MetaNamespaceKeyFunc(obj)
            if err != nil {
                // 有坑:key 生成失败别 panic,否则 controller 直接挂
                klog.Errorf("生成 key 失败: %v", err)
                return
            }
            c.queue.Add(key)
        },
        UpdateFunc: func(old, new interface{}) {
            // 关键:UpdateFunc 触发很频繁,Pod status 心跳每几秒刷一次
            // 这里简单比较 resourceVersion,避免无意义的重复 reconcile
            oldPod, ok1 := old.(*corev1.Pod)
            newPod, ok2 := new.(*corev1.Pod)
            if !ok1 || !ok2 {
                return
            }
            if oldPod.ResourceVersion == newPod.ResourceVersion {
                return
            }
            key, err := cache.MetaNamespaceKeyFunc(new)
            if err != nil {
                return
            }
            c.queue.Add(key)
        },
        DeleteFunc: func(obj interface{}) {
            key, err := cache.DeletionHandlingMetaNamespaceKeyFunc(obj)
            if err != nil {
                return
            }
            c.queue.Add(key)
        },
    })

    return c
}

// Run 启动 controller,workerNum 个 worker 并发消费队列
func (c *PodAnnotateController) Run(ctx context.Context, workers int) error {
    defer runtime.HandleCrash()
    defer c.queue.ShutDown()

    klog.Info("启动 PodAnnotateController")

    // 1. 等 informer 缓存同步完成
    // 关键:不等缓存就 reconcile 会读不到对象,导致空指针
    go c.informer.Run(ctx.Done())
    if !cache.WaitForCacheSync(ctx.Done(), c.informer.HasSynced) {
        return fmt.Errorf("informer 缓存同步超时")
    }
    klog.Info("informer 缓存同步完成")

    // 2. 启动 worker
    for i := 0; i < workers; i++ {
        go wait.UntilWithContext(ctx, c.runWorker, time.Second)
    }

    // 3. 等待退出信号
    <-ctx.Done()
    klog.Info("controller 退出")
    return nil
}

func (c *PodAnnotateController) runWorker(ctx context.Context) {
    for c.processNextItem(ctx) {
    }
}

// processNextItem 取一个 key 处理,处理失败按 rate limiter 退避重试
func (c *PodAnnotateController) processNextItem(ctx context.Context) bool {
    item, shutdown := c.queue.Get()
    if shutdown {
        return false
    }
    // 关键:Done 必须调用,否则队列不释放计数,最后会卡死
    defer c.queue.Done(item)

    key := item.(string)
    if err := c.reconcile(ctx, key); err != nil {
        // 重试有上限:超过次数直接丢弃,避免毒丸消息把队列塞满
        if c.queue.NumRequeues(key) < 5 {
            klog.Warningf("reconcile %s 失败,将重试: %v", key, err)
            c.queue.AddRateLimited(key)
        } else {
            klog.Errorf("reconcile %s 重试 5 次仍失败,放弃: %v", key, err)
        }
        return true
    }
    c.queue.Forget(key) // 成功就清掉重试计数
    return true
}

// reconcile 真正的业务逻辑
func (c *PodAnnotateController) reconcile(ctx context.Context, key string) error {
    namespace, name, err := cache.SplitMetaNamespaceKey(key)
    if err != nil {
        return fmt.Errorf("invalid key %s: %w", key, err)
    }

    // 1. 从本地缓存读 Pod(不走 API,快)
    obj, exists, err := c.informer.GetIndexer().GetByKey(key)
    if err != nil {
        return fmt.Errorf("get from cache: %w", err)
    }
    if !exists {
        return nil // Pod 被删了,没东西可做
    }

    pod, ok := obj.(*corev1.Pod)
    if !ok {
        return fmt.Errorf("invalid pod type in cache")
    }

    // 2. 业务判断:是否需要补注解
    if pod.Annotations != nil && pod.Annotations["created-by"] == "myctl" {
        return nil
    }

    // 3. 执行更新(带冲突重试)
    // 关键:必须重新 Get 拿最新版本,缓存里的可能已经过期
    return retryOnConflict(func() error {
        current, err := c.clientset.CoreV1().Pods(namespace).Get(ctx, name, metav1.GetOptions{})
        if err != nil {
            return err
        }
        if current.Annotations == nil {
            current.Annotations = make(map[string]string)
        }
        current.Annotations["created-by"] = "myctl"
        _, err = c.clientset.CoreV1().Pods(namespace).Update(ctx, current, metav1.UpdateOptions{
            FieldManager: "pod-annotate-controller",
        })
        return err
    })
}

// 用法
func ExampleRun() {
    clientset, _ := BuildClientset("")
    c := NewPodAnnotateController(clientset)
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()
    if err := c.Run(ctx, 3); err != nil {
        panic(err)
    }
}

踩坑提示:

  1. processNextItem 必须在所有 return 路径前 defer c.queue.Done(item)。漏一处队列计数就乱,最后整个 controller 卡死。
  2. reconcile 必须幂等。同一个 key 可能被处理多次,写操作要判断 “是否真的需要改”。
  3. UpdateFunc 一定要做变化检测。不然 Pod status 心跳每 5 秒触发一次,你 reconcile 一直空转,写 API 把限流打满。
  4. 重试上限必须有。否则一个 “找不到的 Pod” 会让队列永久重试,越积越多。
  5. WaitForCacheSync 不能省。我见过有人 controller 启动后 1 秒就 reconcile,结果缓存还没建好,处理了一堆 “not found” 错误。

processNextItemreconcile 的协作可以用一张状态图收尾,这也是 Workqueue 重试逻辑的核心:

flowchart TD
    A[processNextItem 取 key] --> B{队列已关闭?}
    B -- 是 --> Z[返回 false 退出循环]
    B -- 否 --> C[defer queue.Done key]
    C --> D{reconcile 成功?}
    D -- 是 --> E[queue.Forget 清重试计数]
    D -- 否 --> F{重试次数 < 5?}
    F -- 是 --> G[AddRateLimited 退避后重试]
    F -- 否 --> H[放弃 记 error
避免毒丸塞满队列]

这张图在讲什么:每个 key 都有重试计数;成功就 Forget,失败且未超上限就退避重试,超过上限直接丢弃。漏掉 Done 会卡死队列,缺了重试上限会让一个坏 key 无限重试。

选举机制

在高可用控制器部署中,通常需要 Leader Election 保证同一时刻只有一个实例执行 reconcile。client-go 提供了 tools/leaderelection 包,支持基于 Lease 或 ConfigMap/Endpoint 的选主:

flowchart LR
    L[LeaseDuration 15s
lease 有效期] --> R[RenewDeadline 10s
leader 续约截止] R --> P[RetryPeriod 2s
非 leader 检查间隔] A[实例A 成为 Leader
定时续约 lease] -->|lease 过期未续约| B[实例B 抢到 lease
成为新 Leader] B -->|OnStoppedLeading| A2[实例A os.Exit 退出]

这张图在讲什么:三个时长必须严格满足 LeaseDuration > RenewDeadline > RetryPeriod;leader 在 RenewDeadline 内没续约成功,lease 过期,standby 在下一个 RetryPeriod 抢到主。失主的实例必须退出,否则会出现两个实例同时写。

lock := &resourcelock.LeaseLock{
    LeaseMeta: metav1.ObjectMeta{Name: "my-controller", Namespace: "default"},
    Client:    clientset.CoordinationV1(),
    LockConfig: resourcelock.ResourceLockConfig{
        Identity: hostname,
    },
}

这是骨架,下面给完整版,包含选主成功回调、失败处理、健康检查:

package controller

import (
    "context"
    "os"
    "time"

    metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
    "k8s.io/client-go/kubernetes"
    "k8s.io/client-go/tools/leaderelection"
    "k8s.io/client-go/tools/leaderelection/resourcelock"
    "k8s.io/klog/v2"
)

// RunWithLeaderElection 启动带选主的 controller
// 同一时刻只有 leader 实例在跑 reconcile,其他实例 standby
func RunWithLeaderElection(ctx context.Context, clientset *kubernetes.Clientset, runFunc func(ctx context.Context)) error {
    // 1. 生成唯一 identity,通常用 hostname + pid
    // 关键:identity 必须全局唯一,否则两个实例以为自己是同一个,选主逻辑全乱
    hostname, err := os.Hostname()
    if err != nil {
        return fmt.Errorf("get hostname: %w", err)
    }
    identity := fmt.Sprintf("%s-%d", hostname, os.Getpid())

    // 2. 构造 LeaseLock
    lock := &resourcelock.LeaseLock{
        LeaseMeta: metav1.ObjectMeta{
            Name:      "my-controller-leader",
            Namespace: "default",
        },
        Client: clientset.CoordinationV1(),
        LockConfig: resourcelock.ResourceLockConfig{
            Identity: identity,
        },
    }

    // 3. 配置选主回调
    leaderelection.RunOrDie(ctx, leaderelection.LeaderElectionConfig{
        Lock: lock,

        // 关键:LeaseDuration 必须 > RenewDeadline > RetryPeriod
        // 这三个时长是选主的核心,调错了要么脑裂要么频繁切主
        LeaseDuration: 15 * time.Second, // lease 有效期,超过没续约认为 leader 死了
        RenewDeadline: 10 * time.Second, // leader 续约的截止时间
        RetryPeriod:   2 * time.Second,  // 非 leader 检查 lease 的间隔

        Callbacks: leaderelection.LeaderCallbacks{
            OnStartedLeading: func(ctx context.Context) {
                klog.Infof("实例 %s 成为 leader,开始工作", identity)
                // 关键:在子 ctx 里跑业务,失主时 ctx 会被 cancel,业务自然停
                runFunc(ctx)
            },
            OnStoppedLeading: func() {
                klog.Infof("实例 %s 不再是 leader,退出", identity)
                // 失主退出,由 deployment 重启
                os.Exit(0)
            },
            OnNewLeader: func(currentIdentity string) {
                if currentIdentity == identity {
                    return
                }
                klog.Infof("当前 leader 是 %s,我作为 standby 等待", currentIdentity)
            },
        },

        // 关键:开 Watch 防止 API Server 通知丢了导致脑裂
        // 不开的话如果 RenewDeadline 内 API Server 没响应,leader 自己不知道丢了 lease
        WatchDog: leaderelection.NewLeaderHealthzAdaptor(time.Second * 20),
        Name:     "my-controller",
    })

    return nil
}

踩坑提示:

  1. 三个时长 LeaseDuration > RenewDeadline > RetryPeriod 是硬约束。我一般用 15s/10s/2s,太短(比如 5s/3s/1s)网络抖动就切主,太长(60s/30s/10s)leader 挂了 standby 要等一分钟才接管。
  2. OnStoppedLeading 必须 os.Exit。失主后如果不退出,业务代码可能还在跑,造成两个实例同时写。
  3. identity 要稳定。别用随机 UUID,否则 Pod 重启后老 lease 没释放,新 identity 抢不到主。
  4. 跑在集群里时 RBAC 要给 Coordination leases 的读写权限。我第一次部署忘了加 Role,controller 启动就 panic,半天查不出原因。

DiscoveryClient 示例

DiscoveryClient 用来列集群支持哪些 API 资源,类似 kubectl api-resources。写 CLI 工具时很有用——比如让 CLI 自动适配不同 K8s 版本。

package kube

import (
    "context"
    "fmt"

    "k8s.io/client-go/discovery"
)

// ListClusterAPIs 列出集群支持的所有 API 资源
func ListClusterAPIs(ctx context.Context, disco discovery.DiscoveryInterface) error {
    // 1. 拿到所有 API Group
    groupList, err := disco.ServerGroups()
    if err != nil {
        return fmt.Errorf("server groups: %w", err)
    }

    // 2. 每个 group 拿 PreferredVersion 的资源列表
    for _, group := range groupList.Groups {
        if len(group.Versions) == 0 {
            continue
        }
        preferred := group.PreferredVersion
        resources, err := disco.ServerResourcesForGroupVersion(preferred.GroupVersion)
        if err != nil {
            fmt.Printf("获取 %s 资源失败: %v\n", preferred.GroupVersion, err)
            continue
        }

        fmt.Printf("\n=== %s ===\n", preferred.GroupVersion)
        for _, r := range resources.APIResources {
            // 关键:跳过子资源(如 pod/exec),只列真正能 list/get 的
            if r.Verbs.Contains("list") && r.Verbs.Contains("get") {
                fmt.Printf("  %-30s kind=%-20s namespaced=%v\n",
                    r.Name, r.Kind, r.Namespaced)
            }
        }
    }
    return nil
}

// IsResourceSupported 检查集群是否支持某个 GVR,写自适应 CLI 时常用
func IsResourceSupported(ctx context.Context, disco discovery.DiscoveryInterface,
    group, version, resource string) (bool, error) {
    resources, err := disco.ServerResourcesForGroupVersion(group + "/" + version)
    if err != nil {
        return false, err
    }
    for _, r := range resources.APIResources {
        if r.Name == resource {
            return true, nil
        }
    }
    return false, nil
}

踩坑提示:

  1. DiscoveryClient 结果要缓存。每次调用都打 API Server,频繁调用会被限流。建议启动时调一次,结果缓存到内存。
  2. ServerResourcesForGroupVersion 偶尔返回部分错误(某个 API server 不可达),用 *ErrGroupDiscoveryFailed 包装,要遍历 err.(*discovery.ErrGroupDiscoveryFailed).Groups 分别处理。
  3. PreferredVersion 不一定是最新的。有些集群只装了 v1beta1,PreferredVersion 也是 v1beta1,别假设它是 v1。

基于 client-go 的基础设施自动化脚本

借助 client-go,可以将常见运维操作固化为 Go 程序。例如:

  • 批量给指定命名空间添加标签。
  • 定时清理 Evicted 状态的 Pod。
  • 根据注解自动为 Service 创建 Ingress。
  • 导出集群资源配置为 GitOps 仓库。

这类脚本比 Shell 更易于维护、测试和分发,也更容易与 AIOps 平台集成。前面 CLI 版的 clean-evicted 命令稍微改造(去掉 cobra 包装、改成定时调用)就能当 CronJob 用,核心逻辑都是 List + 过滤 + Delete 三步。我建议这类脚本统一用 Go 写,带 dry-run 和审计日志,shell 写到第三次就开始失控。

总结

client-go 是 Kubernetes 生态的编程入口。掌握 Clientset、Informer、Workqueue、Leader Election 等核心概念,是开发控制器、Operator 和 AIOps 工具的关键。下一篇笔记将结合 AIOps 场景,展示如何使用 client-go 实现智能运维脚本。

我自己的体会是:client-go 入门不难,但写好一个 Controller 至少要踩过这几个坑——并发更新冲突、Informer 缓存同步、Workqueue 重试逻辑、Leader Election 时长配置。本笔记里的代码都是我踩完坑之后整理出来的版本,希望能帮你少走点弯路。


自测题与动手练习

自测题(合上书能答出来,才算懂)

  1. Clientset、DynamicClient、DiscoveryClient、RESTClient 分别适合什么场景?为什么写控制器大多是 “Clientset + Informer + Workqueue” 这套?
  2. K8s 资源更新是乐观锁,并发改同一个对象会怎样?UpdateConfigMapWithRetry 为什么每次重试都要先重新 Get
  3. Informer 的 UpdateFunc 为什么一定要做变化检测(比如比较 resourceVersion)?不检测会出什么问题?
  4. Workqueue 的 Done 漏调用、reconcile 不幂等,分别会导致什么后果?
  5. Leader Election 的三个时长 LeaseDuration / RenewDeadline / RetryPeriod 必须满足什么不等式?失主后为什么必须 os.Exit 而不是继续跑?

动手练习(建议真做一遍)

  1. clean-evicted 命令跑起来:先 --dry-run 看会删哪些 Pod,再真删;故意把 FieldSelector 去掉,感受一下"全量拉回本地再筛"和"API 端过滤"的带宽差别。
  2. PodAnnotateController,手动给某个 Pod 删掉 created-by 注解,观察它在几秒内被补回来;然后 ctrl-C 重启,确认 WaitForCacheSync 期间没有刷 “not found” 错误。
  3. 起两个副本的控制器进程,用 kubectl get lease 看只有一个在 renewTime 更新;kill 掉 leader 进程,数一数另一个最多多久(约 LeaseDuration)接管。

本章小结

  • 客户端选型:Clientset 写业务最常用,DynamicClient 操作 CRD,DiscoveryClient 写自适应 CLI,裸 RESTClient 仅作补充。
  • 生产级 Clientset:必须配 QPS/Burst 限流、请求超时、User-Agent;CRUD 要处理 IsConflict 重试、Limit 分页、FieldManager 幂等。
  • 控制器闭环:List-Watch → Informer 本地缓存 → Workqueue 入队 → Worker reconcile → 失败退避重试;Done 必调用、reconcile 必幂等、重试必设上限。
  • Leader ElectionLeaseDuration > RenewDeadline > RetryPeriod 是硬约束,失主必须退出,避免双写脑裂。

下一篇笔记将结合 AIOps 场景,展示如何用 client-go 实现智能运维脚本。

About Me

没什么想介绍的,一个很大众的码农…

喜欢代码,车,马,真的是 🐎

讨厌别人让我给自己的代码写注释 最厌烦别人的程序没有写注释

目标

学AI,加油!加油!