Agent 服务网格化:像治理微服务一样治理智能体集群
Agent 服务网格化像治理微服务一样治理智能体集群一、当你线上跑了 200 个 Agent 实例却没有一个统一的流量治理层Agent 上生产之后最容易被忽视的一层是路由层。大多数团队的做法是一个 Agent 对应一个 Pod前端直接请求这个 Pod 的 endpoint。初期 Agent 数量少3-5 个这么搞没问题。但当 Agent 数量膨胀到几十上百个时你会发现各种问题哪个 Agent 挂了前端不知道没有健康检查流量不均匀有的 Agent 被频繁调用有的闲置版本滚动时新旧 Agent 同时存在没有金丝雀发布。这不是 Agent 特有的问题。五年前微服务也是这么走过来的——从直接 IP 调用到服务注册发现再到Service Mesh 全面治理。Agent 集群现在正站在微服务当年的起点上。解决思路很直接把微服务治理体系中成熟的模式搬过来。不是照搬 Istio 的代码而是搬它的思想sidecar 代理、流量分割、熔断限流、可观测性。Agent 本质上就是一个有状态的微服务所以微服务网格的那套治理框架完全适用。二、底层机制与原理剖析Agent 服务网格的核心架构每个 Agent Pod 内注入一个轻量级 sidecar类似 Envoy由 sidecar 负责所有入站、出站流量的代理。Control Plane 管理全局配置下发路由规则、熔断策略、限流参数。核心设计要点Sidecar 注入每个 Agent Pod 启动时自动注入 sidecar 容器。Sidecar 对 Agent 进程透明——Agent 以为自己直接监听0.0.0.0:8080实际上流量先到 sidecar 的127.0.0.1:15001经过策略检查后再转发到 Agent 的 8080 端口。这层 iptables 或 eBPF 的透明代理是整个网格的基础。流量分割版本发布时Control Plane 下发规则v1 Agent 处理 90% 流量v2 Agent 处理 10% 流量。Sidecar 根据这个规则在本地做权重随机选择。不需要修改 Agent 代码不需要改 DNS单纯的流量层操作。出站流量管控Agent 调用外部 API搜索、计算、数据库时出站流量也经过 sidecar。这里可以加限流防止 Agent 过度调用外部 API、熔断外部 API 异常时快速失败、重试临时故障自动重试避免 Agent 把错误传递给用户。三、生产级代码实现// agent-mesh/main.go // Agent Sidecar 核心实现 // 设计思路不依赖任何厂商如 Istio 的 xDS 协议 // 用标准 gRPC 与 Control Plane 通信降低绑定风险 package main import ( context fmt log net sync time google.golang.org/grpc google.golang.org/grpc/health/grpc_health_v1 ) // --------------------------------------------------------------------------- // 配置定义全部从 Control Plane 拉取不允许本地硬编码 // 原因运行时参数变更不应重启 Sidecar必须支持热更新 // --------------------------------------------------------------------------- type MeshConfig struct { mu sync.RWMutex // 保护热更新时的并发读写 AgentPort int json:agent_port // Agent 实际监听端口 SidecarPort int json:sidecar_port // Sidecar 暴露端口 Routes []RouteRule json:routes CircuitBreakers []CircuitRule json:circuit_breakers RateLimits []RateLimitRule json:rate_limits } type RouteRule struct { AgentType string json:agent_type // 按 Agent 类型路由 Version string json:version // v1 / v2 Weight float64 json:weight // 流量权重 0-1 } type CircuitRule struct { AgentType string json:agent_type MaxConsecutive int json:max_consecutive // 连续失败次数阈值 Timeout time.Duration json:timeout HalfOpenMax int json:half_open_max // 半开状态允许的探测请求数 } type RateLimitRule struct { AgentType string json:agent_type QPS int json:qps // 每秒允许的请求数 } // --------------------------------------------------------------------------- // Sidecar 核心结构 // --------------------------------------------------------------------------- type AgentSidecar struct { config *MeshConfig registry *AgentRegistry // Agent 实例注册表 // 熔断状态机AgentType - 状态 breakerStates map[string]*CircuitState // 令牌桶限流AgentType - Limiter limiterBuckets map[string]*TokenBucket grpcServer *grpc.Server } // --------------------------------------------------------------------------- // Agent 注册表管理所有 Agent 实例的生命周期 // 为什么不用 etcd/consul小规模集群100 实例用内存表 心跳就够了 // 避免引入额外的分布式协调依赖 // --------------------------------------------------------------------------- type AgentRegistry struct { instances map[string]*AgentInstance // key: agent_id mu sync.RWMutex } type AgentInstance struct { ID string AgentType string Version string Address string LastSeen time.Time Status string // READY / DEGRADED / OFFLINE } func (r *AgentRegistry) Register(inst *AgentInstance) { r.mu.Lock() defer r.mu.Unlock() inst.LastSeen time.Now() inst.Status READY r.instances[inst.ID] inst log.Printf([Registry] Agent %s (%s/%s) registered at %s, inst.ID, inst.AgentType, inst.Version, inst.Address) } // 健康检查Sidecar 定期探测挂载的 Agent 进程 // 如果 Agent 无响应标记 DEGRADED触发流量剔除 func (r *AgentRegistry) HealthCheck(ctx context.Context) { ticker : time.NewTicker(5 * time.Second) defer ticker.Stop() for { select { case -ctx.Done(): return case -ticker.C: r.mu.RLock() // 复制一份快照避免长时间持锁 snapshot : make([]*AgentInstance, 0, len(r.instances)) for _, inst : range r.instances { cp : *inst snapshot append(snapshot, cp) } r.mu.RUnlock() for _, inst : range snapshot { if err : probeAgent(inst.Address, 2*time.Second); err ! nil { r.mu.Lock() // 只改状态字段不重新赋值整个对象 if target, ok : r.instances[inst.ID]; ok { target.Status DEGRADED } r.mu.Unlock() log.Printf([HealthCheck] Agent %s unhealthy: %v, inst.ID, err) } else { // 恢复 r.mu.Lock() if target, ok : r.instances[inst.ID]; ok { target.Status READY target.LastSeen time.Now() } r.mu.Unlock() } } } } } // probeAgent 通过 TCP 连接探测 Agent 端口是否可连通 func probeAgent(address string, timeout time.Duration) error { conn, err : net.DialTimeout(tcp, address, timeout) if err ! nil { return err } conn.Close() return nil } // --------------------------------------------------------------------------- // 令牌桶限流实现 // 设计要点使用 atomic 操作 无锁避免高并发下的锁竞争 // 缺点不够精确允许少量超发但在流量保护场景可以接受 // --------------------------------------------------------------------------- type TokenBucket struct { rate float64 // 每秒生成的令牌数 burst int // 突发容量 tokens float64 // 当前令牌数用 float64 避免整数截断 lastRefill time.Time mu sync.Mutex } func NewTokenBucket(rate int, burst int) *TokenBucket { return TokenBucket{ rate: float64(rate), burst: burst, tokens: float64(burst), lastRefill: time.Now(), } } func (tb *TokenBucket) Allow() bool { tb.mu.Lock() defer tb.mu.Unlock() // 补充令牌 now : time.Now() elapsed : now.Sub(tb.lastRefill).Seconds() tb.tokens elapsed * tb.rate if tb.tokens float64(tb.burst) { tb.tokens float64(tb.burst) } tb.lastRefill now // 消耗令牌 if tb.tokens 1.0 { tb.tokens - 1.0 return true } return false } // --------------------------------------------------------------------------- // 主入口启动 Sidecar // --------------------------------------------------------------------------- func main() { ctx, cancel : context.WithCancel(context.Background()) defer cancel() sidecar : AgentSidecar{ config: MeshConfig{ SidecarPort: 15001, AgentPort: 8080, }, registry: AgentRegistry{ instances: make(map[string]*AgentInstance), }, breakerStates: make(map[string]*CircuitState), limiterBuckets: make(map[string]*TokenBucket), } // 启动健康检查循环 go sidecar.registry.HealthCheck(ctx) // 启动 gRPC 入站代理 lis, err : net.Listen(tcp, fmt.Sprintf(:%d, sidecar.config.SidecarPort)) if err ! nil { log.Fatalf(Failed to listen: %v, err) } // 注册健康检查服务给 K8s Probe 用 sidecar.grpcServer grpc.NewServer() grpc_health_v1.RegisterHealthServer(sidecar.grpcServer, healthServer{}) log.Printf([Sidecar] Listening on :%d, proxying to Agent on :%d, sidecar.config.SidecarPort, sidecar.config.AgentPort) if err : sidecar.grpcServer.Serve(lis); err ! nil { log.Fatalf(Failed to serve: %v, err) } } // healthServer 实现 gRPC Health Checking Protocol type healthServer struct{} func (h *healthServer) Check(ctx context.Context, req *grpc_health_v1.HealthCheckRequest) (*grpc_health_v1.HealthCheckResponse, error) { return grpc_health_v1.HealthCheckResponse{ Status: grpc_health_v1.HealthCheckResponse_SERVING, }, nil }四、边界分析与架构权衡适用场景Agent 实例数 20需要统一流量管理多 Agent 类型共存不同类型有不同的路由/限流策略需要金丝雀发布新版本 Agent 先覆盖 10% 流量验证不适用场景Agent 实例 5 个引入 sidecar 增加了部署复杂度收益不大所有 Agent 共享无状态逻辑不需要流量精细控制延迟极为敏感的场景sidecar 增加的 1-2ms hop 都不可接受架构代价Sidecar 注入增加了每个 Pod 的资源开销约 50MB 内存 0.1 CPU 核心Control Plane 本身是高可用系统需要至少 2 个副本运维复杂度上升——出问题时需要排查 sidecar 日志 Agent 日志两层五、总结Agent 服务网格这件事本质是把微服务治理的成熟经验做个针对性的减法。不需要 xDS 的全部复杂度只需要健康检查、流量分割、熔断限流这几个核心能力。用轻量 sidecar 简单的 Control Plane 就能解决 Agent 集群上生产后面临的最基础的流量治理问题。关键是别等到 Agent 数量爆炸了再做——那时迁移成本就不可控了。