☰
Seldon Core v2 Agent API 深度解析:Scheduler 与 Agent 的模型生命周期管理 gRPC 协议
2026/10/6 7:35:03 网站建设 项目流程
  • 人工智能
  • 大模型
  • MLOps
  • LLMOps
  • 模型推理服务
  • 云原生
  • 后端

【免费下载链接】seldon-core

An MLOps framework to package, deploy, monitor and manage thousands of production machine learning models

项目地址:https://gitcode.com/gh_mirrors/se/seldon-core
点击查看免费下载

导读

Agent API 是 Seldon Core v2 中 Scheduler(调度器)与 Agent(与每个推理服务器同 Pod 部署的代理组件)之间的控制平面 gRPC 通信协议,负责模型在推理服务器上的加载与卸载、模型事件上报、服务器排空(drain)以及模型副本自动伸缩。读完本文,你将完整掌握AgentService的四个 RPC 及其消息语义、模型的加载/卸载生命周期状态机、排空与重调度机制,以及 Agent 作为数据平面反向代理的实现原理,并了解如何通过 CLI 参数部署与调优 Agent。

1. Agent API 定位与架构背景

在 Seldon Core v2 的架构中,Scheduler 是控制平面的核心,负责全局的模型调度决策;而 Agent 则"运行在每个推理服务器旁边",承担两项职责:

  • 控制平面:代表服务器向 Scheduler 注册自身(服务器名、副本索引、内存容量、能力清单),接收 Scheduler 下发的模型加载/卸载指令,并上报模型事件与可用内存;
  • 数据平面:作为反向代理(reverse proxy),将来自 Envoy 的推理请求转发到本地推理服务器,并支持"懒加载"(请求到达时模型未在内存则现场加载后重试)。

Agent API 正是连接这两条平面、打通"调度决策"与"落地执行"的契约。它在仓库中的权威定义位于 apis/mlops/agent/agent.proto,生成代码位于 apis/go/mlops/agent。

上图展示了 Seldon Core v2 的整体架构:Agent 位于 Envoy 与推理服务器(MLServer、Triton 等)之间,通过 gRPC 与 Scheduler 双向通信,同时承载数据面的推理请求落地执行。

2. 协议定义:AgentService 与四个 RPC 概览

AgentService是一个包含四个 RPC 的 gRPC 服务(agent.proto):

service AgentService { rpc AgentEvent(ModelEventMessage) returns (ModelEventResponse) {}; rpc Subscribe(AgentSubscribeRequest) returns (stream ModelOperationMessage) {}; rpc ModelScalingTrigger(stream ModelScalingTriggerMessage) returns (ModelScalingTriggerResponse) {}; rpc AgentDrain(AgentDrainRequest) returns (AgentDrainResponse) {}; }
RPC流模式方向用途
Subscribe服务端流式(server streaming)Agent → Scheduler 发起,Scheduler → Agent 持续下发Agent 注册服务器副本并接收模型操作指令(加载/卸载)
AgentEvent一元(unary)Agent → SchedulerAgent 上报模型事件(加载成功/失败、卸载成功/失败、内存信息)
AgentDrain一元(unary)Agent → Scheduler请求排空某个服务器副本,触发模型迁移
ModelScalingTrigger客户端流式(client streaming)Agent → Scheduler上报模型伸缩触发事件(扩容/缩容)

服务端实现位于 scheduler/pkg/agent/server.go,gRPC 服务同时注册了健康检查接口(HealthCheckService),并可通过StartGrpcServer(allowPlainTxt, agentPort, agentTlsPort)同时启动明文与 mTLS 两个端口(server.go),与components/tls模块协同实现传输层安全。

3. 订阅注册:Subscribe 与 ReplicaConfig

Agent 启动后,由 agent_svc_manager.go 中的handleSchedulerSubscription发起Subscribe,携带如下信息:

message AgentSubscribeRequest { string serverName = 1; bool shared = 2; uint32 replicaIdx = 3; ReplicaConfig replicaConfig = 4; repeated ModelVersion loadedModels = 5; uint64 availableMemoryBytes = 6; }

其中ReplicaConfig描述了该服务器副本的完整能力画像(agent.proto):

字段类型含义
inferenceSvcstring推理服务的 DNS 名称
inferenceHttpPortint32推理 HTTP 端口
inferenceGrpcPortint32推理 gRPC 端口
memoryBytesuint64服务器副本的内存容量
capabilitiesrepeated string服务器能力清单,如sklearn、pytorch、xgboost、mlflow
overCommitPercentageuint32允许的超卖内存百分比,设为 0(%)表示禁止超卖

loadedModels是 Agent 启动时已加载模型的快照,availableMemoryBytes是考虑超卖后的可用内存。这两项用于网络抖动后的状态对账:Scheduler 端在Subscribe处理中会调用scheduleModelsFromRequest(server.go),把 Agent 上报的已加载模型重新纳入调度,并重试此前失败的模型(ScheduleFailedModels),从而避免"网络闪断导致模型被调度到其他服务器,而本机其实还加载着"的重复加载。

服务端Subscribe实现的几个关键细节(server.go):

  1. 按副本串行化:使用agentMutex(sync.Map)对同一(serverName, replicaIdx)强制串行——保证旧 Agent 完全断开后新 Agent 才能接入;
  2. 注册副本:将请求存入store.AddServerReplica,纳入全局调度视图;
  3. 流式长连接:阻塞在select上等待fin通道或上下文取消;一旦 Agent 断开,立即删除注册并通过removeServerReplicaImpl(server.go)将该副本上的模型重新调度到其他可用副本,同时再次重试LoadFailed状态的模型。

客户端侧连接失败时使用指数退避重连(boff.RetryNotify+util.GetClientExponentialBackoff),并在成功建立连接后:若运行在 K8s 中,会先等待 Pod 的 IP 发布到Endpoints(HasPublishedIP)再宣告就绪,避免 Envoy 将请求路由到尚未就绪的旧 Pod IP 而出现 503(agent_svc_manager.go)。

4. 模型操作指令:ModelOperationMessage

Subscribe建立的服务端流,持续向 Agent 下发ModelOperationMessage:

message ModelOperationMessage { enum Operation { UNKNOWN_EVENT = 0; LOAD_MODEL = 1; UNLOAD_MODEL = 2; } Operation operation = 1; ModelVersion modelVersion = 2; bool autoscalingEnabled = 3; }

其中ModelVersion内嵌了完整的scheduler.Model(来自 apis/mlops/scheduler/scheduler.proto),包含模型元数据、ModelSpec(存储配置、运行时信息、内存占用)与DeploymentSpec(副本数、min/max replicas)。autoscalingEnabled标记该模型是否启用了基于指标的副本伸缩,供 Agent 侧决定是否挂接伸缩统计。

Scheduler 端通过Sync(modelName)方法(server.go)决定下发何种指令:

  • 对最新版本中处于LoadRequested状态的副本发送LOAD_MODEL,并同步将状态推进到Loading;
  • 对任意版本中处于UnloadRequested状态的副本发送UNLOAD_MODEL,并推进到Unloading;
  • 发送失败时会把状态置为LoadFailed/UnloadFailed并记录错误信息。

Agent 客户端收到指令后在handleSchedulerSubscription的switch operation.Operation中分发到LoadModel与UnloadModel(agent_svc_manager.go)。

4.1 加载模型(LoadModel)

LoadModel的完整流水线(agent_svc_manager.go):

  1. 乱序防护:基于单调时钟的时间戳记录(modelTimestamps)忽略乱序到达的过期指令;
  2. 获取存储配置:从ModelSpec.StorageConfig中解析 rclone 配置或 K8s Secret(getArtifactConfig,支持StorageRcloneConfig与StorageSecretName两种形态);
  3. 下载模型工件:通过ModelRepository.DownloadModelVersion(基于 rclone)将模型拉取到本地/mnt/agent/models;
  4. 加载到推理服务器:调用模型服务器控制面客户端(v2Client.LoadModelVersion,支持 MLServer 与 Triton,工厂实现见 modelserver_controlplane/factory),失败时按maxLoadRetryCount/maxLoadElapsedTime退避重试;
  5. 可选挂接伸缩统计:若AutoscalingEnabled且 Agent 启用了伸缩,则把模型加入StatsAnalyserService;
  6. 上报成功:发送LOADED事件。

4.2 卸载模型(UnloadModel)

UnloadModel(agent_svc_manager.go)与加载对称,但多了一个关键的前置步骤:卸载宽限(unloadGraceTime)。由于 Envoy 是最终一致(eventually consistent)的,立即卸载会导致在途请求打到已卸载的模型上,因此 Agent 先睡眠宽限期让 Envoy 收敛集群变化,再执行卸载、从 rclone 仓库清理模型、发送UNLOADED事件。

5. 事件上报:AgentEvent 与 ModelEventMessage

Agent 通过AgentEvent一元 RPC 把模型状态变化告知 Scheduler,消息体为ModelEventMessage:

message ModelEventMessage { string serverName = 1; uint32 replicaIdx = 2; string modelName = 3; uint32 modelVersion = 4; enum Event { UNKNOWN_EVENT = 0; LOAD_FAIL_MEMORY = 1; LOADED = 2; LOAD_FAILED = 3; UNLOADED = 4; UNLOAD_FAILED = 5; REMOVED = 6; // unloaded and removed from local PVC REMOVE_FAILED = 7; RSYNC = 9; // Ask server for all models that need to be loaded } Event event = 5; string message = 6; uint64 availableMemoryBytes = 7; scheduler.ModelRuntimeInfo runtimeInfo = 8; }

事件枚举中LOAD_FAIL_MEMORY专门表示因内存不足导致的加载失败;REMOVED/REMOVE_FAILED表示"从本地 PVC 中移除"之后的终态;RSYNC用于请求服务器全量同步。runtimeInfo携带模型在服务器上的实际运行时信息(如 MLServer 的parallelWorkers、Triton 的实例数),这些信息在 model_state.go 中用于计算模型占用的真实内存。

服务端AgentEvent的状态映射(server.go)构成了模型副本状态机的核心转移:

上报事件期望前置状态(expected)目标状态(desired)
LOADEDLoadingLoaded
UNLOADEDUnloadingUnloaded
LOAD_FAILED/LOAD_FAIL_MEMORYLoadingLoadFailed
UNLOAD_FAILEDUnloadingUnloadFailed

Scheduler 使用乐观期望状态校验,防止状态乱序回跳;同时将availableMemoryBytes与runtimeInfo一并写入 store,供后续调度决策(如内存感知的模型放置、副本扩展)使用。

6. 排空机制:AgentDrain

当服务器副本需要下线(如节点驱逐、滚动更新)时,Agent 调用AgentDrain(AgentDrainRequest{serverName, replicaIdx})请求 Scheduler 排空该副本。

Scheduler 端drainServerReplicaImpl(server.go)的执行序列:

  1. store.DrainServerReplica将副本上所有模型标记为待迁移;
  2. 用modelRelocatedWaiter为这些模型注册等待组——模型在其他副本上变为Available时被signalModel释放;
  3. 睡眠agentDrainCoolDownPeriod(500ms)作为冷却期,避免多个 Agent 并发排空时调度器把模型调度到同样在排空的服务器上;
  4. 逐个scheduler.Schedule(modelName)把模型重新调度到健康副本;
  5. 阻塞等待所有模型迁移完成(waiter.wait);
  6. 再额外等待EnvoyUpdateDefaultBatchWait + serverDrainingExtraWaitMillis(3000ms),让 Envoy 分批更新最终收敛,之后才认为排空完成。

客户端侧drainOnRequest(agent_svc_manager.go)由 drainservice 的触发信号驱动:一旦触发,Agent 置isDraining=true(此后Ready()返回 false,Pod 从负载均衡摘除)、发送AgentDrain,并释放/terminate的等待。Agent 排空期间若模型加载成功,会被sendAgentEvent主动"取消"为LOAD_FAILED(原因是 server replica is draining),防止状态不一致(agent_svc_manager.go)。

7. 模型自动伸缩:ModelScalingTrigger

ModelScalingTrigger是客户端流式 RPC,Agent 通过它向 Scheduler 上报伸缩事件:

message ModelScalingTriggerMessage { string serverName = 1; uint32 replicaIdx = 2; string modelName = 3; uint32 modelVersion = 4; enum Trigger { SCALE_UP = 0; SCALE_DOWN = 1; } Trigger trigger = 5; uint32 amount = 6; // number of replicas required map<string,uint32> metrics = 7; // optional metrics to expose to the scheduler }

触发来源是 Agent 内的模型伸缩统计服务 modelscaling:基于推理延迟滞后的ScaleUpEvent(阈值由ModelInferenceLagThreshold配置)与基于模型最近使用时间的ScaleDownEvent(阈值由ModelInactiveSecondsThreshold配置)。modelScalingEventsConsumer从事件通道读取并转发到客户端流(agent_svc_manager.go)。

Scheduler 端处理(server.go 与createScalingPseudoRequest/calculateDesiredNumReplicas)要点:

  • 仅当autoscalingModelEnabled时受理,否则返回Unimplemented;
  • 通过createScalingPseudoRequest构造一个"伪加载请求":校验模型存在、事件版本与最新版本一致;
  • 扩容:副本数 +1;缩容:副本数 −1,但缩容前要求模型状态稳定(最近modelScalingCoolingDownSeconds=60s 内没有状态变化),避免在加载/卸载震荡期收缩;
  • checkModelScalingWithinRange强制约束:只有设置了minReplicas/maxReplicas才允许伸缩,目标副本数不得低于minReplicas且不低于 1,不得高于maxReplicas(server.go);
  • 校验通过后写入 store 并重新触发调度。

需要注意:仓库中 cmd/agent/main.go 当前将autoScalingEnabled硬编码为false(注释说明在扩容问题解决前强制禁用),因此上述伸缩链路在默认构建下处于关闭状态,相关 CLI 阈值参数仅作预留。

8. 数据平面:Agent 反向代理

Agent 的数据平面职责由两个反向代理子服务承担,它们都注册为CriticalDataPlaneService,任一失败都会导致 Agent 不可就绪(agent_svc_manager.go)。

8.1 HTTP/REST 反向代理

实现见 rproxy.go,基于httputil.ReverseProxy并自定义lazyModelLoadTransport:

  • 懒加载:请求经addHandlers先调用stateManager.EnsureLoadModel(确保模型在内存),随后通过rewritePath把 URL 中的外部模型名改写为内部模型名(去掉/versions/<ver>段);
  • 404/400 重试:当后端返回404 Not Found(或 Triton 将未加载模型视为400 Bad Request)时,先触发loader加载模型,再用缓存的原请求体重放一次请求(rproxy.go);
  • OpenAI API 翻译:内置了/chat/completions、/embeddings、/images/generations三条路径的 OpenAI 格式翻译器(OpenAIChatCompletionsTranslator等),把 OpenAI 风格的请求/响应在反向代理边界翻译为 Open Inference Protocol(OIP)格式(rproxy.go);
  • 可观测性:通过otelhttp.NewHandler注入 OpenTelemetry 追踪,并统计推理耗时、HTTP 状态码等指标(AddModelInferMetrics),同时透传/生成requestId头用于端到端关联。

8.2 gRPC 反向代理

实现见 rproxy_grpc.go,对外暴露v2_dataplane.GRPCInferenceService,默认监听端口为ReverseGRPCProxyPort = 9998(rproxy_grpc.go)。支持的方法:

  • ModelInfer:一元推理,从 gRPC 元数据中提取SeldonInternalModelHeader/SeldonModelHeader定位模型,EnsureLoadModel后转发,NotFound/Unavailable时懒加载重试;对后端维护 10 连接连接池随机取用;
  • ModelStreamInfer:双向流推理,通过两个forwardStreamgoroutine 分别转发"客户端→后端"与"后端→客户端"两个方向的消息,任一方向出错即取消整体流并释放资源(rproxy_grpc.go);
  • ModelMetadata/ModelReady:同样执行EnsureLoadModel与懒加载重试。

代理还会将请求/响应 Trailer 中透传的requestId写回客户端(setTrailer),并记录 gRPC 状态码指标。

9. Agent 本地状态与缓存管理

Agent 的本地状态由LocalStateManager统一管理(state_manager.go),配套ModelState(model_state.go)维护"已加载模型 → 内存占用"映射,并用 LRU 事务缓存(CacheTransactionManager)实现最近最少使用驱逐:

  • 加载(LoadModelVersion):校验版本、计算内存增量、makeRoomIfNeeded驱逐足够多的 LRU 模型腾出空间、乐观扣减可用内存、调用 v2 控制面加载、加入缓存;
  • 卸载(UnloadModelVersion):从缓存删除、回补内存;若模型已不在缓存(已被驱逐)则只更新记账;
  • 数据面兜底(EnsureLoadModel):推理请求到达时若模型被驱逐出内存,现场重新加载并重放请求;
  • 内存与超卖:availableMainMemoryBytes反映主内存剩余;GetOverCommitMemoryBytes()返回overCommitPercentage/100 × totalMainMemoryBytes的超卖额度,GetAvailableMemoryBytesWithOverCommit在订阅与事件上报中作为"可承诺内存"提供给 Scheduler,使调度器可以按超卖后的容量放置模型(state_manager.go)。

10. Agent 运行配置与 CLI 参数

Agent 进程入口为 scheduler/cmd/agent/main.go,启动序列包括:创建模型仓库与 rclone 目录、启动就绪服务(readyservice)、创建 K8s 客户端(K8s 环境)、初始化 OpenTelemetry tracer、启动 rclone 客户端与模型仓库、创建推理服务器控制面客户端、启动 Prometheus 指标服务、创建 HTTP/gRPC 反向代理、调试服务与排空服务,最后等待子服务就绪后进入StartControlLoop(带退避重连的订阅循环)。

全部启动参数定义于 scheduler/cmd/agent/cli/flags.go,关键参数如下:

参数默认值说明
--server-namemlserver服务器名称,用于在 Scheduler 中唯一标识
--server-idx0服务器副本索引(ReplicaIdx)
--scheduler-host/--scheduler-port/--scheduler-tls-port0.0.0.0/ 默认端口Scheduler 地址与明文/mTLS 端口
--rclone-host/--rclone-port0.0.0.0/ 默认端口rclone 下载服务地址
--inference-host/--inference-http-port/--inference-grpc-port0.0.0.0/ 默认端口推理服务器地址与 HTTP/gRPC 端口
--reverse-proxy-http-port/--reverse-proxy-grpc-port默认 HTTP 端口 /9998数据平面反向代理监听端口
--server-typemlserver模型服务器类型(mlserver或triton),决定仓库处理器与控制面客户端工厂
--memory-bytes1000000服务器可用内存(字节)
--capabilitiessklearn,xgboost服务器能力清单(逗号分隔),如 sklearn、pytorch、xgboost、mlflow
--overcommit-percentage0内存超卖百分比,0 表示关闭
--replica-config空直接以 JSON 传入完整ReplicaConfig(优先于上述散列参数)
--agent-folder/mnt/agent模型仓库根目录(其下含models/与rclone/)
--config-path/mnt/config配置文件目录(agent.yaml / agent.json),见 scheduler/config/agent.yaml
--namespace空K8s 命名空间;非空即视为运行在 K8s 中
--max-load-elapsed-time-minutes/--max-load-retry-count默认值模型加载的超时上限与重试次数
--max-unload-elapsed-time-minutes/--max-unload-retry-count默认值模型卸载的超时上限与重试次数
--unload-grace-seconds默认值卸载前等待 Envoy 收敛的宽限秒数
--log-leveldebug日志级别
--metrics-port/--debug-grpc-port/--drainer-service-port默认端口指标、调试、排空服务端口

当--replica-config未提供时,createReplicaConfig(main.go)会从上述散列参数组装ReplicaConfig,并将InferenceHttpPort/InferenceGrpcPort覆盖为反向代理端口——即 Scheduler 视角中服务器的推理地址是 Agent 的代理端口。

11. 测试与验证

仓库为 Agent API 的各个关键路径提供了完善的单元测试,可作为行为契约参考:

  • scheduler/pkg/agent/server_test.go:覆盖Subscribe注册/断开重调度、AgentEvent状态机映射、Sync加载/卸载下发、AgentDrain排空与ModelScalingTrigger伸缩校验等服务端行为;
  • scheduler/pkg/agent/agent_svc_manager_test.go:验证订阅握手、LoadModel/UnloadModel全流程(含乱序防护、重试、事件上报);
  • scheduler/pkg/agent/rproxy_test.go 与 rproxy_grpc_test.go:验证 HTTP 懒加载重试、路径改写、gRPC 流式转发与元数据透传;
  • scheduler/pkg/agent/state_manager_test.go 与 model_state_test.go:验证内存记账、LRU 驱逐、超卖额度计算与版本管理。

结语

Agent API 是 Seldon Core v2 控制平面与数据平面协同的枢纽协议:Subscribe建立双向通道并完成状态对账,ModelOperationMessage驱动模型的加载与卸载,AgentEvent以事件驱动方式推进模型副本状态机,AgentDrain保障服务器下线时的平滑迁移,ModelScalingTrigger为基于运行时指标的副本伸缩提供数据通道。理解这套协议,是深入 Seldon Core v2 调度、多模型管理(MMS)与数据面路由的关键一步。更完整的系统级说明可继续阅读 docs-gb/apis/internal 下的其他 API 文档与 scheduler/README.md。

  • 人工智能
  • 大模型
  • MLOps
  • LLMOps
  • 模型推理服务
  • 云原生
  • 后端

【免费下载链接】seldon-core

An MLOps framework to package, deploy, monitor and manage thousands of production machine learning models

项目地址:https://gitcode.com/gh_mirrors/se/seldon-core
点击查看免费下载

相关推荐

上一篇:一键找回消失的QQ空间记忆:GetQzonehistory帮你完整备份青春时光
下一篇:Handsontable 仓库 MCP 环境配置实战:ClickUp 任务集成与 code-review-graph 知识图谱

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询