人工智能 智能体 编排与云原生 人工智能 应用部署的渐进迁移方案

将智能体工作流从脚本迁到容器编排平台时,常见风险在于任务状态只保存在进程内、重试缺少幂等约束,以及外部模型调用超时后没有可恢复的检查点。先把这些边界写清,再决定是否引入队列、控制器或定时任务,迁移会更可控。

把 Agent 编排从脚本迁到 Kubernetes,不只是给脚本加一层镜像。迁移前应明确任务状态、重试语义、幂等键和人工介入点,再选择队列与控制器的实现方式。

一、看清旧脚本编排的卡顿瓶颈:为什么死板的定时任务撑不住大模型长链调用?

早期验证时,常见做法是用单个 Python 脚本,借助 celery 或简单的 while True 循环拉取待处理任务。它可以帮助验证业务逻辑,但长期运行前仍要补齐状态持久化、超时处理和重复投递的约束。

首先是大模型调用的长尾延迟问题。一次复杂的 RAG(检索增强生成)或者多步思考 Agent 链条,消耗的时间可能从几秒到几分钟不等。传统的 Task Queue 假设每个任务的执行时间是相对均匀且可预测的。一旦上游模型出现响应变慢,worker 进程就会全部堵塞在 HTTP 等待上,后面的任务瞬间堆积。

其次是状态无法追踪。旧脚本通常把状态写在内存变量或者 Redis 字符串里。当节点资源紧张被 Kubelet 驱逐,或者 Pod 超过 Memory Limit 被系统杀死时,正在执行的 Agent 思考上下文就彻底丢失了。重新启动后,由于不知道上一次到底卡在哪一步,系统只能从头重新调一遍 LLM,造成 API Token 的巨大浪费和数据重复写入。

+-----------------------------------------------------------------------+
|                         旧架构:死板脚本与轮询                             |
|                                                                       |
|  [ CronJob / Python ]  --->  [ 数据库/Redis ]  --->  [ 阻塞调用 LLM API ]  |
|          |                         |                     |            |
|     API Server 压力              状态丢失               OOM / 超时崩溃     |
+-----------------------------------------------------------------------+

二、状态机与事件驱动结合:画出云原生 Agent 状态流转的具体边界。

要解决上述问题,首要任务是将 Agent 的执行逻辑重构为基于事件驱动的状态机。每一个 Agent 动作(比如“检索文档”、“调用工具”、“生成文本”、“校验输出”)都是一个独立的状态节点。节点之间通过 NATS 或 Kafka 等轻量级消息队列进行解耦。

状态变更应具备幂等键、版本号或条件更新,并将可恢复的检查点持久化。消息队列通常提供至少一次投递,消费者还需要处理重复消息、超时任务和外部工具调用已执行但结果未落库的情况。

采用这种架构后,执行单元可以尽量保持无状态,任务进度由事件和持久化检查点承接。扩缩容、进程退出或消息重复投递仍需结合实际的任务租约、终止宽限期和幂等设计验证,不能仅靠架构图推断结果。

三、用 Go 编写自愈型 Agent 调度控制器:拦截超时与脏状态的核心逻辑。

在云原生环境中,我们需要一个自定义控制器(Controller)来监控 Agent 任务的生命周期。下面的代码展示了一个轻量级 Agent 任务调度器核心逻辑,包含上下文超时控制、自动重试机制以及异常捕获处理。

package main

import (
	"context"
	"errors"
	"fmt"
	"log"
	"sync"
	"time"
)

// AgentState 表示 Agent 任务的状态类型
type AgentState string

const (
	StatePending   AgentState = "PENDING"
	StateRunning   AgentState = "RUNNING"
	StateSuccess   AgentState = "SUCCESS"
	StateFailed    AgentState = "FAILED"
	StateRetrying  AgentState = "RETRYING"
)

// TaskContext 保存 Agent 执行上下文
type TaskContext struct {
	ID         string
	Payload    string
	Step       int
	MaxRetries int
	CurrentTry int
	State      AgentState
	Mu         sync.Mutex
}

// AgentController 调度控制器结构体
type AgentController struct {
	taskQueue chan *TaskContext
	timeout   time.Duration
}

// NewAgentController 初始化控制器
func NewAgentController(queueSize int, timeout time.Duration) *AgentController {
	return &AgentController{
		taskQueue: make(chan *TaskContext, queueSize),
		timeout:   timeout,
	}
}

// DispatchTask 分发任务到队列中
func (c *AgentController) DispatchTask(task *TaskContext) error {
	select {
	case c.taskQueue <- task:
		log.Printf("[Info] 任务 %s 已成功入队", task.ID)
		return nil
	default:
		return errors.New("任务队列已满,无法分发")
	}
}

// ExecuteStep 模拟执行单个 Agent 步骤,包含超时与边界处理
func (c *AgentController) ExecuteStep(parentCtx context.Context, task *TaskContext) error {
	task.Mu.Lock()
	task.State = StateRunning
	task.Mu.Unlock()

	// 为当前步骤设置独立超时限制
	ctx, cancel := context.WithTimeout(parentCtx, c.timeout)
	defer cancel()

	errChan := make(chan error, 1)

	go func() {
		defer func() {
			if r := recover(); r != nil {
				errChan <- fmt.Errorf("Agent 内部发生 Panic 恢复: %v", r)
			}
		}()

		// 模拟 Agent 调用外部 LLM 或工具的耗时操作
		if task.Payload == "" {
			errChan <- errors.New("无效的载荷 Payload 不能为空")
			return
		}

		// 假装做一些工作
		time.Sleep(200 * time.Millisecond)
		errChan <- nil
	}()

	select {
	case <-ctx.Done():
		task.Mu.Lock()
		task.CurrentTry++
		if task.CurrentTry <= task.MaxRetries {
			task.State = StateRetrying
			log.Printf("[Warn] 任务 %s 超时,准备第 %d 次重试", task.ID, task.CurrentTry)
		} else {
			task.State = StateFailed
			log.Printf("[Error] 任务 %s 达到最大重试次数,标记为失败", task.ID)
		}
		task.Mu.Unlock()
		return ctx.Err()

	case err := <-errChan:
		task.Mu.Lock()
		defer task.Mu.Unlock()
		if err != nil {
			task.CurrentTry++
			if task.CurrentTry <= task.MaxRetries {
				task.State = StateRetrying
			} else {
				task.State = StateFailed
			}
			return fmt.Errorf("执行步骤失败: %w", err)
		}
		task.State = StateSuccess
		log.Printf("[Info] 任务 %s 步骤 %d 执行成功", task.ID, task.Step)
		return nil
	}
}

func main() {
	controller := NewAgentController(10, 500*time.Millisecond)
	ctx := context.Background()

	task := &TaskContext{
		ID:         "task-agent-001",
		Payload:    "检索云原生最佳实践文档并总结",
		Step:       1,
		MaxRetries: 2,
		CurrentTry: 0,
		State:      StatePending,
	}

	if err := controller.DispatchTask(task); err != nil {
		log.Fatalf("分发失败: %v", err)
	}

	t := <-controller.taskQueue
	if err := controller.ExecuteStep(ctx, t); err != nil {
		log.Printf("任务最终状态处理结果: %v, 当前状态: %s", err, t.State)
	}
}

四、Kubelet 探针与 GPU 资源隔离配置:线上部署时如何避免盲目抢占?

在 Kubernetes 中部署 AI Agent 应用时,资源配置必须极度小心。大模型推理依赖 GPU 或者高配置 CPU,如果配置不当,非常容易触发节点级的连锁反应。

第一,就地配置合理的 HTTP / gRPC 健康检查探针。普通的 LivenessProbe 探针如果直接去调用 Agent 的推理接口,会导致探针本身被阻塞,引发 Kubelet 误判并不断重启 Pod。正确的做法是在应用内部开放一个轻量级的 /healthz 路由,只检查本地数据库连接、消息队列 Channel 状态以及 GPU 驱动初始化情况。

第二,利用 Resource Limits 与 Request 严格限制 Memory 和 Ephemeral-Storage。大模型 Agent 在处理长文本时,内存开销会随 Context Length 呈现非线性增长。必须显式指定 limits.memory,防止单 Pod 疯狂吃光宿主机内存。

下面是一个标准的 Agent Deployment 配置样式:

apiVersion: apps/v1
kind: Deployment
metadata:
  name: ai-agent-worker
  namespace: ai-platform
spec:
  replicas: 3
  selector:
    matchLabels:
      app: ai-agent-worker
  template:
    metadata:
      labels:
        app: ai-agent-worker
    spec:
      containers:
      - name: agent-runner
        image: registry.example.com/ai/agent-runner:v1.2.4
        resources:
          requests:
            cpu: "2"
            memory: "4Gi"
            nvidia.com/gpu: "1"
          limits:
            cpu: "4"
            memory: "8Gi"
            nvidia.com/gpu: "1"
        livenessProbe:
          httpGet:
            path: /healthz
            port: 8080
          initialDelaySeconds: 15
          periodSeconds: 10
          timeoutSeconds: 3
        readinessProbe:
          httpGet:
            path: /ready
            port: 8080
          initialDelaySeconds: 5
          periodSeconds: 5
        env:
        - name: WORKER_CONCURRENCY
          value: "4"
        - name: NATS_URL
          value: "nats://nats-cluster.ai-platform:4222"

五、压测日志与告警观测:压测过程中如何捕捉节点级别的抖动?

当完成了架构迁移并部署到 Kubernetes 集群后,必须通过真实压测来验证系统的稳定性。不能仅仅盯着 HTTP 响应码,而要全方位监控 Kubernetes 节点指标与 Pod 运行状态。

在压测过程中,使用以下命令对集群进行实时排障与诊断:

# 检查 Agent Pod 的资源使用率,查看是否有 Pod 逼近 memory limit
kubectl top pods -n ai-platform --containers

# 实时查看节点级别的 OOM 告警事件
kubectl get events -n ai-platform --field-selector reason=OOMKilled -w

# 查看 GPU 显存与利用率(针对搭载 NVIDIA GPU 的节点)
kubectl exec -it -n ai-platform ai-agent-worker-6789ab-c1d2 -- nvidia-smi --query-gpu=utilization.gpu,utilization.memory,memory.used,memory.free --format=csv -l 1

# 检查 NATS / Kafka 消息队列的堆积情况
kubectl exec -it -n ai-platform nats-0 -- nats-top

# 使用 pprof 诊断 Go 控制器内存与 Goroutine 泄露(若开启了 pprof 端点)
go tool pprof -http=:8081 http://127.0.0.1:6060/debug/pprof/goroutine

这些命令可以帮助定位队列积压、资源压力与调用超时。是否要改为事件驱动状态机,应以任务量、恢复要求和维护成本评估;探针与资源限制只能约束故障影响,不能替代业务状态恢复设计。

更多推荐