- 人工智能
- 大模型
- MLOps
- LLMOps
- 模型推理服务
- 云原生
- 后端
【免费下载链接】seldon-core
An MLOps framework to package, deploy, monitor and manage thousands of production machine learning models
导读
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 → Scheduler | Agent 上报模型事件(加载成功/失败、卸载成功/失败、内存信息) |
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):
| 字段 | 类型 | 含义 |
|---|---|---|
inferenceSvc | string | 推理服务的 DNS 名称 |
inferenceHttpPort | int32 | 推理 HTTP 端口 |
inferenceGrpcPort | int32 | 推理 gRPC 端口 |
memoryBytes | uint64 | 服务器副本的内存容量 |
capabilities | repeated string | 服务器能力清单,如sklearn、pytorch、xgboost、mlflow |
overCommitPercentage | uint32 | 允许的超卖内存百分比,设为 0(%)表示禁止超卖 |
loadedModels是 Agent 启动时已加载模型的快照,availableMemoryBytes是考虑超卖后的可用内存。这两项用于网络抖动后的状态对账:Scheduler 端在Subscribe处理中会调用scheduleModelsFromRequest(server.go),把 Agent 上报的已加载模型重新纳入调度,并重试此前失败的模型(ScheduleFailedModels),从而避免"网络闪断导致模型被调度到其他服务器,而本机其实还加载着"的重复加载。
服务端Subscribe实现的几个关键细节(server.go):
- 按副本串行化:使用
agentMutex(sync.Map)对同一(serverName, replicaIdx)强制串行——保证旧 Agent 完全断开后新 Agent 才能接入; - 注册副本:将请求存入
store.AddServerReplica,纳入全局调度视图; - 流式长连接:阻塞在
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):
- 乱序防护:基于单调时钟的时间戳记录(
modelTimestamps)忽略乱序到达的过期指令; - 获取存储配置:从
ModelSpec.StorageConfig中解析 rclone 配置或 K8s Secret(getArtifactConfig,支持StorageRcloneConfig与StorageSecretName两种形态); - 下载模型工件:通过
ModelRepository.DownloadModelVersion(基于 rclone)将模型拉取到本地/mnt/agent/models; - 加载到推理服务器:调用模型服务器控制面客户端(
v2Client.LoadModelVersion,支持 MLServer 与 Triton,工厂实现见 modelserver_controlplane/factory),失败时按maxLoadRetryCount/maxLoadElapsedTime退避重试; - 可选挂接伸缩统计:若
AutoscalingEnabled且 Agent 启用了伸缩,则把模型加入StatsAnalyserService; - 上报成功:发送
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) |
|---|---|---|
LOADED | Loading | Loaded |
UNLOADED | Unloading | Unloaded |
LOAD_FAILED/LOAD_FAIL_MEMORY | Loading | LoadFailed |
UNLOAD_FAILED | Unloading | UnloadFailed |
Scheduler 使用乐观期望状态校验,防止状态乱序回跳;同时将availableMemoryBytes与runtimeInfo一并写入 store,供后续调度决策(如内存感知的模型放置、副本扩展)使用。
6. 排空机制:AgentDrain
当服务器副本需要下线(如节点驱逐、滚动更新)时,Agent 调用AgentDrain(AgentDrainRequest{serverName, replicaIdx})请求 Scheduler 排空该副本。
Scheduler 端drainServerReplicaImpl(server.go)的执行序列:
store.DrainServerReplica将副本上所有模型标记为待迁移;- 用
modelRelocatedWaiter为这些模型注册等待组——模型在其他副本上变为Available时被signalModel释放; - 睡眠
agentDrainCoolDownPeriod(500ms)作为冷却期,避免多个 Agent 并发排空时调度器把模型调度到同样在排空的服务器上; - 逐个
scheduler.Schedule(modelName)把模型重新调度到健康副本; - 阻塞等待所有模型迁移完成(
waiter.wait); - 再额外等待
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-name | mlserver | 服务器名称,用于在 Scheduler 中唯一标识 |
--server-idx | 0 | 服务器副本索引(ReplicaIdx) |
--scheduler-host/--scheduler-port/--scheduler-tls-port | 0.0.0.0/ 默认端口 | Scheduler 地址与明文/mTLS 端口 |
--rclone-host/--rclone-port | 0.0.0.0/ 默认端口 | rclone 下载服务地址 |
--inference-host/--inference-http-port/--inference-grpc-port | 0.0.0.0/ 默认端口 | 推理服务器地址与 HTTP/gRPC 端口 |
--reverse-proxy-http-port/--reverse-proxy-grpc-port | 默认 HTTP 端口 /9998 | 数据平面反向代理监听端口 |
--server-type | mlserver | 模型服务器类型(mlserver或triton),决定仓库处理器与控制面客户端工厂 |
--memory-bytes | 1000000 | 服务器可用内存(字节) |
--capabilities | sklearn,xgboost | 服务器能力清单(逗号分隔),如 sklearn、pytorch、xgboost、mlflow |
--overcommit-percentage | 0 | 内存超卖百分比,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-level | debug | 日志级别 |
--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
相关推荐
Seldon Core 2 Chainer API 深度解析:Scheduler 与 Dataflow 引擎之间的管道编排 gRPC 协议
Seldon Core 2 Chainer API 深度解析:Scheduler 与 Dataflow 引擎之间的管道编排 gRPC 协议 Chainer AP
人工智能大模型MLOpsLLMOps模型推理服务云原生后端Khoj 代码执行(Code Execution)能力解析:基于 Terrarium 与 E2B 沙箱的本地自托管指南
Khoj 代码执行(Code Execution)能力解析:基于 Terrarium 与 E2B 沙箱的本地自托管指南 Khoj 的代码执行特性允许 AI 在受
人工智能大模型MLOpsLLMOps模型推理服务云原生后端AutoGen Agent 身份(Agent ID)与生命周期管理深度解析
AutoGen Agent 身份(Agent ID)与生命周期管理深度解析 Agent Runtime(运行时)是整个 AutoGen Core 框架的心脏:它
人工智能AI AgentAgent 框架多智能体大模型工具调用
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考