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

如何使用client-go监听Kubernetes特定命名空间的Pod事件

基于client-go实现Pod事件监听与同步的参考方向

1. 核心依赖与客户端初始化

首先引入client-go核心包,通过集群内配置或本地kubeconfig初始化K8s客户端:

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

func getClientset() (*kubernetes.Clientset, error) {
    // 集群内运行时使用InClusterConfig
    config, err := rest.InClusterConfig()
    if err != nil {
        // 本地调试时指定kubeconfig路径
        kubeconfig := "/path/to/your/kubeconfig"
        config, err = clientcmd.BuildConfigFromFlags("", kubeconfig)
        if err != nil {
            return nil, err
        }
    }
    return kubernetes.NewForConfig(config)
}

2. 用Informer实现可靠事件监听

client-go的Informer比原生Watch更稳定,自带本地缓存和重连机制,是监听资源变化的标准方案。

2.1 初始化Pod Informer

指定目标命名空间,创建Pod的SharedIndexInformer:

import (
    "context"
    corev1 "k8s.io/api/core/v1"
    metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
    "k8s.io/apimachinery/pkg/runtime"
    "k8s.io/apimachinery/pkg/watch"
    "k8s.io/client-go/tools/cache"
)

func setupPodInformer(clientset *kubernetes.Clientset, targetNS string) cache.SharedIndexInformer {
    return cache.NewSharedIndexInformer(
        &cache.ListWatch{
            ListFunc: func(opt metav1.ListOptions) (runtime.Object, error) {
                return clientset.CoreV1().Pods(targetNS).List(context.TODO(), opt)
            },
            WatchFunc: func(opt metav1.ListOptions) (watch.Interface, error) {
                return clientset.CoreV1().Pods(targetNS).Watch(context.TODO(), opt)
            },
        },
        &corev1.Pod{},
        0, // 0表示禁用自动重同步
        cache.Indexers{},
    )
}

2.2 注册事件处理函数

分别处理新增、删除、更新事件,其中更新事件需要对比新旧Pod的容器变化:

import "log"

// 处理Pod新增事件
func handlePodAdd(obj interface{}) {
    pod := obj.(*corev1.Pod)
    log.Printf("Pod created: %s/%s", pod.Namespace, pod.Name)
    syncToOtherNS(pod, "created", nil)
}

// 处理Pod删除事件
func handlePodDelete(obj interface{}) {
    pod := obj.(*corev1.Pod)
    log.Printf("Pod deleted: %s/%s", pod.Namespace, pod.Name)
    syncToOtherNS(pod, "deleted", nil)
}

// 处理Pod更新事件:检测新增容器、镜像变更
func handlePodUpdate(oldObj, newObj interface{}) {
    oldPod := oldObj.(*corev1.Pod)
    newPod := newObj.(*corev1.Pod)

    // 跳过资源版本未变化的重复事件
    if oldPod.ResourceVersion == newPod.ResourceVersion {
        return
    }

    // 构建旧容器的名称-镜像映射
    oldContainerMap := make(map[string]string)
    for _, c := range oldPod.Spec.Containers {
        oldContainerMap[c.Name] = c.Image
    }

    // 检查新增容器
    for _, newC := range newPod.Spec.Containers {
        if _, exists := oldContainerMap[newC.Name]; !exists {
            log.Printf("Pod %s added container: %s (image: %s)", newPod.Name, newC.Name, newC.Image)
            syncToOtherNS(newPod, "container_added", map[string]string{"name": newC.Name, "image": newC.Image})
        }
    }

    // 检查容器镜像变更
    for name, oldImage := range oldContainerMap {
        for _, newC := range newPod.Spec.Containers {
            if newC.Name == name && newC.Image != oldImage {
                log.Printf("Pod %s container %s image changed: %s -> %s", newPod.Name, name, oldImage, newC.Image)
                syncToOtherNS(newPod, "image_updated", map[string]string{"name": name, "old_image": oldImage, "new_image": newC.Image})
            }
        }
    }
}

将处理函数注册到Informer:

informer.AddEventHandler(cache.ResourceEventHandlerFuncs{
    AddFunc:    handlePodAdd,
    DeleteFunc: handlePodDelete,
    UpdateFunc: handlePodUpdate,
})

3. 启动Informer并运行

启动Informer控制器,等待缓存同步完成后阻塞主进程:

import "os"
import "os/signal"
import "syscall"

func main() {
    clientset, err := getClientset()
    if err != nil {
        log.Fatalf("Failed to create clientset: %v", err)
    }

    targetNS := "your-target-namespace"
    informer := setupPodInformer(clientset, targetNS)

    // 启动Informer
    stopCh := make(chan struct{})
    defer close(stopCh)
    go informer.Run(stopCh)

    // 等待缓存同步完成
    if !cache.WaitForCacheSync(stopCh, informer.HasSynced) {
        log.Fatal("Failed to sync cache")
    }

    // 监听系统信号,优雅退出
    sigCh := make(chan os.Signal, 1)
    signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
    <-sigCh
}

4. 实现跨命名空间同步逻辑

syncToOtherNS函数根据你的需求实现,比如调用目标命名空间应用的API、更新K8s资源或发送消息到队列:

import (
    "bytes"
    "encoding/json"
    "net/http"
)

func syncToOtherNS(pod *corev1.Pod, eventType string, details interface{}) {
    syncData := map[string]interface{}{
        "pod_namespace": pod.Namespace,
        "pod_name":      pod.Name,
        "event_type":    eventType,
        "details":       details,
    }

    // 替换为目标应用的实际地址
    targetURL := "http://your-app-in-other-ns:8080/sync-event"
    jsonData, _ := json.Marshal(syncData)
    resp, err := http.Post(targetURL, "application/json", bytes.NewBuffer(jsonData))
    if err != nil {
        log.Printf("Failed to sync event: %v", err)
        return
    }
    defer resp.Body.Close()
}

5. 关键注意事项

  • 权限配置:确保应用的ServiceAccount拥有目标命名空间Pod的list、watch权限,以及同步目标的访问权限。
  • 错误处理:所有API调用需添加错误捕获与日志记录,避免进程崩溃。
  • 事件去重:通过对比Pod的ResourceVersion过滤重复更新事件。
  • 资源限制:Informer缓存会占用内存,若监听大量Pod可调整重同步周期或使用索引优化查询。

内容的提问来源于stack exchange,提问作者Jananath Banuka

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 23:41:30