学习目标
学完本章你应该能够:
- 用自己的话讲清"为什么单体 controller 不适合 K8s 故障自动修复",以及多 Agent 拆分的工程动机。
- 画出五大 Agent(观测 / 诊断 / 执行 / 审计 / 协调)的职责边界与协作关系。
- 设计一套跨 Agent 的统一消息信封(Message envelope),说清
trace_id与 UTC 时间戳的作用。 - 用状态机模型描述一次故障修复的完整生命周期,并讲清"迟到的事件"该如何处理。
- 在面试中把这个系统讲成一个"有护栏、可回滚、可审计"的工程故事,而不是一堆脚本。
前置知识:
- Go 基础:
interface、channel、context、sync.Mutex - Kubernetes 基础:Pod / Node / Event、client-go informer 机制
- LLM 基本认知:prompt 构造、JSON 输出、置信度(confidence)
本章你会动手做的事:
- 跑通观测 Agent 的 informer 采集逻辑,观察一个异常 Pod 如何被识别并带
trace_id发出。 - 拿一条真实 Pod 崩溃日志,手工构造一条 Observation 消息喂给诊断 Agent。
- 画出一次 “CrashLoopBackOff → 扩容” 的状态机流转图,并标出超时分支。
问题背景
Kubernetes 集群中的故障排查通常涉及多个维度:指标异常、Pod 状态异常、网络不通、存储挂载失败、配置错误等。单一脚本或单一 Agent 很难覆盖所有场景。多 Agent 协同架构通过将不同能力拆分为独立角色,可以更高效、更可靠地完成复杂故障的自动诊断与修复。
说实话我一开始也不太信这套,觉得"写个 controller 不就完了"。直到有一次线上 Node NotReady 引发连锁告警,我们那个单体 controller 同时打 K8s API、调 LLM、写日志,结果 goroutine 飙到几千,自己 OOM 了。从那以后我才真正理解为啥要把不同职责拆开——不是为了好看,是为了炸的时候不连坐。
多 Agent 系统架构
一个完整的 K8s 故障自动修复系统可划分为以下 Agent:
+-------------+ +-------------+ +-------------+ +-------------+
| 观测 Agent | -> | 诊断 Agent | -> | 执行 Agent | -> | 审计 Agent |
+-------------+ +-------------+ +-------------+ +-------------+
^ |
+---------------------------------------------------+
类比:把这套系统想成一家"故障急诊医院"。观测 Agent 是分诊台护士,把病人症状(指标 / 日志 / 事件)量出来;诊断 Agent 是主治医生,结合病历(RAG)和自己的经验(LLM)下诊断;执行 Agent 是手术团队,动手治疗;审计 Agent 是病历管理员,谁动了什么全记下来;协调 Agent 是急诊科主任,掌控整个流程不跑偏。
下面这张图是五大 Agent 的协作关系——注意最后有一条虚线绕回观测:修复之后还要回头确认指标是否真的恢复,形成闭环。
flowchart LR
OBS[观测 Agent
采集指标/日志/事件] --> DIAG[诊断 Agent
LLM+RAG 推理根因]
DIAG --> EXEC[执行 Agent
调用 K8s API 修复]
EXEC --> AUD[审计 Agent
记录+回滚]
AUD --> COORD[协调 Agent
状态机串联全流程]
COORD -. 分配任务/管理状态 .-> OBS| Agent 角色 | 职责 |
|---|---|
| 观测 Agent | 采集指标、日志、事件、链路数据 |
| 诊断 Agent | 基于 LLM 和知识库推理故障根因 |
| 执行 Agent | 调用 K8s API 或运维工具执行修复 |
| 审计 Agent | 记录操作日志,支持回滚与复盘 |
| 协调 Agent | 分配任务、管理状态、协调多 Agent 协作 |
我个人觉得这五个角色里最容易被低估的是审计 Agent,很多人觉得"不就是写个日志嘛"。但你只要踩过一次"自动修复把集群搞炸了却不知道谁干的"的坑,就会明白审计不是日志,是责任边界。
Agent 间通信协议
在写具体 Agent 之前,我先把消息格式定下来。这点很关键——Agent 之间不约定好协议,后面联调就是灾难现场。我踩过,A Agent 用 snake_case,B Agent 用 camelCase,光是写适配代码就花了我一个下午。
我们采用一个统一的 Message envelope,所有 Agent 收发都用这个结构:
⚠️ 新手必踩的坑:协议不统一。A Agent 用 snake_case、B Agent 用 camelCase,联调时写适配代码能写一下午。所以在写任何业务代码之前,先把消息格式钉死,后面省的是命。
下面用一张时序图示意"观测 → 诊断"这一跳的消息长什么样——注意 trace_id 一路带着走:
sequenceDiagram
participant OBS as 观测 Agent
participant DIAG as 诊断 Agent
OBS->>DIAG: observation 消息(trace_id=inc-001, severity=critical)
Note over DIAG: 解析 Pod 崩溃现场字段
DIAG-->>AUD: diagnosis 消息(根因 + 建议动作)// agent/message.go
package agent
import "time"
// MessageKind 区分消息类型,避免用字符串硬编码到处飞
type MessageKind string
const (
KindObservation MessageKind = "observation" // 观测结果
KindDiagnosis MessageKind = "diagnosis" // 诊断结果
KindAction MessageKind = "action" // 执行动作
KindAudit MessageKind = "audit" // 审计记录
KindError MessageKind = "error" // 错误回传
)
// Message 是所有 Agent 通信的统一信封
type Message struct {
TraceID string `json:"trace_id"` // 贯穿一次故障的 ID,必填
From string `json:"from"` // 发送方 Agent 名
To string `json:"to"` // 接收方 Agent 名
Kind MessageKind `json:"kind"` // 消息类型
Timestamp time.Time `json:"timestamp"` // 发送时间
Severity string `json:"severity"` // info/warn/critical
Payload map[string]interface{} `json:"payload"` // 实际业务数据
}
// NewMessage 构造消息,强制带 traceID 和时间戳
func NewMessage(traceID, from, to string, kind MessageKind) Message {
return Message{
TraceID: traceID,
From: from,
To: to,
Kind: kind,
Timestamp: time.Now().UTC(),
Severity: "info",
Payload: map[string]interface{}{},
}
}
一条具体的 Observation 消息长这样,方便理解:
{
"trace_id": "inc-20251110-001",
"from": "observer",
"to": "diagnoser",
"kind": "observation",
"timestamp": "2025-11-10T06:11:02Z",
"severity": "critical",
"payload": {
"alert": "PodCrashLoop",
"namespace": "payment",
"pod": "api-server-7d8b-x2k4",
"container": "api",
"restart_count": 17,
"last_log": "dial tcp redis:6379: connect: connection refused",
"events": ["Back-off restarting failed container"]
}
}
踩坑提示:trace_id 一定要从告警入口就生成好,一路传下去。我见过有人在每个 Agent 里都 uuid.New(),结果审计日志根本串不起来,出了事只能猜。还有,Timestamp 强制 UTC,否则跨时区集群对账能让你怀疑人生。
观测 Agent
观测 Agent 负责从多个数据源收集信息:
- Prometheus:CPU、内存、QPS、P99 延迟等指标。
- Loki / ELK:容器日志、应用日志、系统日志。
- Jaeger / Tempo:分布式链路追踪。
- Kubernetes Events:Pod 调度、镜像拉取、健康检查等事件。
为啥观测要单独拆出来?因为数据源太多,每个都有自己的客户端、认证、限流。如果让诊断 Agent 直接去拉,那诊断 Agent 就成了一个"什么都懂"的怪物。拆开后观测 Agent 只对"数据是否齐全"负责,诊断 Agent 只对"推理是否合理"负责。
下面是观测 Agent 用 client-go 采集异常 Pod、Node Condition、Event 的完整代码。我只列核心部分,依赖注入和 main 函数你自己补:
// agent/observer/pod_collector.go
package observer
import (
"context"
"fmt"
"time"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/tools/cache"
)
// PodAnomaly 描述一个异常 Pod 的快照
type PodAnomaly 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:"restart_count"`
NodeName string `json:"node_name"`
DetectedAt time.Time `json:"detected_at"`
ContainerMsgs []string `json:"container_msgs"` // 各 container 的 last state message
}
// PodCollector 用 informer 监听 Pod 变化,比轮询省 API
type PodCollector struct {
client kubernetes.Interface
out chan<- PodAnomaly
threshold int32 // 重启次数阈值,超过才上报
}
func NewPodCollector(client kubernetes.Interface, out chan<- PodAnomaly, restartThreshold int32) *PodCollector {
return &PodCollector{client: client, out: out, threshold: restartThreshold}
}
// Run 启动 informer,阻塞直到 ctx 取消
func (c *PodCollector) Run(ctx context.Context) error {
// 默认 30s resync,太短会刷屏,太长错过更新
factory := informers.NewSharedInformerFactoryWithOptions(
c.client,
30*time.Second,
informers.WithNamespace(""), // 全集群监听
)
inf := factory.Core().V1().Pods().Informer()
// 这里有个坑:UpdateFunc 会被频繁触发,比如 Pod 的 status.hostIP 变化
// 所以必须自己判断"是否真的异常"以及"是否已经上报过",否则会刷爆下游
inf.AddEventHandler(cache.ResourceEventHandlerFuncs{
UpdateFunc: func(oldObj, newObj interface{}) {
oldPod, ok1 := oldObj.(*corev1.Pod)
newPod, ok2 := newObj.(*corev1.Pod)
if !ok1 || !ok2 {
return
}
// 资源版本没变就别处理,informer 偶尔会重复回调
if oldPod.ResourceVersion == newPod.ResourceVersion {
return
}
anomaly, ok := c.isAnomaly(newPod)
if !ok {
return
}
select {
case c.out <- *anomaly:
case <-ctx.Done():
// ctx 取消时直接返回,别阻塞在 channel
}
},
})
factory.Start(ctx.Done())
if !cache.WaitForCacheSync(ctx.Done(), inf.HasSynced) {
return fmt.Errorf("pod informer sync timeout")
}
<-ctx.Done()
return nil
}
// isAnomaly 判断 Pod 是否处于异常状态
func (c *PodCollector) isAnomaly(pod *corev1.Pod) (*PodAnomaly, bool) {
// Pending 太久也算异常,但这里只看 Running 阶段的容器异常
if pod.Status.Phase == corev1.PodSucceeded {
return nil, false // Job 正常完成,忽略
}
var reason string
var msgs []string
var maxRestart int32
for _, cs := range pod.Status.ContainerStatuses {
if cs.RestartCount > maxRestart {
maxRestart = cs.RestartCount
}
// 关键:waiting 状态才有 reason,terminated 状态看 reason+message
if cs.State.Waiting != nil && cs.State.Waiting.Reason != "" {
reason = cs.State.Waiting.Reason
if cs.State.Waiting.Message != "" {
msgs = append(msgs, cs.State.Waiting.Message)
}
}
if cs.State.Terminated != nil && cs.State.Terminated.Reason != "" {
reason = cs.State.Terminated.Reason
if cs.State.Terminated.Message != "" {
msgs = append(msgs, cs.State.Terminated.Message)
}
}
}
// 阈值过滤:重启次数不够或者没有明确 reason 就别报
// 我踩过坑:曾经没加 reason 判断,把正常滚动更新的 Pod 全报上去了
if reason == "" && maxRestart < c.threshold {
return nil, false
}
// 拿 Pod 自身的 condition reason 兜底,CrashLoopBackOff 一般在 container status 里
if reason == "" && pod.Status.Reason != "" {
reason = pod.Status.Reason
}
return &PodAnomaly{
Namespace: pod.Namespace,
Name: pod.Name,
Phase: string(pod.Status.Phase),
Reason: reason,
RestartCount: maxRestart,
NodeName: pod.Spec.NodeName,
DetectedAt: time.Now().UTC(),
ContainerMsgs: msgs,
}, true
}
Node Condition 采集相对简单,因为 Node 数量少、变化慢,直接 List-Watch 就行:
// agent/observer/node_collector.go
package observer
import (
"context"
"time"
corev1 "k8s.io/api/core/v1"
"k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/tools/cache"
)
// NodeAnomaly 描述一个异常 Node
type NodeAnomaly struct {
Name string `json:"name"`
Conditions []string `json:"conditions"` // 异常 condition 列表
DetectedAt time.Time `json:"detected_at"`
}
type NodeCollector struct {
client kubernetes.Interface
out chan<- NodeAnomaly
}
func (c *NodeCollector) Run(ctx context.Context) error {
factory := informers.NewSharedInformerFactory(c.client, 60*time.Second)
inf := factory.Core().V1().Nodes().Informer()
inf.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) { c.check(ctx, obj) },
UpdateFunc: func(_, newObj interface{}) { c.check(ctx, newObj) },
})
factory.Start(ctx.Done())
if !cache.WaitForCacheSync(ctx.Done(), inf.HasSynced) {
return context.DeadlineExceeded
}
<-ctx.Done()
return nil
}
func (c *NodeCollector) check(ctx context.Context, obj interface{}) {
node, ok := obj.(*corev1.Node)
if !ok {
return
}
var badConds []string
for _, cond := range node.Status.Conditions {
// 关键判断:NetworkUnavailable / MemoryPressure / DiskPressure / PIDPressure
// 只要不是 True 的,全算异常;OutOfDisk 老版本才有,现在基本忽略
// 这里有个反例:Ready=False 不一定意味着 Node 真挂了,可能是 kubelet 暂时没上报
// 所以我们要再加一个 lastHeartbeatTime 太久的判断
if cond.Status != corev1.ConditionTrue && cond.Type != corev1.NodeReady {
badConds = append(badConds, string(cond.Type))
continue
}
if cond.Type == corev1.NodeReady && cond.Status != corev1.ConditionTrue {
// Ready=False 才是真正的 NotReady,且持续超过 1 分钟才算
if time.Since(cond.LastHeartbeatTime.Time) > time.Minute {
badConds = append(badConds, "Ready=False(stale>1m)")
}
}
}
if len(badConds) == 0 {
return
}
select {
case c.out <- NodeAnomaly{
Name: node.Name,
Conditions: badConds,
DetectedAt: time.Now().UTC(),
}:
case <-ctx.Done():
}
}
Event 采集单独说一下。Event 在 K8s 里是会被自动清理的(默认 1 小时),所以你不能依赖 List 拿历史 Event,必须用 Watch 实时收。还有一个坑:Event 的 involvedObject 才是关联 Pod 的关键,过滤的时候要按这个来。
// agent/observer/event_collector.go
package observer
import (
"context"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/fields"
"k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/tools/cache"
)
// EventCollector 只关心 Warning 事件,Normal 事件太多没必要全收
type EventCollector struct {
client kubernetes.Interface
out chan<- corev1.Event
}
func (c *EventCollector) Run(ctx context.Context) error {
// 只看 Warning 事件,Normal 太多太吵
// 不过这里有个反例:有些健康检查失败的 Event 类型是 Normal 但 reason 是 Unhealthy
// 所以严格点应该再按 reason 过滤,但这会让代码很丑,折中:先按 type 过滤
fieldSelector := fields.OneTermEqualSelector("type", "Warning")
factory := informers.NewSharedInformerFactoryWithOptions(
c.client,
0,
informers.WithTweakListOptions(func(opts *metav1.ListOptions) {
opts.FieldSelector = fieldSelector.String()
}),
)
inf := factory.Core().V1().Events().Informer()
inf.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
ev, ok := obj.(*corev1.Event)
if !ok {
return
}
// 重复事件有 count 字段,count > 1 说明已经发生过多次
// 审计的时候要记 count,否则会以为是单次
select {
case c.out <- *ev:
case <-ctx.Done():
}
},
})
factory.Start(ctx.Done())
if !cache.WaitForCacheSync(ctx.Done(), inf.HasSynced) {
return context.DeadlineExceeded
}
<-ctx.Done()
return nil
}
踩坑提示:informer 的 WithTweakListOptions 配 FieldSelector 在大集群(10w+ Event)下能救命,不然你的 observer 内存会涨到几个 G。还有就是 Event 的 lastTimestamp 才是事件真实发生时间,eventTime 字段有些版本不写,别用错。
诊断 Agent
诊断 Agent 接收观测数据,结合 LLM 与 RAG 知识库进行根因分析。诊断过程可采用 ReAct 模式(思考 - 行动 - 观察循环):
flowchart TD
T[思考: 提出根因假设] --> A[行动: 拉取日志/Redis 状态]
A --> O[观察: 拿到新证据]
O --> T
O --> C{证据是否充分?}
C -->|否| T
C -->|是| R[结论: 根因 + 修复建议]思考1:Pod 处于 CrashLoopBackOff,首先检查最近日志。
行动1:请求日志摘要。
观察1:日志显示 "connection refused to redis:6379"。
思考2:可能是 Redis 服务不可用或配置错误,需要检查 Redis Service 和 Endpoint。
行动2:请求 Redis 状态。
观察2:Redis Endpoint 为空,说明没有 Redis Pod 被选中。
结论:Redis Deployment 副本数为 0,需要扩容。
诊断 Agent 可以维护一个故障模式库,将历史故障案例向量化存储,遇到新告警时先进行相似案例检索。
LLM 根因分析最关键的是 prompt 构造。说实话我刚写的时候 prompt 就一行"分析这个故障",结果模型给我编出各种不存在的 K8s 概念。后来才明白:prompt 必须把上下文塞满,包括 Pod 的 status、events、最近日志、相关资源状态。模型不是算命先生,你给它什么它就用什么。
// agent/diagnoser/llm_diagnoser.go
package diagnoser
import (
"context"
"encoding/json"
"fmt"
"strings"
)
// DiagnosisInput 诊断 Agent 的输入,从 Observation 消息组装
type DiagnosisInput struct {
TraceID string `json:"trace_id"`
Namespace string `json:"namespace"`
Pod string `json:"pod"`
Reason string `json:"reason"` // CrashLoopBackOff 等
ContainerMsgs []string `json:"container_msgs"`
Events []string `json:"events"` // 相关 Event 摘要
RecentLogs string `json:"recent_logs"` // 最近 N 行日志
SimilarCases []string `json:"similar_cases"` // RAG 检索到的相似案例
}
// DiagnosisResult LLM 给出的诊断结论
type DiagnosisResult struct {
TraceID string `json:"trace_id"`
RootCause string `json:"root_cause"` // 根因描述
Confidence float64 `json:"confidence"` // 0~1
RecommendedAction string `json:"recommended_action"` // scale/restart/rollback/...
ActionTarget string `json:"action_target"` // 操作对象,比如 deployment/redis
Reasoning string `json:"reasoning"` // 推理过程,便于审计
}
// LLMClient 抽象一下,方便换 OpenAI/Claude/本地模型
type LLMClient interface {
Chat(ctx context.Context, prompt string) (string, error)
}
type LLMDiagnoser struct {
client LLMClient
}
func NewLLMDiagnoser(client LLMClient) *LLMDiagnoser {
return &LLMDiagnoser{client: client}
}
// buildPrompt 构造 prompt,这是诊断质量的关键
// 我个人觉得 prompt 工程在这里比换大模型还有用
func (d *LLMDiagnoser) buildPrompt(in DiagnosisInput) string {
var sb strings.Builder
// 系统角色:明确告诉模型你是干啥的,输出格式是什么
// 不约束输出格式,下游解析就是地狱
sb.WriteString("你是一个 Kubernetes 故障诊断专家。")
sb.WriteString("根据下面的观测数据,分析故障根因并给出修复建议。\n")
sb.WriteString("必须以 JSON 格式输出,字段固定为 root_cause/confidence/recommended_action/action_target/reasoning。\n")
sb.WriteString("confidence 为 0~1 的浮点数,低于 0.7 的不要建议执行破坏性操作。\n")
sb.WriteString("recommended_action 只能从 scale/restart/rollback/update_config/none 中选一个。\n\n")
// 现场数据:信息越全,幻觉越少
sb.WriteString("## 故障现场\n")
sb.WriteString(fmt.Sprintf("- namespace: %s\n", in.Namespace))
sb.WriteString(fmt.Sprintf("- pod: %s\n", in.Pod))
sb.WriteString(fmt.Sprintf("- reason: %s\n", in.Reason))
if len(in.ContainerMsgs) > 0 {
sb.WriteString(fmt.Sprintf("- container_messages: %s\n", strings.Join(in.ContainerMsgs, "; ")))
}
if len(in.Events) > 0 {
sb.WriteString("- recent_events:\n")
for _, e := range in.Events {
sb.WriteString(fmt.Sprintf(" * %s\n", e))
}
}
if in.RecentLogs != "" {
// 日志只截最近 50 行,太长会让 prompt 超 token 还容易跑偏
sb.WriteString("- recent_logs (tail 50):\n")
sb.WriteString(in.RecentLogs)
sb.WriteString("\n")
}
// RAG 检索的相似案例,给模型一个"参考答案"
// 但要明确告诉模型只是参考,避免它无脑照抄
if len(in.SimilarCases) > 0 {
sb.WriteString("\n## 历史相似案例(仅供参考,不要直接照搬结论)\n")
for i, c := range in.SimilarCases {
sb.WriteString(fmt.Sprintf("案例 %d: %s\n", i+1, c))
}
}
sb.WriteString("\n请给出诊断结论,只输出 JSON,不要任何额外解释。\n")
return sb.String()
}
// Diagnose 调用 LLM 做诊断,并解析结果
func (d *LLMDiagnoser) Diagnose(ctx context.Context, in DiagnosisInput) (*DiagnosisResult, error) {
prompt := d.buildPrompt(in)
raw, err := d.client.Chat(ctx, prompt)
if err != nil {
return nil, fmt.Errorf("llm chat failed: %w", err)
}
// 大模型经常会在 JSON 前后塞一段 markdown,比如 ```json ... ```
// 必须先清掉,否则 json.Unmarshal 直接炸
raw = strings.TrimSpace(raw)
raw = strings.TrimPrefix(raw, "```json")
raw = strings.TrimPrefix(raw, "```")
raw = strings.TrimSuffix(raw, "```")
raw = strings.TrimSpace(raw)
var result DiagnosisResult
if err := json.Unmarshal([]byte(raw), &result); err != nil {
// 解析失败别硬撑,把原始返回原样回传给审计 Agent 让人看
return nil, fmt.Errorf("parse llm output failed: %w, raw=%s", err, raw)
}
result.TraceID = in.TraceID
// 兜底:confidence 异常时强行降级
if result.Confidence < 0 || result.Confidence > 1 {
result.Confidence = 0.3
}
return &result, nil
}
踩坑提示:LLM 输出 JSON 的时候,强烈建议用 OpenAI 的 response_format: json_object 或者函数调用,否则你 10 次里有 1 次会拿到带前缀说明的"垃圾 JSON"。还有 confidence 是模型自报的,不能完全信,我们在执行 Agent 那一层还会再校验一次,破坏性操作低于 0.8 直接拒绝。
执行 Agent
执行 Agent 根据诊断结果执行修复动作。执行前需要校验:
- 动作是否在白名单中。
- 是否超过当日执行次数限制。
- 是否需要人工审批。
⚠️ 新手必踩的坑:破坏性操作零护栏。drain 一个 Node 等于把上面所有 Pod 驱逐走,如果忘了加白名单或人工审批,LLM 一拍脑袋就能把线上节点干掉。护栏(白名单 + 置信度 + 审批)不是可选项,是生死线。
执行 Agent 在真正动手前,会依次过四道闸门,任何一道不通过就直接拒绝:
flowchart TD
S[收到 Action] --> C1{在白名单?}
C1 -->|否| X[拒绝: 未知动作]
C1 -->|是| C2{置信度达标?}
C2 -->|否| X
C2 -->|是| C3{当日次数超限?}
C3 -->|是| X
C3 -->|否| C4{需人工审批?}
C4 -->|是| P[转 PendingApproval 等回执]
C4 -->|否| E[执行修复]常见修复动作包括:
| 故障场景 | 修复动作 |
|---|---|
| Pod 崩溃 | 重启 Pod、回滚镜像版本 |
| CPU 饱和 | 扩容 Deployment |
| 内存泄漏 | 滚动重启、调整资源限制 |
| 配置错误 | 更新 ConfigMap/Secret 并滚动更新 |
| 网络不通 | 检查 Service/Endpoint/NetworkPolicy |
// agent/executor/k8s_executor.go
package executor
import (
"context"
"fmt"
"time"
corev1 "k8s.io/api/core/v1"
policyv1 "k8s.io/api/policy/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/labels"
"k8s.io/client-go/kubernetes"
)
// Action 执行动作的统一结构
type Action struct {
TraceID string `json:"trace_id"`
Type string `json:"type"` // scale/restart/rollback/drain
Namespace string `json:"namespace"`
Target string `json:"target"` // deployment name / node name
Replicas int32 `json:"replicas"` // scale 用
Confidence float64 `json:"confidence"` // 来自诊断 Agent
Reason string `json:"reason"` // 为什么执行
}
// Executor 执行 Agent 主体
type Executor struct {
client kubernetes.Interface
dailyLimit int // 每日每动作上限
executedToday map[string]int // 动作类型 -> 当日次数
whitelist map[string]bool // 允许的动作类型
requireApproval map[string]bool // 需要人工审批的动作
}
func NewExecutor(client kubernetes.Interface, dailyLimit int) *Executor {
return &Executor{
client: client,
dailyLimit: dailyLimit,
executedToday: make(map[string]int),
// 白名单是个护栏,没列进去的动作一律拒绝
// 我踩过坑:曾经忘了把 drain 加白名单,结果 Node 故障时执行 Agent 直接报错
whitelist: map[string]bool{
"scale": true, "restart": true, "rollback": true, "drain": true,
},
// 破坏性动作必须人工审批,别让 LLM 一拍脑袋就把 Node drain 了
requireApproval: map[string]bool{
"drain": true, "rollback": true,
},
}
}
// Execute 执行动作,所有校验失败都返回 error,绝不悄悄吞掉
func (e *Executor) Execute(ctx context.Context, action Action) error {
// 校验 1:白名单
if !e.whitelist[action.Type] {
return fmt.Errorf("action %s not in whitelist", action.Type)
}
// 校验 2:置信度门槛,破坏性动作要求更高
minConfidence := 0.7
if e.requireApproval[action.Type] {
minConfidence = 0.85
}
if action.Confidence < minConfidence {
return fmt.Errorf("confidence %.2f below threshold %.2f for action %s",
action.Confidence, minConfidence, action.Type)
}
// 校验 3:当日执行次数
if e.executedToday[action.Type] >= e.dailyLimit {
return fmt.Errorf("daily limit %d reached for action %s", e.dailyLimit, action.Type)
}
// 校验 4:人工审批(这里只做模拟,实际接钉钉/飞书 webhook)
if e.requireApproval[action.Type] {
// 真实场景下这里应该 block 等审批结果,或者把 Action 入队等审批回执
// 这里简化为直接拒绝,由协调 Agent 转入 PendingApproval 状态
return fmt.Errorf("action %s requires human approval", action.Type)
}
// 所有校验通过,执行
var err error
switch action.Type {
case "scale":
err = e.scale(ctx, action)
case "restart":
err = e.restart(ctx, action)
case "drain":
err = e.drainNode(ctx, action)
default:
err = fmt.Errorf("unsupported action: %s", action.Type)
}
if err == nil {
e.executedToday[action.Type]++
}
return err
}
// scale 扩缩容 Deployment
func (e *Executor) scale(ctx context.Context, action Action) error {
dep, err := e.client.AppsV1().Deployments(action.Namespace).Get(ctx, action.Target, metav1.GetOptions{})
if err != nil {
return fmt.Errorf("get deployment failed: %w", err)
}
// 关键:直接改 spec.replicas,不要用 patch,因为我们要先读后写避免覆盖别人改动
// 这里有竞态——并发两个 scale 会后写覆盖前写,必要时用 resourceVersion 做乐观锁
replicas := action.Replicas
if replicas == *dep.Spec.Replicas {
return nil // 已经是目标副本数,幂等返回
}
dep.Spec.Replicas = &replicas
_, err = e.client.AppsV1().Deployments(action.Namespace).Update(ctx, dep, metav1.UpdateOptions{})
return err
}
// restart 删除 Pod 触发重建,比 rollout restart 更直接
func (e *Executor) restart(ctx context.Context, action Action) error {
// action.Target 是 deployment 名,先按 label 找出所有 Pod
dep, err := e.client.AppsV1().Deployments(action.Namespace).Get(ctx, action.Target, metav1.GetOptions{})
if err != nil {
return err
}
selector := labels.Set(dep.Spec.Selector.MatchLabels).AsSelector()
podList, err := e.client.CoreV1().Pods(action.Namespace).List(ctx, metav1.ListOptions{
LabelSelector: selector.String(),
})
if err != nil {
return err
}
// 关键:删除策略是 Foreground,等 Pod 真的没了再返回
// 不然你前脚删完,后脚 controller 还没创建新 Pod,业务就断一会儿
// 还有,别一次删光所有 Pod,至少留一个健康的
if len(podList.Items) <= 1 {
return fmt.Errorf("only %d pods, refuse to restart all", len(podList.Items))
}
deleted := 0
for _, pod := range podList.Items {
// 跳过健康的 Pod,只删异常的
if pod.Status.Phase == corev1.PodRunning && isPodReady(&pod) {
continue
}
// Foreground 策略:先删 Pod,等 finalizer 清完
// 这里的坑:Foreground 在某些 K8s 版本上对 Pod 不生效,实际行为是立即删除
// 所以更稳的做法是用 metav1.DeleteOptions 里的 GracePeriodSeconds
grace := int64(30) // 给 30s 优雅退出
err := e.client.CoreV1().Pods(action.Namespace).Delete(ctx, pod.Name, metav1.DeleteOptions{
GracePeriodSeconds: &grace,
})
if err != nil {
return fmt.Errorf("delete pod %s failed: %w", pod.Name, err)
}
deleted++
if deleted >= len(podList.Items)/2 {
break // 一次最多删一半,避免雪崩
}
}
return nil
}
// drainNode 驱逐节点上所有 Pod,破坏性最大,必须谨慎
func (e *Executor) drainNode(ctx context.Context, action Action) error {
// 1. 先把节点标记为 unschedulable
_, err := e.client.CoreV1().Nodes().Update(ctx,
&corev1.Node{
ObjectMeta: metav1.ObjectMeta{Name: action.Target},
Spec: corev1.NodeSpec{Unschedulable: true},
}, metav1.UpdateOptions{})
if err != nil {
return fmt.Errorf("cordon node failed: %w", err)
}
// 2. 拿到节点上所有 Pod(排除 daemonset)
podList, err := e.client.CoreV1().Pods("").List(ctx, metav1.ListOptions{
FieldSelector: fmt.Sprintf("spec.nodeName=%s", action.Target),
})
if err != nil {
return err
}
// 3. 逐个 evict
// 这里只做骨架,真实场景还要处理 PDB(PodDisruptionBudget)
// 忽略 PDB 直接驱逐可能导致服务可用副本数不足,是个常见坑
for _, pod := range podList.Items {
if isDaemonSetPod(&pod) {
continue // DaemonSet Pod 跟着节点走,驱逐没意义
}
// Eviction 是个 subresource,比 Delete 多一层 PDB 检查
evict := &policyv1.Eviction{
ObjectMeta: metav1.ObjectMeta{
Namespace: pod.Namespace,
Name: pod.Name,
},
DeleteOptions: metav1.DeleteOptions{
GracePeriodSeconds: ptrInt64(30),
},
}
err := e.client.PolicyV1().Evictions(pod.Namespace).Evict(ctx, evict)
if err != nil {
// 单个 Pod 驱逐失败不阻断整个 drain,记录下来继续
// 真实实现应该把失败 Pod 名字收集起来回传给审计
continue
}
// 给 controller 一点时间重新调度,不然下一秒就驱逐下一个可能造成压力
time.Sleep(2 * time.Second)
}
return nil
}
// isPodReady 简单判断 Pod 是否 Ready
func isPodReady(pod *corev1.Pod) bool {
for _, cond := range pod.Status.Conditions {
if cond.Type == corev1.PodReady && cond.Status == corev1.ConditionTrue {
return true
}
}
return false
}
// isDaemonSetPod 通过 ownerReferences 判断是否 DaemonSet
func isDaemonSetPod(pod *corev1.Pod) bool {
for _, ref := range pod.ObjectMeta.OwnerReferences {
if ref.Kind == "DaemonSet" {
return true
}
}
return false
}
func ptrInt64(v int64) *int64 { return &v }
踩坑提示:Eviction 和 Delete 是两码事,前者会检查 PDB,后者直接干掉。我见过有人图省事直接用 Delete,结果把一个只有 1 副本的服务干没了。还有 drain 一定要先 cordon 再 evict,不然你 evict 完 controller 又把新 Pod 调度到这个节点上,白干。
审计 Agent
审计 Agent 记录每次故障处理的完整上下文:
- 原始告警信息
- 各 Agent 的输入输出
- 最终执行的修复动作
- 执行前后的指标对比
- 人工介入记录
这些数据不仅用于合规审计,也是持续优化诊断模型和修复策略的重要素材。
审计最实用的做法是把决策写入 K8s Event,这样运维直接 kubectl describe pod 就能看到这个 Pod 被自动修复过。同时再落一份到外部存储(ES/ClickHouse)做长期归档。
// agent/auditor/k8s_auditor.go
package auditor
import (
"context"
"encoding/json"
"fmt"
"time"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
)
// AuditRecord 一次完整故障处理的审计记录
type AuditRecord struct {
TraceID string `json:"trace_id"`
StartedAt time.Time `json:"started_at"`
FinishedAt time.Time `json:"finished_at"`
Alert string `json:"alert"`
Namespace string `json:"namespace"`
Target string `json:"target"`
Diagnosis map[string]interface{} `json:"diagnosis"`
Action string `json:"action"`
Result string `json:"result"` // success/failed/skipped
ErrorMessage string `json:"error_message,omitempty"`
}
// Auditor 审计 Agent
type Auditor struct {
client kubernetes.Interface
eventSink EventSink // 外部存储,比如 ES/ClickHouse
}
// EventSink 外部审计存储抽象
type EventSink interface {
Write(ctx context.Context, record AuditRecord) error
}
func NewAuditor(client kubernetes.Interface, sink EventSink) *Auditor {
return &Auditor{client: client, eventSink: sink}
}
// WriteK8sEvent 把审计记录写成 K8s Event,绑到对应资源上
// 这样运维 kubectl describe 的时候直接能看到"这个 Pod 被自动修过"
func (a *Auditor) WriteK8sEvent(ctx context.Context, record AuditRecord) error {
var involvedObject corev1.ObjectReference
// 简化:这里假设 record.Target 是 Pod 名,实际要根据 Action 类型判断是 Pod/Deployment/Node
involvedObject = corev1.ObjectReference{
Kind: "Pod",
Namespace: record.Namespace,
Name: record.Target,
}
// Event 的 message 控制在 1KB 以内,太长会被截断
// 完整数据走 eventSink 落外部存储
summary := fmt.Sprintf("auto-remediation: alert=%s action=%s result=%s trace=%s",
record.Alert, record.Action, record.Result, record.TraceID)
// Event 名字带 traceID,方便关联
// 这里有个坑:Event 名字长度限制 253 字符,namespace/traceID 拼起来可能超
// 不过实际 traceID 一般 30 字符内,没问题
event := &corev1.Event{
ObjectMeta: metav1.ObjectMeta{
Name: fmt.Sprintf("remediation-%s", record.TraceID),
// 关键:不带 namespace 的话 cluster-scoped,绑不上 involvedObject
Namespace: record.Namespace,
},
InvolvedObject: involvedObject,
// 用自定义 reason 方便后续筛选
Reason: "AutoRemediation",
Message: summary,
// 关键:Type 是 Normal/Warning,别瞎填别的值,dashboard 会显示不出来
Type: eventTypeByResult(record.Result),
Source: corev1.EventSource{
Component: "aiops-remediation-agent",
},
FirstTimestamp: metav1.Time{Time: record.StartedAt},
LastTimestamp: metav1.Time{Time: record.FinishedAt},
Count: 1,
}
_, err := a.client.CoreV1().Events(record.Namespace).Create(ctx, event, metav1.CreateOptions{})
if err != nil {
return fmt.Errorf("create k8s event failed: %w", err)
}
return nil
}
// WriteExternal 同时落一份到外部存储,做长期归档和分析
func (a *Auditor) WriteExternal(ctx context.Context, record AuditRecord) error {
if a.eventSink == nil {
return nil
}
return a.eventSink.Write(ctx, record)
}
// WriteAll 一把梭,K8s Event 和外部存储都写
// K8s Event 写失败不影响外部存储,反之亦然,互不阻塞
func (a *Auditor) WriteAll(ctx context.Context, record AuditRecord) error {
// 先写外部存储,因为这个最关键,K8s Event 只是给运维看的
if err := a.WriteExternal(ctx, record); err != nil {
return fmt.Errorf("write external audit failed: %w", err)
}
// K8s Event 写失败只 warn,不影响主流程
if err := a.WriteK8sEvent(ctx, record); err != nil {
fmt.Printf("warn: write k8s event failed: %v\n", err)
}
return nil
}
// eventTypeByResult 根据执行结果返回 Event type
func eventTypeByResult(result string) string {
switch result {
case "success":
return "Normal"
case "failed", "skipped":
return "Warning"
default:
return "Normal"
}
}
// MarshalForLog 把审计记录序列化成 JSON,方便写日志或发到消息队列
func (r AuditRecord) MarshalForLog() ([]byte, error) {
return json.Marshal(r)
}
踩坑提示:K8s Event 默认会被 kube-controller-manager 在 1 小时后清理(--event-ttl),所以审计数据必须有外部存储兜底。我见过有人只写 K8s Event,结果一周后想复盘发现啥都没了。还有 Event 的 message 字段超长会被截断到 1KB 左右,别在里面塞完整诊断 JSON。
协调 Agent 与状态机
协调 Agent 负责管理一次故障修复的完整生命周期。可以使用状态机模型:
Detected -> Observing -> Diagnosing -> PendingApproval -> Executing -> Verifying -> Resolved/Failed
状态转换通过消息总线或工作流引擎驱动。各 Agent 之间通过定义良好的消息协议通信,避免紧耦合。
类比:状态机就像快递的物流轨迹——“已揽收 → 运输中 → 派送中 → 已签收 / 拒收”。每个包裹(一次故障)都有自己独立的轨迹,且只能按允许的路径推进,不能从"已签收"突然跳回"运输中"。
下面这张状态机图把所有合法流转画全了,包括任意"进行中"状态超时都可直接进 Failed:
stateDiagram-v2
[*] --> Observing
Observing --> Diagnosing: observe_done
Diagnosing --> Executing: diagnose_done
Diagnosing --> PendingApproval: need_approval
PendingApproval --> Executing: approved
PendingApproval --> Failed: rejected
Executing --> Verifying: execute_done
Executing --> Failed: execute_failed
Verifying --> Resolved: verify_passed
Verifying --> Failed: verify_failed
Observing --> Failed: timeout
Diagnosing --> Failed: timeout
Executing --> Failed: timeout
Verifying --> Failed: timeout
Resolved --> [*]
Failed --> [*]光看这张图其实挺好懂,但真正实现的时候状态机是最容易写烂的。我个人的经验是:状态机千万别用一堆 if state == "xxx" 串起来,三个月后你自己都看不懂。用一个明确的 transition 表,谁能转到谁、转移条件是什么,全列出来,代码即文档。
// agent/coordinator/state_machine.go
package coordinator
import (
"context"
"errors"
"fmt"
"sync"
"time"
)
// State 故障处理的状态
type State string
const (
StateDetected State = "detected"
StateObserving State = "observing"
StateDiagnosing State = "diagnosing"
StatePendingApproval State = "pending_approval"
StateExecuting State = "executing"
StateVerifying State = "verifying"
StateResolved State = "resolved"
StateFailed State = "failed"
)
// Event 触发状态转移的事件
type Event string
const (
EventObserveDone Event = "observe_done"
EventDiagnoseDone Event = "diagnose_done"
EventNeedApproval Event = "need_approval"
EventApproved Event = "approved"
EventRejected Event = "rejected"
EventExecuteDone Event = "execute_done"
EventExecuteFailed Event = "execute_failed"
EventVerifyPassed Event = "verify_passed"
EventVerifyFailed Event = "verify_failed"
EventTimeout Event = "timeout"
)
// transition 状态转移定义
type transition struct {
from State
event Event
to State
// action 在转移时执行,返回 error 则状态机进入 Failed
action func(ctx context.Context, inc *Incident) error
}
// Incident 一次故障处理的完整上下文
type Incident struct {
TraceID string `json:"trace_id"`
State State `json:"state"`
Alert string `json:"alert"`
Namespace string `json:"namespace"`
Target string `json:"target"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
// 各阶段的产物,避免 Agent 间重复传参
Observation map[string]interface{} `json:"observation,omitempty"`
Diagnosis map[string]interface{} `json:"diagnosis,omitempty"`
Action map[string]interface{} `json:"action,omitempty"`
Audit map[string]interface{} `json:"audit,omitempty"`
mu sync.Mutex
}
// Coordinator 协调 Agent 主体
type Coordinator struct {
transitions []transition
incidents map[string]*Incident // traceID -> Incident
mu sync.RWMutex
}
func NewCoordinator() *Coordinator {
c := &Coordinator{incidents: make(map[string]*Incident)}
c.transitions = []transition{
{StateDetected, EventObserveDone, StateDiagnosing, c.doDiagnose},
{StateObserving, EventObserveDone, StateDiagnosing, c.doDiagnose},
{StateDiagnosing, EventDiagnoseDone, StateExecuting, c.doExecute},
{StateDiagnosing, EventNeedApproval, StatePendingApproval, nil},
{StatePendingApproval, EventApproved, StateExecuting, c.doExecute},
{StatePendingApproval, EventRejected, StateFailed, c.doAuditOnReject},
{StateExecuting, EventExecuteDone, StateVerifying, c.doVerify},
{StateExecuting, EventExecuteFailed, StateFailed, c.doAuditOnFail},
{StateVerifying, EventVerifyPassed, StateResolved, c.doAuditOnResolve},
{StateVerifying, EventVerifyFailed, StateFailed, c.doAuditOnFail},
// 超时从任何"进行中"状态都可以直接跳 Failed
{StateObserving, EventTimeout, StateFailed, c.doAuditOnFail},
{StateDiagnosing, EventTimeout, StateFailed, c.doAuditOnFail},
{StateExecuting, EventTimeout, StateFailed, c.doAuditOnFail},
{StateVerifying, EventTimeout, StateFailed, c.doAuditOnFail},
}
return c
}
// HandleEvent 处理一个状态转移事件
// 这是协调 Agent 的核心入口,所有 Agent 完成工作后都通过这里通知协调者
func (c *Coordinator) HandleEvent(ctx context.Context, traceID string, event Event, payload map[string]interface{}) error {
c.mu.RLock()
inc, ok := c.incidents[traceID]
c.mu.RUnlock()
if !ok {
return fmt.Errorf("incident %s not found", traceID)
}
inc.mu.Lock()
defer inc.mu.Unlock()
// 把 payload 塞到对应字段,方便后续 action 用
c.mergePayload(inc, event, payload)
// 找到匹配的转移规则
var matched *transition
for i := range c.transitions {
t := &c.transitions[i]
if t.from == inc.State && t.event == event {
matched = t
break
}
}
if matched == nil {
// 没有匹配的转移,可能是个迟到的消息,记一下但不报错
// 我踩过坑:曾经这里直接返回 error,结果上游疯狂重试
return fmt.Errorf("no transition from %s on event %s", inc.State, event)
}
// 执行转移动作
if matched.action != nil {
if err := matched.action(ctx, inc); err != nil {
// action 失败直接进入 Failed
inc.State = StateFailed
inc.UpdatedAt = time.Now().UTC()
return fmt.Errorf("action failed: %w", err)
}
}
inc.State = matched.to
inc.UpdatedAt = time.Now().UTC()
return nil
}
// StartIncident 开启一次故障处理
func (c *Coordinator) StartIncident(traceID, alert, namespace, target string) *Incident {
inc := &Incident{
TraceID: traceID,
State: StateObserving, // 直接进入 observing
Alert: alert,
Namespace: namespace,
Target: target,
CreatedAt: time.Now().UTC(),
UpdatedAt: time.Now().UTC(),
}
c.mu.Lock()
c.incidents[traceID] = inc
c.mu.Unlock()
return inc
}
// doDiagnose 转移到诊断阶段,触发诊断 Agent
// 这里只是骨架,真实实现是发消息到消息队列让诊断 Agent 消费
func (c *Coordinator) doDiagnose(ctx context.Context, inc *Incident) error {
if inc.Observation == nil {
return errors.New("observation is empty, cannot diagnose")
}
// 实际场景:发送消息到 diagnoser queue
// 这里只做占位
return nil
}
// doExecute 触发执行 Agent
func (c *Coordinator) doExecute(ctx context.Context, inc *Incident) error {
if inc.Diagnosis == nil {
return errors.New("diagnosis is empty, cannot execute")
}
return nil
}
// doVerify 触发验证:观察修复后指标是否恢复
func (c *Coordinator) doVerify(ctx context.Context, inc *Incident) error {
// 验证阶段的关键是"等"——给系统一点时间收敛
// 但不能无限等,必须设超时,超时就触发 EventTimeout -> Failed
return nil
}
// doAuditOnResolve 成功后写审计
func (c *Coordinator) doAuditOnResolve(ctx context.Context, inc *Incident) error {
inc.Audit = map[string]interface{}{
"result": "success",
"finished_at": time.Now().UTC(),
}
return nil
}
// doAuditOnFail 失败后写审计
func (c *Coordinator) doAuditOnFail(ctx context.Context, inc *Incident) error {
inc.Audit = map[string]interface{}{
"result": "failed",
"finished_at": time.Now().UTC(),
}
return nil
}
// doAuditOnReject 人工拒绝后写审计
func (c *Coordinator) doAuditOnReject(ctx context.Context, inc *Incident) error {
inc.Audit = map[string]interface{}{
"result": "rejected",
"finished_at": time.Now().UTC(),
}
return nil
}
// mergePayload 把事件附带的数据塞到 Incident 对应字段
func (c *Coordinator) mergePayload(inc *Incident, event Event, payload map[string]interface{}) {
if payload == nil {
return
}
switch event {
case EventObserveDone:
inc.Observation = payload
case EventDiagnoseDone, EventNeedApproval:
inc.Diagnosis = payload
case EventExecuteDone, EventExecuteFailed:
inc.Action = payload
}
}
踩坑提示:状态机最坑的是"迟到的事件"。比如诊断 Agent 已经回了结果,但状态机因为超时进入了 Failed,这时候 EventDiagnoseDone 来了,你要决定是丢弃还是复活。我个人倾向于丢弃+审计记录,因为复活很容易把状态搞乱。还有一个常见问题:同一个 traceID 的并发事件,必须加锁(我用了 inc.mu),不然两个事件同时改 state 就热闹了。
关键技术挑战
- 幻觉与误判:LLM 可能给出错误诊断,需要设置置信度阈值,并在关键操作前要求人工确认。
- 权限控制:执行 Agent 需要最小权限原则,避免越权操作。
- 并发与竞态:多个告警同时触发时,需要避免冲突操作。
- 可解释性:每次自动修复都应能解释"为什么做"和"做了什么"。
这四条里我特别想强调并发。多个告警同时来的时候,最常见的事故是"两个 Agent 都觉得 Redis 副本数不够,一个扩到 3,另一个扩到 5,最后一个写覆盖前一个"。解法是协调 Agent 用 traceID 维度做串行化——同一个 target 同时只允许一个 incident 在 Executing 状态。代码不复杂,但忘了加就是事故。
收尾的一些碎碎念
多 Agent 协同是实现复杂 K8s 故障自动修复的有效架构。通过观测、诊断、执行、审计等角色的分工协作,系统可以在保证安全可控的前提下,大幅提升故障响应速度和恢复效率。这也是云原生 AIOps 平台的核心能力之一。
写完这套系统我自己最大的感受是:别让 Agent 之间耦合得太紧。一开始为了图省事,我们让诊断 Agent 直接调执行 Agent 的函数,结果诊断 Agent 跟着执行 Agent 一起挂过好几次。后来强制改成消息驱动,谁也不认识谁,只认识 Message,整个系统稳定性立刻上一个台阶。这玩意儿跟微服务拆分是一个道理——边界比实现更重要。
自测题与动手练习
自测题(合上书能答出来,才算懂):
作者为什么说"单体 controller 不适合 K8s 故障自动修复"?它最致命的风险是什么?
答:职责耦合导致一处 OOM / 打爆 API 就全链路挂掉;拆成多 Agent 是为了故障隔离、炸的时候不连坐。
五大 Agent 里哪个最容易被低估?为什么说它"不是日志,是责任边界"?
答:审计 Agent。没有它,自动修复把集群搞炸了都不知道谁干的,审计记录是事后复盘与问责的依据。
Message envelope 里的
trace_id为什么必须从告警入口就生成、一路传下去?强制 UTC 时间戳又是为什么?答:
trace_id串联一次故障的所有 Agent 输入输出,便于审计与排障;UTC 避免跨时区集群对账时时间错乱。执行 Agent 的四道校验闸门分别是什么?为什么
drain这类破坏性动作要强制人工审批?答:白名单 → 置信度门槛 → 当日次数上限 → 人工审批。drain 会驱逐整节点 Pod,误执行后果严重,必须人确认。
协调 Agent 状态机里"迟到的事件"应该怎么处理?为什么作者倾向于"丢弃 + 审计"而不是"复活"?
答:状态已因超时进入 Failed 后到达的
diagnose_done应丢弃,避免把状态搞乱;复活容易引发状态机不一致。
动手练习(建议真做一遍):
- 用 client-go 起一个 Pod informer,故意制造一个 CrashLoopBackOff 的 Pod,验证观测 Agent 能否识别异常、并带
trace_id发出 Observation 消息。 - 拿一条真实 Pod 崩溃日志,手工拼出一条 Observation JSON,喂给诊断 Agent 的
buildPrompt,看生成的 prompt 是否包含完整现场(status / events / 最近日志 / 相似案例)。 - 在状态机里新增一个
ManualPause状态并补上对应 transition 规则,跑单元测试验证"迟到的 diagnose_done"不会被错误复活。
本章小结
- 多 Agent 协同的本质是故障隔离 + 职责边界,比单体 controller 更抗炸、更好维护。
- 统一消息信封 +
trace_id+ UTC 时间戳,是跨 Agent 协作的命脉。 - 执行 Agent 的护栏(白名单 / 置信度 / 审批)+ 审计 Agent 的外部存储兜底,决定系统能不能"安全地自动"。
- 状态机用显式 transition 表驱动,避免 if-else 串状态;迟到的事件宁可丢弃也别复活。
- 下一篇可深入每个 Agent 的内部算法,比如诊断 Agent 的 RAG 召回与 prompt 工程、执行 Agent 的并发串行化。