Lesson 07 of 9 · Part 3 — Go and Kubernetes
client-go: talking to Kubernetes
Talk to Kubernetes from Go with client-go: load kubeconfig or in-cluster config, list and get objects, patch them safely, watch for changes, and use informers and listers so long-running tools don't hammer the API server.
Connect: kubeconfig or in-cluster
package main
import (
"context"
"fmt"
"os"
"path/filepath"
"time"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/clientcmd"
)
func config() (*rest.Config, error) {
if cfg, err := rest.InClusterConfig(); err == nil {
return cfg, nil // running in a pod
}
path := os.Getenv("KUBECONFIG")
if path == "" {
home, _ := os.UserHomeDir()
path = filepath.Join(home, ".kube", "config")
}
return clientcmd.BuildConfigFromFlags("", path)
}
func main() {
cfg, err := config()
if err != nil {
fmt.Fprintln(os.Stderr, "error: load config:", err)
os.Exit(1)
}
cfg.Timeout = 15 * time.Second
cs, err := kubernetes.NewForConfig(cfg)
if err != nil {
fmt.Fprintln(os.Stderr, "error:", err)
os.Exit(1)
}
ctx := context.Background()
nodes, err := cs.CoreV1().Nodes().List(ctx, metav1.ListOptions{})
if err != nil {
fmt.Fprintln(os.Stderr, "error: list nodes:", err)
os.Exit(1)
}
for _, n := range nodes.Items {
ready := "Unknown"
for _, c := range n.Status.Conditions {
if c.Type == "Ready" {
ready = string(c.Status)
}
}
fmt.Printf("%-30s ready=%-7s kubelet=%s zone=%s\n", n.Name, ready,
n.Status.NodeInfo.KubeletVersion, n.Labels["topology.kubernetes.io/zone"])
}
}
The clientset has a typed client per API group and version (CoreV1(), AppsV1(), NetworkingV1()), so fields and types are checked by the compiler. Exec credential plugins in kubeconfig (EKS aws eks get-token, GKE gke-gcloud-auth-plugin) work transparently.
client-go is your phone line to the cluster's front desk. You can ask for a list once (List), ask to be called back whenever something changes (Watch), or hire an assistant who keeps an up-to-date notebook of everything (informer), so you read the notebook instead of phoning every minute.
Read: selectors and pagination
pods, err := cs.CoreV1().Pods("").List(ctx, metav1.ListOptions{
LabelSelector: "app=web",
FieldSelector: "status.phase!=Running",
Limit: 500, // page size; follow pods.Continue for the next page
})
Filter on the server with label and field selectors instead of fetching everything. Large clusters need pagination (Limit + Continue).
Write: patch, don't overwrite
patch := []byte(`{"metadata":{"labels":{"audit.example.com/checked":"true"}}}`)
_, err = cs.CoreV1().Nodes().Patch(ctx, "n1", types.MergePatchType, patch, metav1.PatchOptions{})
if apierrors.IsNotFound(err) {
// the node is gone; usually fine for an audit tool
}
Imports: k8s.io/apimachinery/pkg/types and apierrors "k8s.io/apimachinery/pkg/api/errors". Server-side apply (Apply methods with ApplyConfiguration types) is the modern way for tools that own a set of fields.
Watch and informers
A raw Watch streams changes, but you must handle reconnects and resync yourself. Informers do it for you and keep a local cache:
factory := informers.NewSharedInformerFactory(cs, 10*time.Minute)
podInformer := factory.Core().V1().Pods()
podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
UpdateFunc: func(oldObj, newObj any) {
p := newObj.(*corev1.Pod)
if p.Status.Phase == corev1.PodFailed {
fmt.Printf("pod %s/%s failed\n", p.Namespace, p.Name)
}
},
})
factory.Start(ctx.Done())
factory.WaitForCacheSync(ctx.Done())
lister := podInformer.Lister() // reads from the cache, no API calls
pending, _ := lister.List(labels.Everything())
(Imports: k8s.io/client-go/informers, k8s.io/client-go/tools/cache, corev1 "k8s.io/api/core/v1", k8s.io/apimachinery/pkg/labels.) Objects from the cache are shared: never modify them; DeepCopy() first.
RBAC for your tool
A tool running in a pod gets the permissions of its ServiceAccount. Give it exactly the verbs it uses:
apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRole
metadata:
name: kaudit
rules:
- apiGroups: [ "" ]
resources: [ "nodes", "pods" ]
verbs: [ "get", "list", "watch" ]
- apiGroups: [ "" ]
resources: [ "nodes" ]
verbs: [ "patch" ]
Check with kubectl auth can-i list pods --as system:serviceaccount:tools:kaudit -A.
Versions
Match client-go's minor version to your clusters' (client-go v0.34.x ↔ Kubernetes 1.34); one minor either way works for common APIs. Upgrade client-go with the clusters, and pin apimachinery and api to the same version.
Try it: kaudit nodes, for real
- Create a kind cluster with 3 nodes and build the program above; compare its output with
kubectl get nodes -L topology.kubernetes.io/zone. - Add a
podssubcommand that lists Pending and CrashLoopBackOff pods across namespaces, using a field selector for the phase. - Add
--labelto patch anaudit.example.com/checked=truelabel onto every Ready node; check withkubectl get nodes --show-labels. - Run it in-cluster as a Job with a ServiceAccount and the ClusterRole above; remove
patchfrom the role and read the error. - Turn the pods report into a long-running watcher with an informer that prints failed pods as they happen.
Going deeper: dynamic and generated clients
- The dynamic client works with any resource (including CRDs) as unstructured maps: good for generic tools.
- For CRDs you own, generate typed clients (code-generator) or use controller-runtime's client (lesson 08), which handles typed objects and caching for you.
- Set
QPSandBurston the rest config deliberately for tools that make many calls; the defaults are conservative.
Recap
- Config from in-cluster or kubeconfig; a typed clientset per API group.
- Selectors and pagination on reads; patch (or server-side apply) for writes;
apierrorshelpers for status. - Informers + listers for long-running tools; never mutate cached objects.
- Least-privilege RBAC and matching versions.
This site is a public version of my personal engineering knowledge hub. It intentionally excludes confidential company information and internal operational details.