学习目标
学完本章,你应该能够:
- 用 client-go 写一个集群巡检程序,自动发现异常 Pod(
CrashLoopBackOff/ImagePullBackOff/Evicted)、Node 资源瓶颈与Warning事件,把"人肉kubectl“变成定时脚本。 - 把采集到的异常数据构造成 Prompt 喂给 LLM,拿到根因分析与修复建议,并正确解析其 Function Calling(Tool Calling)输出。
- 用 client-go 执行 LLM 给出的修复动作(
restart_pod/scale_deployment/cordon_node),并讲清操作白名单、dry-run、审计日志这几道安全护栏为什么缺一不可。 - 讲清"观测 → 决策 → 执行"三段式 AIOps 闭环,以及为什么只在低风险场景让程序先处理一轮、处理不了的再升级叫人。
- 总结生产落地的安全要点:特殊资源保护(StatefulSet /
kube-system)、参数范围限制、低置信度跳过、操作限流,并能说清"安全控制比 LLM 能力更重要"的含义。
前置知识:
- 上一篇的 client-go 基础:Clientset、Informer/Watch、REST 操作的基本认知。
- Go 基础,能读能写并运行一个命令行程序。
- Kubernetes 资源对象认知:Pod、Node、Deployment、Event 的作用与关系。
- LLM 的 OpenAI 兼容 API 概念(特别是 Function Calling / Tool Calling 的入参与出参结构)。
本章你会动手做的事:
- 跑通一个集群巡检程序,自动列出当前集群里的异常 Pod 与 Node 瓶颈。
- 把一段异常信息喂给 LLM,拿到"根因 + 修复建议"的 Function Calling 输出并解析出要执行的操作。
- 在 dry-run 模式下用 client-go 真正执行一次自动修复,并到审计日志里查出这一条操作记录。
实战目标
上一篇把 client-go 的架构和基本用法过了一遍,这篇我直接上手写代码。说实话 client-go 学完不拿来干点实事儿,就跟买了车不开一样——你知道它能跑,但你不知道它能帮你省多少事。
本篇要做的事:
- 写一个完整的集群巡检程序,自动发现异常 Pod、Node 资源瓶颈、Warning 事件。
- 把采集到的异常数据喂给 LLM,让它做根因分析并给出修复建议。
- 解析 LLM 的 Function Calling 输出,用 client-go 执行修复动作。
- 加上安全控制——操作白名单、dry-run 模式、审计日志,不然 LLM 乱来你哭都来不及。
- 最后讲讲生产环境跑这套东西踩的坑。
整套代码我在测试集群上跑过,K8s 1.28,Go 1.22,LLM 用的 OpenAI 兼容接口(实际跑的是通义千问,但接口一样)。你换成任何 OpenAI 兼容的 API 都行。
场景设计
假设我们要解决以下运维问题:
- 集群中经常出现
CrashLoopBackOff、ImagePullBackOff、Evicted等异常 Pod,值班同学每次都得手动kubectl delete pod清理。 - 部分 Pod CPU 或内存使用接近 limit,存在 OOM 风险,但 HPA 来不及反应。
- Node 偶发 NotReady,上面跑的 Pod 全跟着抖,告警刷屏但没人去 cordon。
- 传统告警只能通知问题,修复仍依赖人工介入,凌晨告警响应慢得要命。
说白了就是"发现问题→人去处理"这个流程太慢。我想用程序替代中间那个人——不是完全替代,是在低风险场景下让程序先处理一轮,处理不了的再叫人。
系统架构
类比:整套系统像一个"带 AI 参谋的运维值班室”。观测阶段是巡检员每隔一段时间把机房情况记下来;决策阶段是把记录交给 AI 参谋,它查手册给出处理方案;执行阶段是班长按方案动手,但必须先过"操作白名单"这道闸——不许动的机器绝不碰。三个阶段循环跑,就构成了"发现问题 → 自动处理"的闭环。
下面用一张图把"观测 → 决策 → 执行"的主控制循环画清楚:
flowchart LR
A[观测阶段
采集异常 Pod/Node/事件] --> B[决策阶段
构造 Prompt 调 LLM]
B --> C[LLM 返回
根因和修复建议]
C --> D[执行阶段
白名单校验和 dry-run]
D -->|允许| E[client-go 执行修复]
D -->|拒绝| F[记录并告警人工]
E --> G[审计日志]
F --> G
G -->|下一轮巡检| A┌──────────────────────────────────────────────────────────────┐
│ 主控制循环 (每 60s) │
│ │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ 观测阶段 │───▶│ 决策阶段 │───▶│ 执行阶段 │ │
│ │ │ │ │ │ │ │
│ │ · 异常 Pod │ │ · 构造 Prompt │ │ · 解析 FC │ │
│ │ · Node 瓶颈 │ │ · 调 LLM API │ │ · 白名单校验 │ │
│ │ · Warning 事件│ │ · 解析建议 │ │ · dry-run │ │
│ │ · 资源使用率 │ │ │ │ · 执行 + 审计 │ │
│ └──────────────┘ └──────────────┘ └──────────────┘ │
│ │ │
│ ▼ │
│ ┌──────────────┐ │
│ │ 审计日志 │ │
│ │ (JSON 文件) │ │
│ └──────────────┘ │
└──────────────────────────────────────────────────────────────┘
三个阶段的核心职责:
- 观测阶段:用 client-go 拉取 Pod、Node、Event 等数据,过滤出异常项。数据量要控制——别把整个集群状态丢给 LLM,它会懵。
- 决策阶段:把异常数据组织成结构化 Prompt,调用 LLM 获取修复建议。LLM 输出必须是 Function Calling 格式,不能让它自由发挥。
- 执行阶段:解析 LLM 的 Function Calling 输出,校验白名单,执行修复动作,写审计日志。
观测阶段:采集集群状态
异常 Pod 检测
先定义一个结构体把异常信息统一打包,后面构造 Prompt 要用:
// pkg/collector/pod.go
package collector
import (
"context"
"fmt"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
)
// PodProblem 描述一个异常 Pod 的详细信息
// 这个结构体会被序列化成 JSON 塞进 LLM 的 prompt
type PodProblem struct {
Namespace string `json:"namespace"`
Name string `json:"name"`
Phase string `json:"phase"` // Pending/Running/Failed/Unknown
Reason string `json:"reason"` // CrashLoopBackOff/ImagePullBackOff 等
RestartCount int32 `json:"restartCount"` // 重启次数,>5 就值得注意
NodeName string `json:"nodeName"` // 调度到哪个节点
Age string `json:"age"` // Pod 存活时长
LastEvent string `json:"lastEvent"` // 最近一条 Warning 事件
ContainerName string `json:"containerName"` // 出问题的容器名
Image string `json:"image"` // 镜像地址,排查 ImagePullBackOff 用
}
// ListProblematicPods 扫描全集群异常 Pod
// 返回的 list 会直接进 LLM prompt,所以要做过滤——别把 200 个正常 Pod 也塞进去
func ListProblematicPods(ctx context.Context, clientset *kubernetes.Clientset) ([]PodProblem, error) {
// ListOptions 里可以加 FieldSelector 过滤,但 K8s API Server 对 field selector 支持有限
// 常见可用的是 status.phase,但 status.containerStatuses 不支持过滤
// 所以还是全量 List 然后在客户端过滤
pods, err := clientset.CoreV1().Pods("").List(ctx, metav1.ListOptions{})
if err != nil {
return nil, fmt.Errorf("list pods failed: %w", err)
}
var problems []PodProblem
for _, pod := range pods.Items {
// 跳过 kube-system 命名空间——系统 Pod 的异常通常需要特殊处理
// 不适合让 LLM 自动修复,容易把集群搞炸
if pod.Namespace == "kube-system" {
continue
}
for _, cs := range pod.Status.ContainerStatuses {
// 检查 Waiting 状态:CrashLoopBackOff、ImagePullBackOff、ErrImagePull 等
if cs.State.Waiting != nil {
reason := cs.State.Waiting.Reason
// 只关注需要介入的异常,Pending 状态的 ContainerCreating 是正常的
if isProblematicReason(reason) {
problems = append(problems, PodProblem{
Namespace: pod.Namespace,
Name: pod.Name,
Phase: string(pod.Status.Phase),
Reason: reason,
RestartCount: cs.RestartCount,
NodeName: pod.Spec.NodeName,
Age: formatAge(pod.CreationTimestamp.Time),
ContainerName: cs.Name,
Image: cs.Image,
LastEvent: getPodLastEvent(ctx, clientset, pod.Namespace, pod.Name),
})
}
}
// 检查 Terminated 状态:容器异常退出
if cs.State.Terminated != nil && cs.State.Terminated.ExitCode != 0 {
problems = append(problems, PodProblem{
Namespace: pod.Namespace,
Name: pod.Name,
Phase: string(pod.Status.Phase),
Reason: fmt.Sprintf("ExitCode:%d", cs.State.Terminated.ExitCode),
RestartCount: cs.RestartCount,
NodeName: pod.Spec.NodeName,
Age: formatAge(pod.CreationTimestamp.Time),
ContainerName: cs.Name,
Image: cs.Image,
})
}
}
// Pod 整体失败
if pod.Status.Phase == corev1.PodFailed {
problems = append(problems, PodProblem{
Namespace: pod.Namespace,
Name: pod.Name,
Phase: string(pod.Status.Phase),
Reason: "PodFailed",
NodeName: pod.Spec.NodeName,
Age: formatAge(pod.CreationTimestamp.Time),
})
}
}
return problems, nil
}
// isProblematicReason 判断 Waiting 状态的 reason 是否需要介入
func isProblematicReason(reason string) bool {
// 这些是常见的需要人工/自动介入的异常状态
switch reason {
case "CrashLoopBackOff", "ImagePullBackOff", "ErrImagePull",
"CreateContainerError", "RunContainerError",
"InvalidImageName", "ContainerConfigError":
return true
default:
return false
}
}
// getPodLastEvent 查询 Pod 关联的最近一条 Warning 事件
// 注意:Event 的 API 在 client-go 中是 CoreV1().Events(namespace)
func getPodLastEvent(ctx context.Context, clientset *kubernetes.Clientset, ns, podName string) string {
events, err := clientset.CoreV1().Events(ns).List(ctx, metav1.ListOptions{
// 按涉及的对象名过滤,减少返回量
// 注意:FieldSelector 对 events 的支持因 K8s 版本而异
// 1.27+ 可以用 involvedObject.name
FieldSelector: fmt.Sprintf("involvedObject.name=%s,type=Warning", podName),
Limit: 1,
})
if err != nil || len(events.Items) == 0 {
return ""
}
// 返回最后一条事件的 message
return events.Items[0].Message
}
func formatAge(t metav1.Time) string {
d := metav1.Now().Sub(t)
if d.Hours() > 24 {
return fmt.Sprintf("%.0fd", d.Hours()/24)
}
if d.Hours() > 1 {
return fmt.Sprintf("%.1fh", d.Hours())
}
return fmt.Sprintf("%.1fm", d.Minutes())
}
踩坑提示:
Pods("").List跨所有 namespace,需要 cluster-admin 权限。生产环境建议用 ServiceAccount + RBAC 限定能操作的 namespace。getPodLastEvent每次查 Event 都打一次 API,Pod 多的时候 QPS 会爆。实际用的时候我会在主循环里一次性拉全集群 Event,然后在内存里做关联。RestartCount > 5这个阈值要跟业务团队商量。有些服务 OOM 重启是常态(比如 Java 堆设小了),你不能一看到重启就报异常。- 别拿
pod.Status.Phase == Pending直接当异常。Pending 可能只是还没调度到节点,等几秒就好了。我一开始就是这么写的,结果 LLM 每次都收到一堆"假阳性"。
Node 资源瓶颈检测
光看 Pod 不够,Node 层面的问题往往更致命——一个 Node 挂了,上面几十个 Pod 全跟着抖:
// pkg/collector/node.go
package collector
import (
"context"
"fmt"
"strings"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/resource"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
)
// NodeProblem 描述一个 Node 的异常状态
type NodeProblem struct {
Name string `json:"name"`
Ready string `json:"ready"` // True/False/Unknown
Conditions string `json:"conditions"` // 关键条件摘要
PodCount int `json:"podCount"` // 节点上运行的 Pod 数
CPUUsage string `json:"cpuUsage"` // CPU 分配率
MemUsage string `json:"memUsage"` // 内存分配率
DiskPressure bool `json:"diskPressure"` // 磁盘压力
MemPressure bool `json:"memPressure"` // 内存压力
PIDPressure bool `json:"pidPressure"` // PID 压力
}
// ListProblematicNodes 扫描异常节点
// 关键判断:NotReady、DiskPressure、MemPressure
func ListProblematicNodes(ctx context.Context, clientset *kubernetes.Clientset) ([]NodeProblem, error) {
nodes, err := clientset.CoreV1().Nodes().List(ctx, metav1.ListOptions{})
if err != nil {
return nil, fmt.Errorf("list nodes failed: %w", err)
}
// 同时拉所有 Pod,用于计算每个 Node 上的 Pod 数量
// 不逐 Node 查是因为那样 API 调用太多
allPods, err := clientset.CoreV1().Pods("").List(ctx, metav1.ListOptions{
FieldSelector: "status.phase=Running",
})
if err != nil {
return nil, fmt.Errorf("list pods for node check failed: %w", err)
}
// 统计每个 Node 上的 Pod 数
podCountByNode := make(map[string]int)
for _, pod := range allPods.Items {
if pod.Spec.NodeName != "" {
podCountByNode[pod.Spec.NodeName]++
}
}
var problems []NodeProblem
for _, node := range nodes.Items {
np := NodeProblem{
Name: node.Name,
PodCount: podCountByNode[node.Name],
}
// 解析 Node Conditions
// 常见 Condition: Ready, MemoryPressure, DiskPressure, PIDPressure, NetworkUnavailable
for _, cond := range node.Status.Conditions {
switch cond.Type {
case corev1.NodeReady:
if cond.Status == corev1.ConditionTrue {
np.Ready = "True"
} else {
np.Ready = string(cond.Status) // False 或 Unknown
}
case corev1.NodeDiskPressure:
np.DiskPressure = cond.Status == corev1.ConditionTrue
case corev1.NodeMemoryPressure:
np.MemPressure = cond.Status == corev1.ConditionTrue
case corev1.NodePIDPressure:
np.PIDPressure = cond.Status == corev1.ConditionTrue
}
}
// 计算 CPU 和内存分配率
// 用 requests 而不是 actual usage——这是"已承诺"的资源
// actual usage 要查 metrics-server,后面进阶部分会讲
np.CPUUsage, np.MemUsage = calcNodeAllocatable(node, allPods.Items, node.Name)
// 判断是否异常
// 1. NotReady
// 2. 有 DiskPressure / MemPressure / PIDPressure
// 3. Pod 数接近上限(默认 110)
isProblematic := np.Ready != "True" ||
np.DiskPressure || np.MemPressure || np.PIDPressure ||
np.PodCount > 100 // 接近 K8s 默认每节点 110 Pod 上限
if isProblematic {
var conds []string
if np.DiskPressure {
conds = append(conds, "DiskPressure")
}
if np.MemPressure {
conds = append(conds, "MemPressure")
}
if np.PIDPressure {
conds = append(conds, "PIDPressure")
}
np.Conditions = strings.Join(conds, ",")
problems = append(problems, np)
}
}
return problems, nil
}
// calcNodeAllocatable 计算 Node 上所有 Pod 的 requests 总和占 allocatable 的比例
func calcNodeAllocatable(node corev1.Node, pods []corev1.Pod, nodeName string) (string, string) {
// 节点的 allocatable 资源
cpuAllocatable := node.Status.Allocatable.Cpu()
memAllocatable := node.Status.Allocatable.Memory()
if cpuAllocatable.IsZero() || memAllocatable.IsZero() {
return "unknown", "unknown"
}
// 汇总该节点上所有 Running Pod 的 requests
var totalCPU, totalMem resource.Quantity
for _, pod := range pods {
if pod.Spec.NodeName != nodeName || pod.Status.Phase != corev1.PodRunning {
continue
}
for _, c := range pod.Spec.Containers {
totalCPU.Add(*c.Resources.Requests.Cpu())
totalMem.Add(*c.Resources.Requests.Memory())
}
}
// 算百分比
cpuPercent := float64(totalCPU.MilliValue()) / float64(cpuAllocatable.MilliValue()) * 100
memPercent := float64(totalMem.Value()) / float64(memAllocatable.Value()) * 100
return fmt.Sprintf("%.1f%%", cpuPercent), fmt.Sprintf("%.1f%%", memPercent)
}
踩坑提示:
- Node 的
allocatable和capacity不一样:capacity是物理总量,allocatable是扣掉 system-reserved 和 kube-reserved 之后给 Pod 用的。算分配率要用allocatable。 podCount > 100这个阈值别写死。有些节点配置高,跑 150 个 Pod 也没问题。最好做成配置项,按节点规格区分。- Node NotReady 的告警延迟很重要。K8s 默认
node-monitor-grace-period是 40 秒,也就是说 Node 挂了最快也要 40 秒才被检测到。你的巡检间隔如果设 60 秒,那从故障到检测最慢 100 秒——对某些业务来说太慢了。
Warning 事件采集
Events 是 K8s 里的"黑匣子记录仪",很多异常的根因就藏在 Warning 事件里:
// pkg/collector/event.go
package collector
import (
"context"
"fmt"
"time"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
)
// EventSummary 描述一条重要的 Warning 事件
type EventSummary struct {
Namespace string `json:"namespace"`
Name string `json:"name"`
Reason string `json:"reason"` // 事件原因,如 FailedScheduling, BackOff
Message string `json:"message"` // 事件详情
Count int32 `json:"count"` // 发生次数,多次重复的事件优先级更高
LastTime string `json:"lastTime"` // 最近一次发生时间
Type string `json:"type"` // Warning / Normal
}
// ListWarningEvents 拉取近 N 分钟内的 Warning 事件
// 重点过滤:FailedScheduling, BackOff, FailedMount, Unhealthy, EvictedThreshold
func ListWarningEvents(ctx context.Context, clientset *kubernetes.Clientset, since time.Duration) ([]EventSummary, error) {
// Events API 在 K8s 1.25+ 有变化:events.k8s.io/v1 成为正式版本
// 但 client-go 的 CoreV1().Events() 仍然可用,兼容性好
events, err := clientset.CoreV1().Events("").List(ctx, metav1.ListOptions{
// 只拉 Warning 类型,减少返回量
// 注意:FieldSelector 对 type 的支持因版本而异
// 保险起见全拉然后在客户端过滤
})
if err != nil {
return nil, fmt.Errorf("list events failed: %w", err)
}
cutoff := time.Now().Add(-since)
var summaries []EventSummary
for _, evt := range events.Items {
// 过滤时间范围
if evt.LastTimestamp.Time.Before(cutoff) {
continue
}
// 只关注 Warning
if evt.Type != "Warning" {
continue
}
// 过滤掉低价值事件
// 比如 LeaderElection 之类的是噪音,不是真正的异常
if isNoiseEvent(evt.Reason) {
continue
}
summaries = append(summaries, EventSummary{
Namespace: evt.Namespace,
Name: evt.Name,
Reason: evt.Reason,
Message: evt.Message,
Count: evt.Count,
LastTime: evt.LastTimestamp.Format("15:04:05"),
Type: evt.Type,
})
}
return summaries, nil
}
// isNoiseEvent 判断是否为噪音事件
// 这些事件虽然也是 Warning,但通常不需要自动修复
func isNoiseEvent(reason string) bool {
switch reason {
case "LeaderElectionLock", "FailedScheduling":
// FailedScheduling 噪音特别多,很多时候只是资源暂时不够,等一会儿就好了
// 但如果是持续性的 FailedScheduling 就需要关注了
// 这里简单过滤,生产环境可以加"持续 N 分钟才报"的逻辑
return true
default:
return false
}
}
踩坑提示:
- Events 在 K8s 里默认只保留 1 小时(
--event-ttl控制)。如果你要查历史事件,必须接 EventRouter 把事件落 ES/Loki。 evt.Count很有用——同一条事件重复发生 100 次和发生 1 次的紧急程度完全不同。我会在 Prompt 里把这个信息带上,LLM 会根据重复次数调整建议的优先级。- 全量 List Events 在大集群上返回量很大(几万条)。生产环境必须加
FieldSelector或者用ResourceVersion做增量 watch。
决策阶段:构造 Prompt 并调用 LLM
定义 Function Calling 工具集
先定义 LLM 能调用的"工具"——也就是它能执行的修复动作。这一步很关键,工具定义直接决定 LLM 的行为边界:
// pkg/llm/tools.go
package llm
// ToolDefinition 定义 LLM 可以调用的工具
// 这个结构体会被序列化成 OpenAI Function Calling 的 JSON schema
type ToolDefinition struct {
Name string `json:"name"`
Description string `json:"description"`
Parameters map[string]interface{} `json:"parameters"`
}
// GetAvailableTools 返回所有允许 LLM 调用的工具
// 安全核心:这里就是白名单,LLM 能干啥全由这个列表决定
// 没在这里定义的操作,LLM 就算建议了也不会执行
func GetAvailableTools() []ToolDefinition {
return []ToolDefinition{
{
Name: "restart_pod",
Description: "删除指定 Pod,由 Deployment/ReplicaSet 控制器重新创建。适用于 CrashLoopBackOff、容器僵死等场景。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"namespace": map[string]interface{}{
"type": "string",
"description": "Pod 所在的命名空间",
},
"pod_name": map[string]interface{}{
"type": "string",
"description": "Pod 名称",
},
},
"required": []string{"namespace", "pod_name"},
},
},
{
Name: "scale_deployment",
Description: "调整 Deployment 副本数。适用于扩容缓解压力或缩容释放资源。注意:缩容到 0 会停服务。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"namespace": map[string]interface{}{"type": "string", "description": "命名空间"},
"deployment": map[string]interface{}{"type": "string", "description": "Deployment 名称"},
"replicas": map[string]interface{}{"type": "integer", "description": "目标副本数", "minimum": 0, "maximum": 20},
},
"required": []string{"namespace", "deployment", "replicas"},
},
},
{
Name: "cordon_node",
Description: "将节点标记为不可调度,新 Pod 不会被调度到该节点。适用于节点异常但需要优雅排空的场景。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"node_name": map[string]interface{}{"type": "string", "description": "节点名称"},
},
"required": []string{"node_name"},
},
},
{
Name: "no_action",
Description: "当前情况不需要自动修复,或风险过高需要人工介入。LLM 判断无法安全处理时选择此项。",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"reason": map[string]interface{}{"type": "string", "description": "不执行修复的原因"},
},
"required": []string{"reason"},
},
},
}
}
注意我特意加了 no_action 工具——让 LLM 有一个"什么都不做"的选项。这很重要,没有这个选项 LLM 会倾向于"总得干点啥",而这种冲动行为在运维场景下很危险。
构造 Prompt 并调用 LLM
// pkg/llm/client.go
package llm
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"time"
"yourapp/pkg/collector"
)
// LLMClient 封装 LLM API 调用
type LLMClient struct {
APIKey string
Endpoint string // OpenAI 兼容的 endpoint,如 https://api.openai.com/v1/chat/completions
Model string // 模型名,如 gpt-4o / qwen-max
HTTP *http.Client
}
// NewLLMClient 创建客户端,超时设 30 秒
// LLM 推理慢,超时太短会频繁失败
func NewLLMClient(apiKey, endpoint, model string) *LLMClient {
return &LLMClient{
APIKey: apiKey,
Endpoint: endpoint,
Model: model,
HTTP: &http.Client{Timeout: 30 * time.Second},
}
}
// AnalysisRequest 把采集到的异常数据打包
type AnalysisRequest struct {
ProblematicPods []collector.PodProblem `json:"problematicPods"`
ProblematicNodes []collector.NodeProblem `json:"problematicNodes"`
WarningEvents []collector.EventSummary `json:"warningEvents"`
}
// AnalysisResult LLM 返回的决策结果
type AnalysisResult struct {
ToolName string `json:"tool_name"`
Parameters map[string]interface{} `json:"parameters"`
Explanation string `json:"explanation"` // LLM 给出的分析理由
Confidence string `json:"confidence"` // high/medium/low
}
// chatRequest OpenAI 兼容的请求体
type chatRequest struct {
Model string `json:"model"`
Messages []chatMessage `json:"messages"`
Tools []toolDef `json:"tools"`
ToolChoice string `json:"tool_choice"` // "auto" 让 LLM 自己决定
}
type chatMessage struct {
Role string `json:"role"`
Content interface{} `json:"content"`
ToolCalls []toolCall `json:"tool_calls,omitempty"`
}
type toolDef struct {
Type string `json:"type"` // "function"
Function ToolDefinition `json:"function"`
}
type toolCall struct {
ID string `json:"id"`
Type string `json:"type"` // "function"
Function toolCallFunc `json:"function"`
}
type toolCallFunc struct {
Name string `json:"name"`
Arguments string `json:"arguments"` // JSON 字符串
}
// Analyze 把异常数据送给 LLM,让它选一个工具执行
func (c *LLMClient) Analyze(ctx context.Context, req AnalysisRequest) (*AnalysisResult, error) {
// 构造 system prompt:定义角色 + 规则
systemPrompt := `你是一名资深 Kubernetes SRE,负责分析集群异常并给出修复建议。
# 你的职责
1. 分析异常 Pod、Node 和事件的根因
2. 从可用工具中选择最合适的修复动作
3. 如果不确定或风险高,选择 no_action
# 决策规则
- CrashLoopBackOff 且重启次数 < 5:优先 restart_pod
- ImagePullBackOff:选择 no_action(镜像问题需要人工检查),在 reason 里说明
- Node NotReady:cordon_node 阻止新 Pod 调度上去
- Node MemPressure/DiskPressure:cordon_node 并建议人工排查
- 如果同时有多个异常,选择影响最大的那个处理
- 永远不要在同一轮中执行多个操作
# 安全约束
- 不允许删除 StatefulSet 的 Pod(数据风险)
- 不允许缩容到 0
- 不允许修改 ConfigMap/Secret`
// 把异常数据序列化成 JSON 塞进 user message
// 用 JSON 而不是自然语言描述,因为 LLM 对结构化数据的理解更准确
userData, err := json.MarshalIndent(req, "", " ")
if err != nil {
return nil, fmt.Errorf("marshal analysis request failed: %w", err)
}
userPrompt := fmt.Sprintf("当前集群异常快照:\n\n```json\n%s\n```\n\n请分析并选择一个修复动作。", userData)
// 构造请求
tools := make([]toolDef, 0)
for _, t := range GetAvailableTools() {
tools = append(tools, toolDef{
Type: "function",
Function: t,
})
}
body := chatRequest{
Model: c.Model,
Messages: []chatMessage{
{Role: "system", Content: systemPrompt},
{Role: "user", Content: userPrompt},
},
Tools: tools,
ToolChoice: "auto",
}
bodyBytes, err := json.Marshal(body)
if err != nil {
return nil, fmt.Errorf("marshal request failed: %w", err)
}
// 发 HTTP 请求
httpReq, err := http.NewRequestWithContext(ctx, "POST", c.Endpoint, bytes.NewReader(bodyBytes))
if err != nil {
return nil, fmt.Errorf("create http request failed: %w", err)
}
httpReq.Header.Set("Content-Type", "application/json")
httpReq.Header.Set("Authorization", "Bearer "+c.APIKey)
resp, err := c.HTTP.Do(httpReq)
if err != nil {
return nil, fmt.Errorf("call LLM API failed: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode != 200 {
respBody, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("LLM API returned %d: %s", resp.StatusCode, string(respBody))
}
// 解析响应
var chatResp struct {
Choices []struct {
Message chatMessage `json:"message"`
} `json:"choices"`
}
if err := json.NewDecoder(resp.Body).Decode(&chatResp); err != nil {
return nil, fmt.Errorf("decode LLM response failed: %w", err)
}
if len(chatResp.Choices) == 0 {
return nil, fmt.Errorf("LLM returned no choices")
}
choice := chatResp.Choices[0]
if len(choice.Message.ToolCalls) == 0 {
// LLM 没调用工具,直接返回了文本——可能是它觉得不需要修复
// 这种情况当作 no_action 处理
return &AnalysisResult{
ToolName: "no_action",
Explanation: fmt.Sprintf("%v", choice.Message.Content),
Confidence: "low",
}, nil
}
// 取第一个 tool call
tc := choice.Message.ToolCalls[0]
var params map[string]interface{}
// arguments 是 JSON 字符串,需要再反序列化
if err := json.Unmarshal([]byte(tc.Function.Arguments), ¶ms); err != nil {
return nil, fmt.Errorf("parse tool arguments failed: %w", err)
}
return &AnalysisResult{
ToolName: tc.Function.Name,
Parameters: params,
Confidence: "medium", // 默认 medium,实际可以从 LLM 回复里解析
}, nil
}
踩坑提示:
ToolChoice: "auto"让 LLM 自己决定调不调工具。你也可以设成"required"强制它调,但那样它每次都会调一个工具,哪怕情况根本不需要修复。我推荐"auto",然后在工具列表里加no_action。- System Prompt 里的安全约束非常重要——这是你的最后一道防线。LLM 不遵守 prompt 里的规则是常态,所以执行阶段还要再做一层代码级别的白名单校验。别只靠 prompt。
arguments字段是 JSON 字符串而不是 JSON 对象,这是 OpenAI API 的设计,第一次用的人十有八九会被坑。- LLM 调用耗时不稳定,快的 2 秒慢的 20 秒。主循环里务必用 context 加超时,别让一个 LLM 调用卡住整个巡检流程。
- prompt 里塞的异常数据别太多。我试过一次塞了 50 个异常 Pod 的信息,LLM 直接懵了,选了个
no_action说"情况太复杂"。后来我做了分批,每次最多 10 个异常,效果好很多。
执行阶段:解析并调用 client-go
执行器与安全控制
这是整个系统最敏感的部分——LLM 的决策在这里变成真实操作。安全控制必须比 LLM 更可靠:
// pkg/executor/executor.go
package executor
import (
"context"
"encoding/json"
"fmt"
"log"
"os"
"time"
autoscalingv1 "k8s.io/api/autoscaling/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
)
// Executor 执行 LLM 的修复建议
type Executor struct {
Clientset *kubernetes.Clientset
DryRun bool // true = 只打印不执行,调试用
AuditLog *AuditLogger
}
// AuditEntry 审计日志条目
type AuditEntry struct {
Timestamp string `json:"timestamp"`
Action string `json:"action"`
Parameters map[string]interface{} `json:"parameters"`
DryRun bool `json:"dryRun"`
Result string `json:"result"` // success / failed / skipped
Error string `json:"error,omitempty"`
}
// AuditLogger 审计日志写入器
// 写到 JSON 文件,一行一条,方便用 jq 查询
type AuditLogger struct {
file *os.File
}
func NewAuditLogger(path string) (*AuditLogger, error) {
f, err := os.OpenFile(path, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0644)
if err != nil {
return nil, fmt.Errorf("open audit log failed: %w", err)
}
return &AuditLogger{file: f}, nil
}
func (a *AuditLogger) Log(entry AuditEntry) error {
entry.Timestamp = time.Now().Format(time.RFC3339)
data, err := json.Marshal(entry)
if err != nil {
return err
}
data = append(data, '\n')
_, err = a.file.Write(data)
return err
}
func (a *AuditLogger) Close() error {
return a.file.Close()
}
// Execute 执行 LLM 的决策
// 这是整个系统的"手"——所有操作都从这里过
func (e *Executor) Execute(ctx context.Context, toolName string, params map[string]interface{}) error {
// 第一步:白名单校验(代码级别,不信任 LLM 遵守 prompt)
if !isAllowedTool(toolName) {
err := fmt.Errorf("tool %s is not in whitelist", toolName)
e.logAudit(toolName, params, "skipped", err)
return err
}
// 第二步:dry-run 模式
if e.DryRun {
msg := fmt.Sprintf("[DRY-RUN] would execute %s with params: %v", toolName, params)
log.Println(msg)
e.logAudit(toolName, params, "skipped", nil)
return nil
}
// 第三步:执行
var err error
switch toolName {
case "restart_pod":
err = e.restartPod(ctx, params)
case "scale_deployment":
err = e.scaleDeployment(ctx, params)
case "cordon_node":
err = e.cordonNode(ctx, params)
case "no_action":
// 不执行任何操作,记录原因即可
reason, _ := params["reason"].(string)
log.Printf("LLM chose no_action: %s", reason)
e.logAudit(toolName, params, "success", nil)
return nil
default:
err = fmt.Errorf("unknown tool: %s", toolName)
}
if err != nil {
e.logAudit(toolName, params, "failed", err)
return err
}
e.logAudit(toolName, params, "success", nil)
return nil
}
// isAllowedTool 代码级白名单
// 跟 LLM prompt 里的约束是双保险
func isAllowedTool(name string) bool {
allowed := map[string]bool{
"restart_pod": true,
"scale_deployment": true,
"cordon_node": true,
"no_action": true,
}
return allowed[name]
}
// restartPod 删除 Pod,让控制器重建
func (e *Executor) restartPod(ctx context.Context, params map[string]interface{}) error {
namespace, ok := params["namespace"].(string)
if !ok || namespace == "" {
return fmt.Errorf("missing namespace parameter")
}
podName, ok := params["pod_name"].(string)
if !ok || podName == "" {
return fmt.Errorf("missing pod_name parameter")
}
// 安全检查:不允许删除 kube-system 的 Pod
if namespace == "kube-system" {
return fmt.Errorf("refuse to restart pod in kube-system namespace")
}
// 安全检查:先查 Pod 属于什么——如果是 StatefulSet 的 Pod,拒绝
// StatefulSet Pod 删除后虽然会重建,但有序性和 PVC 关联可能出问题
pod, err := e.Clientset.CoreV1().Pods(namespace).Get(ctx, podName, metav1.GetOptions{})
if err != nil {
return fmt.Errorf("get pod %s/%s failed: %w", namespace, podName, err)
}
for _, owner := range pod.OwnerReferences {
if owner.Kind == "StatefulSet" {
return fmt.Errorf("refuse to restart StatefulSet pod %s/%s (data risk)", namespace, podName)
}
}
// 执行删除
// grace period 设 30 秒,给应用优雅退出时间
// 如果 Pod 已经 CrashLoopBackOff 了,大概率也收不到 SIGTERM,但设了总比不设好
gracePeriod := int64(30)
err = e.Clientset.CoreV1().Pods(namespace).Delete(ctx, podName, metav1.DeleteOptions{
GracePeriodSeconds: &gracePeriod,
})
if err != nil {
return fmt.Errorf("delete pod %s/%s failed: %w", namespace, podName, err)
}
log.Printf("restarted pod %s/%s", namespace, podName)
return nil
}
// scaleDeployment 调整 Deployment 副本数
func (e *Executor) scaleDeployment(ctx context.Context, params map[string]interface{}) error {
namespace, _ := params["namespace"].(string)
deployment, _ := params["deployment"].(string)
// JSON 数字反序列化后是 float64,需要转成 int32
replicasFloat, ok := params["replicas"].(float64)
if !ok {
return fmt.Errorf("missing or invalid replicas parameter")
}
replicas := int32(replicasFloat)
// 安全检查:不允许缩容到 0
if replicas <= 0 {
return fmt.Errorf("refuse to scale deployment to 0")
}
// 安全检查:不允许超过 20 副本(防止 LLM 给你 scale 到 9999)
if replicas > 20 {
return fmt.Errorf("refuse to scale deployment to %d (max 20)", replicas)
}
// 用 Scale subresource 来改副本数,比 Update 整个 Deployment 更轻量
scale := &autoscalingv1.Scale{
ObjectMeta: metav1.ObjectMeta{
Name: deployment,
Namespace: namespace,
},
Spec: autoscalingv1.ScaleSpec{
Replicas: replicas,
},
}
_, err := e.Clientset.AppsV1().Deployments(namespace).UpdateScale(ctx, deployment, scale, metav1.UpdateOptions{})
if err != nil {
return fmt.Errorf("scale deployment %s/%s to %d failed: %w", namespace, deployment, replicas, err)
}
log.Printf("scaled deployment %s/%s to %d replicas", namespace, deployment, replicas)
return nil
}
// cordonNode 将节点标记为不可调度
func (e *Executor) cordonNode(ctx context.Context, params map[string]interface{}) error {
nodeName, ok := params["node_name"].(string)
if !ok || nodeName == "" {
return fmt.Errorf("missing node_name parameter")
}
// 拿到 Node 对象,修改 spec.unschedulable 字段
node, err := e.Clientset.CoreV1().Nodes().Get(ctx, nodeName, metav1.GetOptions{})
if err != nil {
return fmt.Errorf("get node %s failed: %w", nodeName, err)
}
// 已经 cordon 过了就不重复操作
if node.Spec.Unschedulable {
log.Printf("node %s already cordoned, skip", nodeName)
return nil
}
node.Spec.Unschedulable = true
_, err = e.Clientset.CoreV1().Nodes().Update(ctx, node, metav1.UpdateOptions{})
if err != nil {
return fmt.Errorf("cordon node %s failed: %w", nodeName, err)
}
log.Printf("cordoned node %s", nodeName)
return nil
}
func (e *Executor) logAudit(action string, params map[string]interface{}, result string, err error) {
entry := AuditEntry{
Action: action,
Parameters: params,
DryRun: e.DryRun,
Result: result,
}
if err != nil {
entry.Error = err.Error()
}
if e.AuditLog != nil {
if logErr := e.AuditLog.Log(entry); logErr != nil {
log.Printf("failed to write audit log: %v", logErr)
}
}
}
踩坑提示:
UpdateScale用的是 Scale subresource,比Update整个 Deployment 对象安全得多——不会覆盖其他字段的变更。但要注意 K8s 1.x 各版本的 API 路径可能有差异,appsV1是 1.9+ 的标准。cordonNode用Update改unschedulable字段有个并发问题:如果你和其他控制器同时改 Node 对象,会有冲突。更安全的做法是用Patch:
// 注意:使用 Patch 需要额外导入 "k8s.io/apimachinery/pkg/types"
patch := `{"spec":{"unschedulable":true}}`
_, err = e.Clientset.CoreV1().Nodes().Patch(ctx, nodeName, types.MergePatchType, []byte(patch), metav1.PatchOptions{})
- StatefulSet Pod 检查那段是血泪教训。我有一次让 LLM 自动重启了一个 StatefulSet 的 Pod,结果 PVC 没释放干净,新 Pod 挂不上存储,数据直接丢了。从那以后我加了这道检查,宁可漏修也不碰有状态服务。
replicas参数从 JSON 反序列化后是float64而不是int32,这是 Go 的encoding/json的行为。第一次写的时候直接params["replicas"].(int32)会 panic。- 审计日志一定要写文件,别只打 stdout。容器重启后 stdout 日志可能被截断,但文件持久化(挂 PVC)可以保留完整记录。
完整主程序
把上面三个阶段串起来:
// cmd/aiops-agent/main.go
package main
import (
"context"
"flag"
"log"
"os"
"os/signal"
"syscall"
"time"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/tools/clientcmd"
"yourapp/pkg/collector"
"yourapp/pkg/executor"
"yourapp/pkg/llm"
)
func main() {
var (
kubeconfig string
llmAPIKey string
llmEndpoint string
llmModel string
interval int
dryRun bool
auditLogPath string
)
// 命令行参数
flag.StringVar(&kubeconfig, "kubeconfig", "", "kubeconfig 文件路径,空则用 in-cluster config")
flag.StringVar(&llmAPIKey, "llm-api-key", os.Getenv("LLM_API_KEY"), "LLM API Key")
flag.StringVar(&llmEndpoint, "llm-endpoint", "https://dashscope.aliyuncs.com/compatible-mode/v1/chat/completions", "LLM API endpoint")
flag.StringVar(&llmModel, "llm-model", "qwen-max", "LLM 模型名")
flag.IntVar(&interval, "interval", 60, "巡检间隔(秒)")
flag.BoolVar(&dryRun, "dry-run", false, "dry-run 模式:只分析不执行")
flag.StringVar(&auditLogPath, "audit-log", "/var/log/aiops-audit.jsonl", "审计日志路径")
flag.Parse()
if llmAPIKey == "" {
log.Fatal("LLM_API_KEY is required")
}
// 初始化 K8s client
config, err := clientcmd.BuildConfigFromFlags("", kubeconfig)
if err != nil {
log.Fatalf("build kubeconfig failed: %v", err)
}
clientset, err := kubernetes.NewForConfig(config)
if err != nil {
log.Fatalf("create clientset failed: %v", err)
}
// 初始化 LLM client
llmClient := llm.NewLLMClient(llmAPIKey, llmEndpoint, llmModel)
// 初始化执行器
auditLogger, err := executor.NewAuditLogger(auditLogPath)
if err != nil {
log.Fatalf("create audit logger failed: %v", err)
}
defer auditLogger.Close()
exec := &executor.Executor{
Clientset: clientset,
DryRun: dryRun,
AuditLog: auditLogger,
}
if dryRun {
log.Println("=== DRY-RUN MODE: no actions will be executed ===")
}
// 信号处理:Ctrl-C 优雅退出
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer cancel()
ticker := time.NewTicker(time.Duration(interval) * time.Second)
defer ticker.Stop()
log.Printf("AIOps agent started, interval=%ds, dryRun=%v", interval, dryRun)
// 首次立即执行一次,不用等第一个 tick
runOnce(ctx, clientset, llmClient, exec, dryRun)
for {
select {
case <-ctx.Done():
log.Println("shutting down...")
return
case <-ticker.C:
runOnce(ctx, clientset, llmClient, exec, dryRun)
}
}
}
// runOnce 执行一轮完整的"观测 -> 决策 -> 执行"流程
func runOnce(ctx context.Context, clientset *kubernetes.Clientset, llmClient *llm.LLMClient, exec *executor.Executor, dryRun bool) {
log.Println("--- starting inspection cycle ---")
// ===== 阶段 1:观测 =====
// 并发采集可以加速,但为了简单先串行
// 生产环境用 errgroup 并发采集,能省一半时间
pods, err := collector.ListProblematicPods(ctx, clientset)
if err != nil {
log.Printf("collect problematic pods failed: %v", err)
// 采集失败不退出,继续采集其他维度
pods = nil
}
nodes, err := collector.ListProblematicNodes(ctx, clientset)
if err != nil {
log.Printf("collect problematic nodes failed: %v", err)
nodes = nil
}
events, err := collector.ListWarningEvents(ctx, clientset, 5*time.Minute)
if err != nil {
log.Printf("collect warning events failed: %v", err)
events = nil
}
// 如果没有任何异常,跳过 LLM 调用,省钱
if len(pods) == 0 && len(nodes) == 0 && len(events) == 0 {
log.Println("no problems found, skipping LLM analysis")
return
}
log.Printf("found %d problematic pods, %d problematic nodes, %d warning events",
len(pods), len(nodes), len(events))
// 异常太多时分批处理,避免 prompt 过长
// 每批最多 10 个 Pod + 3 个 Node + 20 个 Event
batches := batchProblems(pods, nodes, events, 10)
for i, batch := range batches {
log.Printf("processing batch %d/%d", i+1, len(batches))
// ===== 阶段 2:决策 =====
result, err := llmClient.Analyze(ctx, llm.AnalysisRequest{
ProblematicPods: batch.Pods,
ProblematicNodes: batch.Nodes,
WarningEvents: batch.Events,
})
if err != nil {
log.Printf("LLM analysis failed: %v", err)
continue
}
log.Printf("LLM decision: tool=%s, params=%v, confidence=%s",
result.ToolName, result.Parameters, result.Confidence)
if result.Explanation != "" {
log.Printf("LLM explanation: %s", result.Explanation)
}
// 低置信度的决策在非 dry-run 模式下也跳过
// 宁可漏修不可错修
if result.Confidence == "low" && !dryRun {
log.Println("confidence is low, skipping execution (would require human review)")
continue
}
// ===== 阶段 3:执行 =====
if err := exec.Execute(ctx, result.ToolName, result.Parameters); err != nil {
log.Printf("execute %s failed: %v", result.ToolName, err)
}
// 每批之间间隔 5 秒,避免操作太快触发 K8s 限流
time.Sleep(5 * time.Second)
}
}
// batchProblem 批次结构
type batchProblem struct {
Pods []collector.PodProblem
Nodes []collector.NodeProblem
Events []collector.EventSummary
}
// batchProblems 把大量异常分批
func batchProblems(pods []collector.PodProblem, nodes []collector.NodeProblem, events []collector.EventSummary, batchSize int) []batchProblem {
var batches []batchProblem
for i := 0; i < len(pods); i += batchSize {
end := i + batchSize
if end > len(pods) {
end = len(pods)
}
batch := batchProblem{
Pods: pods[i:end],
}
// Node 和 Event 每批都带上,因为它们是全局上下文
if len(nodes) > 0 {
batch.Nodes = nodes
}
if len(events) > 0 {
batch.Events = events
}
batches = append(batches, batch)
}
// 如果没有异常 Pod 但有异常 Node,单独成一批
if len(pods) == 0 && (len(nodes) > 0 || len(events) > 0) {
batches = append(batches, batchProblem{
Nodes: nodes,
Events: events,
})
}
return batches
}
运行方式:
# dry-run 模式先跑一遍看输出
go run cmd/aiops-agent/main.go \
--kubeconfig ~/.kube/config \
--llm-api-key sk-xxx \
--dry-run
# 确认无误后去掉 dry-run
go run cmd/aiops-agent/main.go \
--kubeconfig ~/.kube/config \
--llm-api-key sk-xxx
审计日志可以用 jq 查询:
# 查看所有执行过的操作
cat /var/log/aiops-audit.jsonl | jq 'select(.result == "success")'
# 统计每种操作的执行次数
cat /var/log/aiops-audit.jsonl | jq -r '.action' | sort | uniq -c
# 查看失败的执行
cat /var/log/aiops-audit.jsonl | jq 'select(.result == "failed")'
踩坑提示:
- 我强烈建议先 dry-run 跑一周,把审计日志和 LLM 决策都看一遍,确认没有离谱的操作再开真实执行。我第一次直接开真实模式,LLM 把一个正常 Pod 给 restart 了,因为那个 Pod 重启了 2 次(Java 启动慢,健康检查超时),LLM 觉得它"异常"。加了一周 dry-run 后这种误判就容易发现了。
low confidence跳过执行这个策略救过我好几次。LLM 在信息不全的时候会给出低置信度的建议,这种建议很多时候是错的。宁可让人来看也不要让低置信度的决策自动执行。- 巡检间隔 60 秒是经验值。太短(<30s)API 压力大,太长(>120s)反应慢。60 秒刚好能覆盖大部分 CrashLoopBackOff 的 back-off 周期(默认 5 分钟指数退避)。
进阶:资源画像与自动扩缩容
除了异常修复,client-go 还可以用于采集资源画像,为 AIOps 模型提供训练数据:
// pkg/collector/metrics.go
package collector
import (
"context"
"fmt"
"time"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/rest"
metricsv1beta1 "k8s.io/metrics/pkg/client/clientset/versioned"
)
// ResourceSnapshot 资源使用快照
type ResourceSnapshot struct {
Timestamp time.Time `json:"timestamp"`
Pods []PodMetric `json:"pods"`
}
// PodMetric 单个 Pod 的资源使用
type PodMetric struct {
Namespace string `json:"namespace"`
PodName string `json:"podName"`
NodeName string `json:"nodeName"`
CPUm int64 `json:"cpuMilliCores"` // millicores
MemMB int64 `json:"memMB"` // MB
}
// CollectMetrics 从 metrics-server 采集 Pod 级别的资源使用数据
// metrics-server 是 K8s 原生指标采集组件,HPA 就靠它
// restConfig 通常是创建 clientset 时用的那个 *rest.Config,直接传进来复用
func CollectMetrics(ctx context.Context, restConfig *rest.Config) (*ResourceSnapshot, error) {
// 创建 metrics client
// 注意:metrics-server 需要单独部署,不是 K8s 自带的
metricsClient, err := metricsv1beta1.NewForConfig(restConfig)
if err != nil {
return nil, fmt.Errorf("create metrics client failed: %w", err)
}
// 拉取所有 Pod 的指标
podMetrics, err := metricsClient.MetricsV1beta1().PodMetricses("").List(ctx, metav1.ListOptions{})
if err != nil {
return nil, fmt.Errorf("list pod metrics failed: %w", err)
}
snapshot := &ResourceSnapshot{
Timestamp: time.Now(),
}
for _, pm := range podMetrics.Items {
// 一个 Pod 可能有多个容器,累加
var totalCPU, totalMem int64
for _, container := range pm.Containers {
totalCPU += container.Usage.Cpu().MilliValue()
totalMem += container.Usage.Memory().Value() / (1024 * 1024) // 转成 MB
}
snapshot.Pods = append(snapshot.Pods, PodMetric{
Namespace: pm.Namespace,
PodName: pm.Name,
// 注意:PodMetrics 结构体不包含 NodeName 字段
// 如需 NodeName,需要用 clientset 查 Pod 对象获取
NodeName: "",
CPUm: totalCPU,
MemMB: totalMem,
})
}
return snapshot, nil
}
| 指标项 | 采集方式 |
|---|---|
| Pod CPU/内存使用 | metrics-server API(上面代码) |
| 请求 QPS | Ingress Controller 指标(Prometheus 查询) |
| 启动时间 | Pod 对象的 status.startTime |
| 节点负载 | Node 对象的 status.capacity 与 status.allocatable |
| 网络收发流量 | /proc/net/dev 或 CNI 插件指标 |
把这些数据持久化到时序数据库(InfluxDB/TDengine/Prometheus),就能用于训练流量预测模型,实现预测性自动扩容。下一篇文章会专门讲这个。
安全与审计
自动化修复必须考虑安全。我把它分成三层防御:
┌─────────────────────────────────────────────────┐
│ 第一层:Prompt 约束 │
│ 在 system prompt 里告诉 LLM 哪些不能做 │
│ ↓ 可被绕过(LLM 不一定遵守) │
├─────────────────────────────────────────────────┤
│ 第二层:代码白名单 │
│ executor 里的 isAllowedTool + 具体操作的安全检查 │
│ ↓ 不可绕过(代码级硬限) │
├─────────────────────────────────────────────────┤
│ 第三层:运行时保护 │
│ dry-run、审计日志、低置信度跳过、限流 │
│ ↓ 最后的安全网 │
└─────────────────────────────────────────────────┘
具体来说:
- 操作白名单:
isAllowedTool只允许restart_pod、scale_deployment、cordon_node、no_action。其他一律拒绝。 - 特殊资源保护:
restartPod里检查StatefulSetPod 直接拒绝。kube-systemnamespace 硬编码拒绝。 - 参数范围限制:
scaleDeployment副本数限制 1-20,不允许 0 或超大值。 - dry-run 模式:上线前先跑一周 dry-run,看 LLM 决策合不合理。
- 审计日志:每次操作都写 JSON 文件,一行一条,方便用
jq查。 - 低置信度跳过:LLM 返回
lowconfidence 的决策不自动执行,只记录。 - 限流:每轮巡检最多执行一个操作,每批之间间隔 5 秒,避免短时间内大量操作。
- 人工确认机制(可选):高风险操作通过飞书/钉钉通知,人工点击确认后才执行。
我个人的经验是:自动修复要从小范围开始。先只开 restart_pod 一个操作,跑两周没问题再开 scale_deployment。别一上来就全开,出事的概率比你想象的高。
总结
client-go 不仅是 K8s 的编程接口,更是 AIOps 系统与集群交互的"手脚"。结合 LLM 的推理能力,我们可以构建从观测、决策到执行的完整智能运维闭环。
这篇笔记里写的代码不是生产级——它缺了 Informer/Watch 机制、缺了多集群支持、缺了 Prometheus 指标采集、缺了优雅的配置管理。但核心流程是完整的:采集异常 -> LLM 分析 -> 安全执行 -> 审计追溯。
后续文章会在此基础上进一步学习如何开发 Kubernetes Operator,把这套能力以声明式方式沉淀到集群中。Operator 相比独立 Agent 的优势在于:它可以跟 K8s 的 reconcile 机制深度结合,不需要外部 cron 调度,也不需要单独部署和监控。
最后一句忠告:AIOps 自动修复落地时,安全控制比 LLM 能力更重要。你的 LLM 模型再强、prompt 写得再好,只要有一道安全检查没做好,一个误操作就可能让线上挂半小时。把安全控制做到"即使 LLM 疯了也不会出大事"的程度,再考虑上线。
自测题与动手练习
自测题(合上书能答出来,才算懂):
- 本文的 AIOps 闭环拆成了"观测 / 决策 / 执行"三段,请说出每一段 client-go 具体要干什么?为什么是"每 60 秒一轮"而不是实时 Watch?
- 异常 Pod 检测通常会关注哪几类状态(
CrashLoopBackOff/ImagePullBackOff/Evicted)?它们各自对应什么运维动作? - LLM 返回的是 Function Calling 输出,执行阶段在真正调 client-go 之前要做哪几道校验?
isAllowedTool白名单的意义是什么? - 本文对
StatefulSetPod 和kube-systemnamespace 做了特殊保护,为什么这两类资源不能随便自动重启/缩容? - dry-run 模式和审计日志分别解决了自动修复的什么问题?为什么"低置信度跳过"是必要的安全阀?
动手练习(建议真做一遍):
- 改造巡检程序,让它除了列异常 Pod 外,还输出每个 Node 的 CPU/内存使用率,并标出"接近 limit"的高风险 Pod。
- 在 dry-run 模式下跑一遍主循环,把 LLM 给出的所有修复建议落到一个 JSON 审计文件里,用
jq统计它想执行哪些操作、哪些被白名单拦了。 - 给
scaleDeployment加一条"缩容不低于当前副本 50%“的护栏,观察在流量抖动时是否能避免误缩容;再试着把操作限流从"每轮 1 个"放宽,看会不会引发抖动。
本章小结
- AIOps 用 client-go 把"发现问题 → 人去处理"的慢链路,压缩成"观测 → 决策 → 执行"的自动闭环,目标是让程序在低风险场景先处理一轮。
- 观测阶段靠 Clientset 拉 Pod/Node/Event;决策阶段把异常喂给 LLM 拿根因和修复建议;执行阶段解析 Function Calling 并用 client-go 落地。
- 安全护栏是自动修复的生命线:操作白名单、特殊资源保护、参数范围限制、dry-run、审计日志、低置信度跳过、操作限流,七道闸缺一不可。
- 这套代码不是生产级——缺 Informer/Watch、多集群、Prometheus 指标和配置管理,但核心三段式流程是完整的。
- 落地铁律:安全控制比 LLM 能力更重要,先小范围(只开
restart_pod)跑稳再逐步放开。
用 client-go 把异常自动捞出来只是第一步;异常捞到之后,如何结合流量预测做"预测式扩缩容"来从根上减少异常,正是下一篇要展开的工程链路。