学习目标
学完本章,你应该能够:
- 讲清 client-go 的几大客户端怎么选:Clientset、RESTClient、DynamicClient、DiscoveryClient 分别解决什么场景,90% 的控制器为什么是 “Clientset + Informer + Workqueue” 这套组合。
- 写出一个带超时、限流、dry-run 的运维 CLI:用 cobra 搭骨架,用
clientcmd加载 kubeconfig,用FieldSelector在 API 端过滤,而不是把几万个对象拉回本地再筛。 - 写出带冲突重试的 CRUD:理解 K8s 的乐观锁(
resourceVersion)和errors.IsConflict重试、Limit分页、FieldManager幂等这些生产必踩点。 - 写出一个能跑的生产级 Controller:List-Watch → 入队 → Worker 消费 → Reconcile 的完整闭环,理解 Informer 本地缓存、
WaitForCacheSync、Workqueue 重试与幂等。 - 配好 Leader Election 避免脑裂:讲清
LeaseDuration > RenewDeadline > RetryPeriod的硬约束,以及失主为什么必须os.Exit。
前置知识:
- Go 基础:goroutine、
channel、context、接口、defer。 - 一点 Kubernetes 常识:Pod / ConfigMap / Deployment 是什么,kubectl 基本命令(
get/apply/delete)。 - 一个能连的 K8s 集群(或 kind/minikube),以及本地的
~/.kube/config。
本章你会动手做的事:
- 把文章里的
clean-evictedCLI 跑起来,先加--dry-run看它会删哪些 Pod,再真删一遍。 - 把
PodAnnotateController编译运行,故意改一个 Pod 让它缺注解,观察控制器自动补上created-by=myctl。 - 起两个副本的控制器进程,看 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/apimachinery、k8s.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,比如自定义的
Tenant、AppConfig资源) - 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 常见的使用场景有:
- 运维脚本:批量查询、修改、删除资源。
- 自定义 CLI 工具:为团队封装内部使用的 kubectl 插件或独立工具。
- 控制器与 Operator:监听资源变化并调谐期望状态。
- 平台集成:将 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)
}
}
踩坑提示:
RunE比Run好。Run不返回 error,你只能 panic 或者 os.Exit(1),错误处理不优雅。SilenceUsage: true必加。否则每次运行出错都打印一大坨 usage,烦死人。PersistentFlagsvsFlags:要所有子命令都能用的放PersistentFlags,只当前命令用的放Flags。kubeconfig/namespace 用 Persistent,dry-run 用 Flags。- 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
}
踩坑提示:
errors.IsConflict一定要处理。多控制器并发改同一个对象太常见了,不重试就间歇性失败。GracePeriodSeconds要谨慎。强杀(=0)适合调试,生产建议 30s 以上,给应用收尾时间。StatefulSet 的 Pod 尤其不能强杀。Limit字段别忘。我见过有人List全集群 Pod 把 API Server 内存打爆的,集群直接雪崩。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
}
踩坑提示:
- GVR 的
Resource字段是复数小写(tenants),不是Tenant也不是tenant。搞错直接 404。 NestedInt64三返回值必须检查found。否则字段不存在时返回 0,你以为是默认值其实是没设。- DynamicClient 没有 Lister。要用 Informer 缓存 CRD,得自己用
dynamicinformer包,类型安全不如 typed client。 - 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)
}
}
踩坑提示:
processNextItem必须在所有 return 路径前defer c.queue.Done(item)。漏一处队列计数就乱,最后整个 controller 卡死。reconcile必须幂等。同一个 key 可能被处理多次,写操作要判断 “是否真的需要改”。- UpdateFunc 一定要做变化检测。不然 Pod status 心跳每 5 秒触发一次,你 reconcile 一直空转,写 API 把限流打满。
- 重试上限必须有。否则一个 “找不到的 Pod” 会让队列永久重试,越积越多。
WaitForCacheSync不能省。我见过有人 controller 启动后 1 秒就 reconcile,结果缓存还没建好,处理了一堆 “not found” 错误。
processNextItem 和 reconcile 的协作可以用一张状态图收尾,这也是 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
}
踩坑提示:
- 三个时长
LeaseDuration > RenewDeadline > RetryPeriod是硬约束。我一般用 15s/10s/2s,太短(比如 5s/3s/1s)网络抖动就切主,太长(60s/30s/10s)leader 挂了 standby 要等一分钟才接管。 OnStoppedLeading必须os.Exit。失主后如果不退出,业务代码可能还在跑,造成两个实例同时写。- identity 要稳定。别用随机 UUID,否则 Pod 重启后老 lease 没释放,新 identity 抢不到主。
- 跑在集群里时 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
}
踩坑提示:
- DiscoveryClient 结果要缓存。每次调用都打 API Server,频繁调用会被限流。建议启动时调一次,结果缓存到内存。
ServerResourcesForGroupVersion偶尔返回部分错误(某个 API server 不可达),用*ErrGroupDiscoveryFailed包装,要遍历err.(*discovery.ErrGroupDiscoveryFailed).Groups分别处理。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 时长配置。本笔记里的代码都是我踩完坑之后整理出来的版本,希望能帮你少走点弯路。
自测题与动手练习
自测题(合上书能答出来,才算懂):
- Clientset、DynamicClient、DiscoveryClient、RESTClient 分别适合什么场景?为什么写控制器大多是 “Clientset + Informer + Workqueue” 这套?
- K8s 资源更新是乐观锁,并发改同一个对象会怎样?
UpdateConfigMapWithRetry为什么每次重试都要先重新Get? - Informer 的
UpdateFunc为什么一定要做变化检测(比如比较resourceVersion)?不检测会出什么问题? - Workqueue 的
Done漏调用、reconcile不幂等,分别会导致什么后果? - Leader Election 的三个时长
LeaseDuration / RenewDeadline / RetryPeriod必须满足什么不等式?失主后为什么必须os.Exit而不是继续跑?
动手练习(建议真做一遍):
- 把
clean-evicted命令跑起来:先--dry-run看会删哪些 Pod,再真删;故意把FieldSelector去掉,感受一下"全量拉回本地再筛"和"API 端过滤"的带宽差别。 - 跑
PodAnnotateController,手动给某个 Pod 删掉created-by注解,观察它在几秒内被补回来;然后 ctrl-C 重启,确认WaitForCacheSync期间没有刷 “not found” 错误。 - 起两个副本的控制器进程,用
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 Election:
LeaseDuration > RenewDeadline > RetryPeriod是硬约束,失主必须退出,避免双写脑裂。
下一篇笔记将结合 AIOps 场景,展示如何用 client-go 实现智能运维脚本。