feat(flexlb): support logical multi-engine workers - #1332
Conversation
LLLLKKKK
left a comment
There was a problem hiding this comment.
AI Code Review - PR #1332
Status: LGTM
Summary: P0/0 · P1/0 · P2/16 · P3/7
Reviewed: commit c47bdf062ad3 · 2026-08-25 22:38 UTC+8
lgtm ready to ci
Non-blocking Suggestions
P2
- KVCM host_ip_port 键语义在 proto 契约与 Java 消费方之间自相矛盾 @
rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/match/kvcm/KvcmCacheMatchProvider.java:31- 建议:同步更新
kvcm_meta_service.proto中host_ip_port的注释为逻辑ip:port@engineIndex(并注明单引擎旧格式在 FlexLB 侧有回退兼容),使 proto 成为唯一契约来源;在 PR description 中把「启用multi_engine_num > 1前 KVCM Subscriber 必须先上报逻辑身份」列为显式发布前置条件与灰度顺序。
- 建议:同步更新
- KVCM 键与逻辑 worker 全不匹配时静默按零命中处理,缺可观测手段 @
rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/domain/CacheMatchResult.java:56- 建议:补两个计数(或限频 warn):legacy 物理键回退命中次数;「KVCM 返回非空 hostMatches 但无任何键匹配已知 logical worker」的次数(tag 带
multiEngineNum,日志附一个示例键)。并在 04/06 文档把「启用multi_engine_num > 1前 KVCM 必须上报 logical key,否则命中恒为 0」列为显式运维前置条件,使灰度阶段即可暴露契约不一致。也可考虑把该兼容判定从 domain record 上移到 KVCM provider 或查询编排层,避免CacheMatchResult直接耦合运行态WorkerStatus。
- 建议:补两个计数(或限频 warn):legacy 物理键回退命中次数;「KVCM 返回非空 hostMatches 但无任何键匹配已知 logical worker」的次数(tag 带
- logical 到 physical 映射生命周期不完整,孤儿条目单向泄漏且缺失映射走隐式兜底 @
rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/match/localsync/KvCacheManager.java:164- 建议:把待清理集合改为
getAllEngineIpPorts()与physicalIpPortByLogicalIpPort.keySet()的并集再做物理活跃性过滤,让孤儿映射同轮回收;把getOrDefault兜底改为显式分支,缺失时先 warn 再决定是否判 stale(当前该分支实际不可达,显式化可防后续调用方绕过登记)。也可改由WorkerStatusProvider当前活跃 logical 集合驱动清理以免维护第二份状态,并把该 map size 纳入app.cache.*便于线上确认回收生效。
- 建议:把待清理集合改为
- removeStaleEngineCaches 对空活跃地址列表无保护,服务发现瞬时返回空会清空整份索引 @
rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/match/localsync/KvCacheManager.java:158- 建议:对空活跃列表直接跳过清理并打 warn,或要求连续多轮观测到缺失后才执行清理,使缓存元数据清理与 worker 摘除采用一致的宽限期语义。
- 指标身份契约多处变更(tag key 重命名 + 取值格式 + port 维度丢失)缺少迁移清单 @
rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/telemetry/CacheMetricsReporter.java:333- 建议:在 06 文档或发布说明中补一节指标迁移清单:列出被重命名的 tag key(
engineIpPort→engineIp、ip→engineIp)、取值格式变化涉及的指标族(至少app.cache.*、app.engine.*、app.routing.*)与新增的engineIndextag,并注明回滚该版本时看板需回退旧格式;同时评估同 IP 多 frontend 部署形态,若存在则为丢失 port 维度的两个 cache 指标补独立porttag 以避免标签冲突。
- 建议:在 06 文档或发布说明中补一节指标迁移清单:列出被重命名的 tag key(
- cache-hit PV 的 worker 字段取值由裸 IP 变为逻辑身份,文档却描述为「继续使用」 @
rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/match/localstandby/LocalStandbyComparisonService.java:141- 建议:将该句改为明确的变更声明(旧值为裸 IP,新值为
ip:port@engineIndex,含 N=1),并在 PR 描述的兼容性小节列出该 PV 字段变更;若下游消费方较多,可短期并存一个保留裸 IP 语义的独立字段,给切换留出一个发布周期后再下线。另CacheHitFeedback.workerIp()现已无生产调用方,可一并评估是否随之收敛。
- 建议:将该句改为明确的变更声明(旧值为裸 IP,新值为
- 路由决策 PV 快照身份迁移不完整,同一记录内 logical 与 physical 混用 @
rtp_llm/flexlb/flexlb-common/src/main/java/org/flexlb/dao/pv/ShortestTtftDecision.java:47- 建议:给
WorkerDecision增加engineIndex(或直接改用 logical 身份字段),使workers[]行可唯一定位逻辑引擎并能与cacheAffinityDecision的两个字段字符串对齐;同时在 06 文档/PR 描述中列出cacheLeaderIpPort/shortestTtftWorkerIpPort的取值格式变更。若短期不改字段,至少在文档说明多引擎下应改用cacheLeader/shortestTtftWorker布尔位而非字符串关联。
- 建议:给
- 已有不可变 WorkerIdentity 却仍以相邻同类型 String 传递三种身份表示 @
rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/match/localsync/KvCacheManager.java:105- 建议:让
updateEngineCache/calculateDiff直接接收WorkerIdentity(或最小只读身份视图),在方法内部取getLogicalIpPort()/getPhysicalIpPort()/getIpIndex(),把三种表示的一致性交给类型系统而非形参顺序;CacheHitComparisonResult同理可只持有 identity 并派生两个字段。若倾向最小改动,可让KvCacheManager拿到DiffResult后自行上报 diff 指标,calculateDiff保持仅依赖存储 key 的原签名。
- 建议:让
- multi_engine_num 缺少显式上限,配置笔误会放大同步扇出并使整机永久不可路由 @
rtp_llm/flexlb/flexlb-common/src/main/java/org/flexlb/config/ModelServiceConfiguration.java:64- 建议:在
multiEngineNum < 1判断处一并检查显式上界常量(取值对齐现实部署上限),错误信息同时给出配置值与上限;补一条超上限被拒绝的单测,并在 06 文档字段说明同步该上限,使误配在启动期失败而非退化为运行期静默摘流。
- 建议:在
- 启动日志未输出多引擎派生信息,误配后缺少可对照的排查入口 @
rtp_llm/flexlb/flexlb-common/src/main/java/org/flexlb/config/ModelServiceConfiguration.java:36- 建议:在
modelMetaConfig的 info 日志中补充每个 endpoint 的address、multi_engine_num与派生端口段(或展开后的逻辑 worker 总数);该日志启动期仅执行一次,无热路径噪声风险。可一并考虑在selectRoutableModelWorkerStatus过滤掉不完整物理组时输出限频 warn 或计数,使 sibling 不齐导致的整机摘流可观测。
- 建议:在
- 端口段上界「接受」用例依赖字符串 replace,模板变更时会静默退化为永真断言 @
rtp_llm/flexlb/flexlb-common/src/test/java/org/flexlb/config/ModelServiceConfigurationTest.java:85- 建议:在
run回调内追加对生效配置的断言(取出 endpoint 校验getWorkerStatusPort()为 65534、getMultiEngineNum()为 2);或参照同文件modelConfig(...)的"""...""".formatted(...)写法用带占位符的模板生成 endpoint JSON,不再依赖 replace 命中。
- 建议:在
- CacheMetricsReporterTest 只锁定一个上报方法,其余 engineIp tag 语义变更无覆盖 @
rtp_llm/flexlb/flexlb-cache/src/test/java/org/flexlb/cache/telemetry/CacheMetricsReporterTest.java:26- 建议:为上述四个方法各补一条 tag 断言(复用
FlexMetricTags.of(...)+eq(...)),并补ipIndex == null时不上报的用例,使「每引擎一条指标」的可观测契约可回归。
- 建议:为上述四个方法各补一条 tag 断言(复用
- WorkerStatusTest 未覆盖 engineIndex 传播到 CacheHitFeedback 的路径 @
rtp_llm/flexlb/flexlb-common/src/test/java/org/flexlb/dao/master/WorkerStatusTest.java:34- 建议:在产出 feedback 的用例中设置
setEngineIndex(1)(或新增一条),断言cacheHitFeedbacks().getFirst().logicalWorkerId()为127.0.0.1:8080@1、ipIndex()为127.0.0.1@1,把「engineIndex 端到端进入 feedback」这一契约钉死。
- 建议:在产出 feedback 的用例中设置
- KvCacheManagerTest 缺少同物理机多逻辑引擎的失效清理覆盖 @
rtp_llm/flexlb/flexlb-cache/src/test/java/org/flexlb/cache/match/localsync/KvCacheManagerTest.java:51- 建议:补一条用例:
getAllEngineIpPorts()返回10.0.0.1:8080@0与10.0.0.1:8080@1并分别登记映射,断言物理机仍在活跃列表时两者都不清理、消失时两者都清理;另补一条「logical key 在 view 中但映射缺失且物理机存活」的用例,明确该兜底分支的预期行为。
- 建议:补一条用例:
- 缺少「预测集合中不存在该 logical worker」的静默降级负例 @
rtp_llm/flexlb/flexlb-cache/src/test/java/org/flexlb/cache/match/localstandby/LocalStandbyComparisonServiceTest.java:110- 建议:补一条用例:预测仅含
10.0.0.1:8080@0而 feedback 的 engineIndex 为 1,断言localStandby().hit()为 0 且delta()等于actualHitTokens,把该降级语义显式固化,后续 key 格式漂移即可被区分出来。
- 建议:补一条用例:预测仅含
- CacheMatchResultTest 缺少空 worker 与退化输入的边界覆盖 @
rtp_llm/flexlb/flexlb-cache/src/test/java/org/flexlb/cache/domain/CacheMatchResultTest.java:39- 建议:补
hostMatch((WorkerStatus) null)返回 null 的用例,以及matchedTokens(0, 4096, 100)、matchedTokens(2, 0, 100)断言为 0,使 legacy 回退与退化输入的适用范围在测试中无歧义。可另加一条engineIndex>0 且 multiEngineNum==1的防御性用例(该组合由normalizeHosts:109保证不可达,仅用于固化预期)。
- 建议:补
P3
- CacheMatchResult 两个 hostMatch 重载语义分歧,易被后续调用方误用为严格版本 @
rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/domain/CacheMatchResult.java:45- 建议:给两个方法取语义化的不同名字(如
exactHostMatch(String)与resolveHostMatch(WorkerStatus)),或把hostMatch(String)收窄为包级可见,并在 javadoc 顶部显式标注「不含 KVCM legacy 回退,路由路径请使用 WorkerStatus 重载」。
- 建议:给两个方法取语义化的不同名字(如
- WorkerCacheUpdateResult 的 javadoc 交叉引用指向已不存在的 WorkerStatus 字段 @
rtp_llm/flexlb/flexlb-cache/src/main/java/org/flexlb/cache/domain/WorkerCacheUpdateResult.java:23- 建议:改为引用
org.flexlb.dao.master.WorkerIdentity#logicalIpPort(或WorkerStatus#getLogicalIpPort()),与本 PR 确立的「身份统一由 WorkerIdentity 提供」约定保持一致。
- 建议:改为引用
- LOCAL_STANDBY 测试夹具仍用物理 key,与生产 key 契约脱节 @
rtp_llm/flexlb/flexlb-cache/src/test/java/org/flexlb/cache/match/localstandby/LocalStandbyCacheMatchProviderTest.java:48- 建议:将该用例的 mock key 与断言统一改为逻辑格式(如
10.0.0.1:8080@0),物理 key 用例仅保留在 KVCM 兼容场景;并可把LocalSyncCacheMatchProviderTest各用例重复的 mock/provider 构造抽到@BeforeEach,减少后续身份格式变更的漏改面。
- 建议:将该用例的 mock key 与断言统一改为逻辑格式(如
- 新增 engine-index 用例的期望值与 feedback 字段重合,判别力弱于命名声明 @
rtp_llm/flexlb/flexlb-cache/src/test/java/org/flexlb/cache/match/localstandby/LocalStandbyComparisonServiceTest.java:133- 建议:把 feedback 的
predictedHitTokens/deltaHitTokens改成与预测推导值不同的数值(例如 4096/4904),使断言只能由「按索引查预测」这一条路径满足。
- 建议:把 feedback 的
- WorkerIdentityTest 混入 WorkerHost 断言,且未覆盖多引擎新字段 @
rtp_llm/flexlb/flexlb-common/src/test/java/org/flexlb/dao/master/WorkerIdentityTest.java:32- 建议:将该用例移入新建的
WorkerHostTest并补getWorkerStatusPort()/getMultiEngineNum()断言,WorkerIdentityTest仅保留WorkerIdentity自身用例;如 equals/hashCode 用于集合去重或参数匹配,另补一条等值用例。
- 建议:将该用例移入新建的
- 端口上界常量与既有字面量重复,端口范围校验错误语义不统一 @
rtp_llm/flexlb/flexlb-common/src/main/java/org/flexlb/config/ModelServiceConfiguration.java:23- 建议:将 TCP 端口上界提升为共享常量(本 PR 已在
CommonConstants增补LOGICAL_WORKER_ENGINE_INDEX_SEPARATOR,可顺带放入,或提供isValidTcpPort小工具)供两处共用;并考虑在validateOptimizer中对端口补一致的范围校验,让配置面统一走启动期 fail-fast。
- 建议:将 TCP 端口上界提升为共享常量(本 PR 已在
- address 为空时端点级校验的错误信息退化为 null @
rtp_llm/flexlb/flexlb-common/src/main/java/org/flexlb/config/ModelServiceConfiguration.java:65- 建议:在
validateEngineEndpointConfiguration开头先校验endpoint.getAddress()非空,或在拼接消息时用 group/role 等可辨识信息兜底,保证任何端点级错误信息都能指向唯一 endpoint。
- 建议:在
Checklist Findings (15 fail / 26 total)
General Principles Checklist
- [6.1] Architecture — 兼容性:外部 HTTP/RPC API、持久数据、配置、环境迁移安全 → issue
路由决策 PV 快照身份迁移不完整,同一记录内 logical 与 physical 混用
CacheAffinityDecision的 javadoc(:34-35)已声明两个 worker 字段使用逻辑ip:port@engineIndex,CacheAffinityFirstStrategy:108-109/139-140也确实改用getLogicalIpPort()。但同一 record 内WorkerDecision(:47-50)仍是ip+port两个字段、无 engineIndex(ShortestTTFTStrategy:604-607填worker.getIp()/getPort())。该记录经PvLogData.shortestTtftDecisions落 PV 日志:N>1 时workers[]会出现多行 ip/port 完全相同、无法区分引擎的候选;且原先按ip+":"+port关联cacheLeaderIpPort的离线分析会因新增@index后缀失配。 - [6.1] Architecture — 可观测性:日志/指标/超时可操作、非噪声 → issue
CacheMetricsReporterTest 只锁定一个上报方法,其余 engineIp tag 语义变更无覆盖
该文件是本 PR 为CacheMetricsReporter新增的唯一测试,全文只有一条reportsCacheUpdateLatencyWithIndexedEngineIp(:26-33),断言reportUpdateEngineBlockCacheRT的 tag(key 由engineIpPort改为engineIp、值为10.0.0.8@1)。同一提交把reportEngineLocalMetrics:127、reportCacheHitMetrics:158、reportKvcmSelectedMatch:171、reportCacheDiffMetrics:353的engineIp取值同样改为ipIndex,这四个方法以及ipIndex == null的早退分支(:120、:348)均无任何断言。多引擎下若某条链路仍传入无索引 IP,指标会把兄弟引擎聚合成一条曲线而测试不失败。 - [6.1] Architecture — 回滚路径:风险行为存在运维回滚手段 → issue
指标身份契约多处变更(tag key 重命名 + 取值格式 + port 维度丢失)缺少迁移清单
本 PR 同时做了三类指标契约变更:一是 tag key 重命名,CACHE_UPDATE_ENGINE_BLOCK_CACHE_RT由engineIpPort改为engineIp(:333),ExpirationCleaner的TASK_REMOVED由ip改为engineIp;二是约 20 个指标的engineIp取值由ip/ip:port改为ip@engineIndex(reportEngineLocalMetrics:127、reportCacheHitMetrics:158、reportKvcmSelectedMatch:171、reportCacheDiffMetrics:353、EngineHealthReporter、ResourceMonitorReporter等),N=1 也带@0;三是 select detail 新增engineIndextag。按旧 key/取值精确匹配或按:拆端口的看板与告警会静默断流。`06-configuratio - [6.1] Architecture — 状态不变量:创建/更新/失败/重试/回滚路径有效 → issue
multi_engine_num 缺少显式上限,配置笔误会放大同步扇出并使整机永久不可路由
校验只拒绝multiEngineNum < 1(:64),唯一隐式天花板是端口段检查(:84),故"worker_status_port":1,"multi_engine_num":60000可通过。该值随后线性放大三处资源:RoutingServiceDiscovery:101的discoveredHosts.size() * multiEngineNum、每个逻辑 worker 一条GrpcWorkerStatusRunner/GrpcCacheStatusCheckRunner连接、EngineWorkerStatus:91的new boolean[expected]。更严重的是EngineWorkerStatus.isHealthyPhysicalGroup:82-96要求 siblings 数恰好等于 N,20 误写成 200 时该物理机组永远不健康、被整体摘除路由,而启动期不报错。 - [6.1] Architecture — 错误语义:fail-fast/retry/fallback/silent 行为显式 → issue
address 为空时端点级校验的错误信息退化为 null
validateServiceRoute的循环先调用validateEngineEndpointConfiguration(endpoint)(:54)再调用serviceDiscovery.validate(endpoint)(:55),而 address 非空校验位于RoutingServiceDiscovery:70-71。因此当某 endpoint 同时缺失address且multi_engine_num非法时,抛出的是... multi_engine_num must be greater than zero: null(:65-67),运维在多 endpoint 配置里无法据此定位是哪个 endpoint 出错;worker_status_port的两条错误信息(:71-74、:85-88)同样以 address 结尾,存在相同问题。 - [6.1] Quality — PR description 说明动机与设计 → issue
cache-hit PV 的 worker 字段取值由裸 IP 变为逻辑身份,文档却描述为「继续使用」
result()构造CacheHitComparisonResult的worker实参由feedback.workerIp()(裸 IP,如10.0.0.1)改为feedback.logicalWorkerId()(10.0.0.1:8080@0,N=1 也带@0)。该 JSON 由GrpcWorkerStatusRunner直接pvLogger.info(json)落盘,是命中率离线分析的 worker 维度,所有存量部署该字段取值都会变化。而06-configuration-and-observability.md:169-170写的是「路由、KVCM、cache key 和 cache-hit PV 的 worker 字段继续使用完整ip:port@engineIndex」,把一次取值格式变更描述成保持不变,读者无法据此判断需要改造下游解析,按 IP 精确匹配或聚合该字段的报表/离线任务会在发布后静默失配。 - [6.1] Software Engineering — DRY:重复非平凡逻辑被抽取或显式复用 → issue
端口上界常量与既有字面量重复,端口范围校验错误语义不统一
新增私有常量MAX_TCP_PORT = 65_535(:23),而flexlb-sync/.../optimizer/OptimizerAddressResolver.java:105-111对同一领域概念仍用裸字面量port <= 0 || port > 65535做等价判定,并以log.warn静默跳过。同一约束在两个模块各自编码,后续调整(如收窄到非特权端口)需多点同步;且同为MODEL_SERVICE_CONFIG字段,worker_status_port越界会启动失败,optimizer 端口越界只在运行期告警丢弃,错误语义不一致。 - [6.1] Software Engineering — ISP:调用方不依赖无关大接口 → issue
已有不可变 WorkerIdentity 却仍以相邻同类型 String 传递三种身份表示
updateEngineCache(engineIPort, physicalIpPort, ipIndex, role, newCacheBlocks)(:105-110)把已聚合的身份拆回 4 个相邻 String;EngineLocalView.calculateDiff(engineIPort, ipIndex, newCacheBlocks, role)(:57-58)同样有两个相邻同类型 String,其中ipIndex不参与任何存储或差分、仅透传给reportCacheDiffMetrics:85-86,而调用方本已持有该值并直接上报reportEngineLocalMetrics。CacheHitComparisonResult亦把worker与ipIndex作为相邻 record 组件。任意两处调换实参都能通过编译,后果是指标 tag 维度错位、或physicalIpPortByLogicalIpPort记下错误映射进而误删活跃引擎缓存,只能靠单测断言字面量发现。 - [6.1] Software Engineering — KISS/YAGNI:无投机性抽象 → issue
已有不可变 WorkerIdentity 却仍以相邻同类型 String 传递三种身份表示
updateEngineCache(engineIPort, physicalIpPort, ipIndex, role, newCacheBlocks)(:105-110)把已聚合的身份拆回 4 个相邻 String;EngineLocalView.calculateDiff(engineIPort, ipIndex, newCacheBlocks, role)(:57-58)同样有两个相邻同类型 String,其中ipIndex不参与任何存储或差分、仅透传给reportCacheDiffMetrics:85-86,而调用方本已持有该值并直接上报reportEngineLocalMetrics。CacheHitComparisonResult亦把worker与ipIndex作为相邻 record 组件。任意两处调换实参都能通过编译,后果是指标 tag 维度错位、或physicalIpPortByLogicalIpPort记下错误映射进而误删活跃引擎缓存,只能靠单测断言字面量发现。 - [6.1] Software Engineering — SRP:模块/类职责单一 → issue
WorkerIdentityTest 混入 WorkerHost 断言,且未覆盖多引擎新字段
workerHostExposesItsStoredIdentityRepresentations(:32-43)实际验证的是WorkerHost的 10 参构造与身份 getter 委托,与文件名声明的被测类不一致;后续维护者按类名检索WorkerHost的测试会误判其无覆盖。该用例只断言身份三元组与 engineIndex,未断言构造入参中的workerStatusPort=18003与multiEngineNum=2——这两个值正是多引擎端口派生与 sibling 完整性判定(isHealthyPhysicalGroup:87-96)的输入,位置参数错位不会被发现;WorkerIdentity由 Lombok 生成的 equals/hashCode(WorkerIdentity.java:22,被CacheHitFeedback相等性间接依赖)亦无覆盖。 - [6.1] Tests — 分布式/跨平台变更有对应覆盖 → issue
KvCacheManagerTest 缺少同物理机多逻辑引擎的失效清理覆盖
新增的keepsLogicalCacheWhenItsPhysicalEngineRemainsDiscoverable(:51)与removesLogicalCacheWhenItsPhysicalEngineDisappears(:68)都只注册一个逻辑引擎10.0.0.1:8080@0。而多引擎正是本次改造目标:removeStaleEngineCaches:158-166入参是物理地址、内部靠physicalIpPortByLogicalIpPort反查,同一物理机下@0/@1两条映射必须同时保留或同时清理。当前用例无法发现「只清理了一个兄弟索引」或「映射未按 logical key 建立导致误删存活缓存」;getOrDefault兜底的「映射缺失且物理机存活」方向亦未覆盖。 - [6.1] Tests — 新逻辑有聚焦单测 + 相关集成/smoke 测试 → issue
新增 engine-index 用例的期望值与 feedback 字段重合,判别力弱于命名声明
feedback 传入的predictedHitTokens=8192、deltaHitTokens=808(:133-135),与由@1预测(local(2)× blockSize 4096,inputTokens 12000,actual 9000)推导出的 localStandby hit/delta 完全相同。因此assertEquals(8192, result.localStandby().hit())(:141)虽能排除「错误回退到 index 0」(@0会得到 4096/4904),却无法区分「按 engineIndex 命中预测」与「直接复制 routing 预测字段」两种实现。 - [6.1] Tests — 边界 case 覆盖(空、单元素、最大值) → issue
WorkerIdentityTest 混入 WorkerHost 断言,且未覆盖多引擎新字段
workerHostExposesItsStoredIdentityRepresentations(:32-43)实际验证的是WorkerHost的 10 参构造与身份 getter 委托,与文件名声明的被测类不一致;后续维护者按类名检索WorkerHost的测试会误判其无覆盖。该用例只断言身份三元组与 engineIndex,未断言构造入参中的workerStatusPort=18003与multiEngineNum=2——这两个值正是多引擎端口派生与 sibling 完整性判定(isHealthyPhysicalGroup:87-96)的输入,位置参数错位不会被发现;WorkerIdentity由 Lombok 生成的 equals/hashCode(WorkerIdentity.java:22,被CacheHitFeedback相等性间接依赖)亦无覆盖。
RTP-LLM Checklist
- [I] 代码质量 — 删除或重命名内部 file、registry entry、model name、metric enum、op binding、plugin symbol 时,必须全仓搜索消费者,并提供替代实现、迁移说明或 smoke 覆盖;只有暴露到 HTTP/RPC/config/persisted format 时才按外部兼容性处理 → issue
WorkerCacheUpdateResult 的 javadoc 交叉引用指向已不存在的 WorkerStatus 字段
新增注释在 :23 引用org.flexlb.dao.master.WorkerStatus#engineIndex,但本 PR 已把engineIndex从WorkerStatus移入WorkerIdentity:全仓检索显示该字段仅声明于WorkerIdentity.java:31,WorkerStatus只剩委托给workerIdentity的getEngineIndex():112/setEngineIndex():116。该引用为悬空链接(当前 pom 未启用 maven-javadoc-plugin,故不影响构建,仅误导读者并在开启 javadoc 校验时告警)。 - [I] 代码质量 — 同一功能用统一工具函数 → issue
端口上界常量与既有字面量重复,端口范围校验错误语义不统一
新增私有常量MAX_TCP_PORT = 65_535(:23),而flexlb-sync/.../optimizer/OptimizerAddressResolver.java:105-111对同一领域概念仍用裸字面量port <= 0 || port > 65535做等价判定,并以log.warn静默跳过。同一约束在两个模块各自编码,后续调整(如收窄到非特权端口)需多点同步;且同为MODEL_SERVICE_CONFIG字段,worker_status_port越界会启动失败,optimizer 端口越界只在运行期告警丢弃,错误语义不一致。
Strengths
WorkerIdentity作为不可变值对象一次性预计算 physical / logical / ipIndex 三种表示(WorkerIdentity.java:39-50),WorkerStatus的setIp/setPort/setEngineIndex整份原子替换 volatile 引用(WorkerStatus.java:35,94-120),从根上消除「三种派生表示来自不同快照」的竞态,优于在各调用点分散拼接字符串。- KVCM 兼容策略收敛克制、边界清晰:
hostMatch(WorkerStatus)仅在「logical 未命中 +source == KVCM+multiEngineNum == 1」三条件同时成立时回退 physical key(CacheMatchResult.java:56-67),N>1 与非 KVCM source 一律 exact match,避免把整台 frontend 的命中错误归属到某个 index。 - 身份切换无半迁移:LOCAL_SYNC 两级索引、LOCAL_STANDBY 候选集/rollback/routed-request 记账、cache-hit 对比查表、worker 状态表、路由回滚全部切到 logical key,写入侧与读取侧同源;全仓已无
getEngineIpPort残留消费者。 - 多引擎机制自洽:
normalizeHosts按workerStatusBasePort + engineIndex展开逻辑 worker,两个 gRPC runner 各自连自己的workerStatusPort,每个逻辑引擎的 cache/worker 状态天然独立,不会出现 N 个逻辑 worker 抓同一份数据。 - N=1 语义保守:内部保留
@0身份、仅省略 wireengine_index(ServerStatus.java:36-47),rollback 与 cache 归属在单引擎下同样精确,没有留下「索引可空」分支;selectRoutableModelWorkerStatus在 N=1 时退化为原有 alive 过滤,存量部署行为不变。 - 配置校验严谨:
(long) workerStatusPort + multiEngineNum - 1 > MAX_TCP_PORT(ModelServiceConfiguration.java:84)显式上转 long 规避 int 溢出绕过;N=1 提前 return 不强制基准端口,与normalizeHosts:106-108的 grpcPort 回退闭合,存量单引擎配置零改动可启动;四类错误信息各自独立并附带 endpoint address。 - 兄弟引擎串号这一最核心风险有正面覆盖:
isolatesMatchesBetweenSiblingEngineIndexes、appliesCacheMatchRollbackOnlyToMatchingSiblingIndex、matchesLocalStandbyPredictionByEngineIndex构造同 IP 同端口、仅 engineIndex 不同的两个 worker,实现若退回物理 key 会立即失败;ModelServiceConfigurationTest成对覆盖了端口段上界两侧(65534+2-1通过、65535+2-1拒绝)。 - 身份语义在 javadoc 中逐参数标注(
CacheMatchProvider、EngineLocalView.java:51-55、KvCacheManager.java:97-103、GlobalCacheIndex、KvcmCacheMatchProvider),并同步更新 5 篇架构文档的三种身份表示对照表,后续维护者不需要靠猜字符串格式。
| } | ||
|
|
||
| /** | ||
| * Returns KVCM matches without rewriting their keys. KVCM {@code host_ip_port} values follow |
There was a problem hiding this comment.
[P2] KVCM host_ip_port 键语义在 proto 契约与 Java 消费方之间自相矛盾
本 PR 把 KVCM 返回键的语义改为逻辑身份:KvcmCacheMatchProvider.java:31-32 的 javadoc 写 "KVCM host_ip_port values follow the logical ip:port@engineIndex identity reported by the Subscriber",04-worker-sync-and-cache.md:103 同步声明。但键的生产方契约定义 flexlb-grpc/src/main/proto/kvcm_meta_service.proto:80 未改,仍写 host_ip_port "Must match WorkerStatus#getIpPort()",而 WorkerStatus.java:588-589 明确返回物理 ip:port。KvcmGrpcClient 原样透传该键作 map key,因此 N>1 时 KVCM 侧若按 proto 注释实现即全量 miss。跨服务契约在同一仓库内给出两个互斥定义。
建议: 同步更新 kvcm_meta_service.proto 中 host_ip_port 的注释为逻辑 ip:port@engineIndex(并注明单引擎旧格式在 FlexLB 侧有回退兼容),使 proto 成为唯一契约来源;在 PR description 中把「启用 multi_engine_num > 1 前 KVCM Subscriber 必须先上报逻辑身份」列为显式发布前置条件与灰度顺序。
| * worker has an unambiguous logical {@code @0} identity, so it can use that legacy entry | ||
| * when no exact logical match exists. Multi-engine workers require an exact logical match. | ||
| */ | ||
| public HostCacheMatch hostMatch(WorkerStatus workerStatus) { |
There was a problem hiding this comment.
[P2] KVCM 键与逻辑 worker 全不匹配时静默按零命中处理,缺可观测手段
物理回退仅在 source == KVCM && multiEngineNum == 1 生效(:61-66)。若 KVCM 仍上报物理键而运维先把 multi_engine_num 调到 2,则 logical 全 miss 且回退条件不成立,ShortestTTFTStrategy:234/306/889 与 WeightedCacheLoadBalancer:144 拿到 null,整组 worker 的 cache 亲和归零、TTFT 退化为全量 prefill。该路径既无日志也无指标:MetricConstant 仅有 app.cache.match.standby.fallback.qps(语义是 KVCM 不可用)与 app.cache.kvcm.query.retry.qps,既无法区分「确实无命中」与「键格式不匹配」,也没有 legacy 回退命中计数以判断迁移是否完成、兼容分支何时可删。
建议: 补两个计数(或限频 warn):legacy 物理键回退命中次数;「KVCM 返回非空 hostMatches 但无任何键匹配已知 logical worker」的次数(tag 带 multiEngineNum,日志附一个示例键)。并在 04/06 文档把「启用 multi_engine_num > 1 前 KVCM 必须上报 logical key,否则命中恒为 0」列为显式运维前置条件,使灰度阶段即可暴露契约不一致。也可考虑把该兼容判定从 domain record 上移到 KVCM provider 或查询编排层,避免 CacheMatchResult 直接耦合运行态 WorkerStatus。
| Set<String> activePhysicalIpPorts = new HashSet<>(activeEngineIpPorts); | ||
| Set<String> staleEngineIpPorts = new HashSet<>(engineLocalView.getAllEngineIpPorts()); | ||
| staleEngineIpPorts.removeAll(new HashSet<>(activeEngineIpPorts)); | ||
| staleEngineIpPorts.removeIf(engineIpPort -> |
There was a problem hiding this comment.
[P2] logical 到 physical 映射生命周期不完整,孤儿条目单向泄漏且缺失映射走隐式兜底
physicalIpPortByLogicalIpPort 仅在 updateEngineCache:115-117 写入,在 stale 清理命中(:171)与 clear() 时删除;待清理集合只来自 engineLocalView.getAllEngineIpPorts():163。而 EngineLocalView.removeCacheBlock:138-140 在某引擎 block 集合清空后会 engineViews.remove(engineIPort),该 logical key 此后不再出现于该集合,映射条目永久残留;pod 轮转/扩缩容下条目数按历史 ip:port@index 去重集合单向增长且无指标暴露(功能不受影响,下轮 update 会重新登记)。同时 getOrDefault(engineIpPort, engineIpPort):166 把「映射缺失」隐式当物理地址处理,落入该分支即判 stale 并清空该引擎全部缓存元数据,无日志、无显式分支。
建议: 把待清理集合改为 getAllEngineIpPorts() 与 physicalIpPortByLogicalIpPort.keySet() 的并集再做物理活跃性过滤,让孤儿映射同轮回收;把 getOrDefault 兜底改为显式分支,缺失时先 warn 再决定是否判 stale(当前该分支实际不可达,显式化可防后续调用方绕过登记)。也可改由 WorkerStatusProvider 当前活跃 logical 集合驱动清理以免维护第二份状态,并把该 map size 纳入 app.cache.* 便于线上确认回收生效。
| * | ||
| * @param activeEngineIpPorts active physical engine addresses in {@code ip:port} format | ||
| */ | ||
| public void removeStaleEngineCaches(Collection<String> activeEngineIpPorts) { |
There was a problem hiding this comment.
[P2] removeStaleEngineCaches 对空活跃地址列表无保护,服务发现瞬时返回空会清空整份索引
EngineAddressResolver.updateEndpointHosts:96-97 在某 endpoint 的 hostList 为 null/空时 domainHostsMap.remove(endpoint),随后以聚合结果通知 listener;LocalCacheEngineAddressListener.onAddressUpdate:34-38 只拦 null、不拦空列表。因此聚合为空时全部 logical key 都会被判 stale 并清空两级索引,命中率骤降且需一整轮同步才能重建。对比 worker 状态摘除路径存在宽限期(EngineSyncRunner 按实际同步间隔的倍数延迟摘除),本清理路径没有等价保护。该行为在本次改动前已存在,但本 PR 正在重写该方法的 staleness 判定语义,适合一并收敛。
建议: 对空活跃列表直接跳过清理并打 warn,或要求连续多轮观测到缺失后才执行清理,使缓存元数据清理与 worker 摘除采用一致的宽限期语义。
| FlexMetricTags tags = FlexMetricTags.of("engineIpPort", engineIpPort, "role", role, "success", success); | ||
| public void reportUpdateEngineBlockCacheRT( | ||
| String ipIndex, String role, long startTime, String success) { | ||
| FlexMetricTags tags = FlexMetricTags.of("engineIp", ipIndex, "role", role, "success", success); |
There was a problem hiding this comment.
[P2] 指标身份契约多处变更(tag key 重命名 + 取值格式 + port 维度丢失)缺少迁移清单
本 PR 同时做了三类指标契约变更:一是 tag key 重命名,CACHE_UPDATE_ENGINE_BLOCK_CACHE_RT 由 engineIpPort 改为 engineIp(:333),ExpirationCleaner 的 TASK_REMOVED 由 ip 改为 engineIp;二是约 20 个指标的 engineIp 取值由 ip / ip:port 改为 ip@engineIndex(reportEngineLocalMetrics:127、reportCacheHitMetrics:158、reportKvcmSelectedMatch:171、reportCacheDiffMetrics:353、EngineHealthReporter、ResourceMonitorReporter 等),N=1 也带 @0;三是 select detail 新增 engineIndex tag。按旧 key/取值精确匹配或按 : 拆端口的看板与告警会静默断流。`06-configura...
建议: 在 06 文档或发布说明中补一节指标迁移清单:列出被重命名的 tag key(engineIpPort → engineIp、ip → engineIp)、取值格式变化涉及的指标族(至少 app.cache.*、app.engine.*、app.routing.*)与新增的 engineIndex tag,并注明回滚该版本时看板需回退旧格式;同时评估同 IP 多 frontend 部署形态,若存在则为丢失 port 维度的两个 cache 指标补独立 port tag 以避免标签冲突。
Checklist: [6.1] 回滚路径:风险行为存在运维回滚手段
| @@ -60,8 +62,10 @@ void waitsForStandbyHashBeforeMatching() throws Exception { | |||
| } | |||
There was a problem hiding this comment.
📍 实际位置 rtp_llm/flexlb/flexlb-cache/src/test/java/org/flexlb/cache/match/localstandby/LocalStandbyCacheMatchProviderTest.java:48(不在 diff 展示范围内,就近挂载)
[P3] LOCAL_STANDBY 测试夹具仍用物理 key,与生产 key 契约脱节
waitsForStandbyHashBeforeMatching 让 cacheManager.findMatchingEngines 返回 Map.of("10.0.0.1:8080", 1)(:48),并按 result.hostMatch("10.0.0.1:8080")(:57)断言。但 LocalStandbyCacheManager.findMatchingEngines 的 javadoc 与实现已明确 key 为 ip:port@engineIndex,且 hostMatch(WorkerStatus) 对非 KVCM 来源不再回退物理 key,因此该夹具与生产数据形态不符。同一 PR 内 LocalStandbyComparisonServiceTest 与同文件 updatesRequestDerivedCacheMetadataAsynchronously(:66)已统一使用逻辑 key,此处遗漏容易被后续用例照抄。
建议: 将该用例的 mock key 与断言统一改为逻辑格式(如 10.0.0.1:8080@0),物理 key 用例仅保留在 KVCM 兼容场景;并可把 LocalSyncCacheMatchProviderTest 各用例重复的 mock/provider 构造抽到 @BeforeEach,减少后续身份格式变更的漏改面。
|
|
||
| CacheHitFeedback feedback = new CacheHitFeedback( | ||
| "cache_hit_comparison", "request-index-1", "LOCAL_STANDBY", "PREFILL", "default", | ||
| "10.0.0.1", 8080, 1, "running", 12000, 4096, 8192, |
There was a problem hiding this comment.
[P3] 新增 engine-index 用例的期望值与 feedback 字段重合,判别力弱于命名声明
feedback 传入的 predictedHitTokens=8192、deltaHitTokens=808(:133-135),与由 @1 预测(local(2) × blockSize 4096,inputTokens 12000,actual 9000)推导出的 localStandby hit/delta 完全相同。因此 assertEquals(8192, result.localStandby().hit())(:141)虽能排除「错误回退到 index 0」(@0 会得到 4096/4904),却无法区分「按 engineIndex 命中预测」与「直接复制 routing 预测字段」两种实现。
建议: 把 feedback 的 predictedHitTokens / deltaHitTokens 改成与预测推导值不同的数值(例如 4096/4904),使断言只能由「按索引查预测」这一条路径满足。
Checklist: [6.1] 新逻辑有聚焦单测 + 相关集成/smoke 测试
| } | ||
|
|
||
| @Test | ||
| void workerHostExposesItsStoredIdentityRepresentations() { |
There was a problem hiding this comment.
[P3] WorkerIdentityTest 混入 WorkerHost 断言,且未覆盖多引擎新字段
workerHostExposesItsStoredIdentityRepresentations(:32-43)实际验证的是 WorkerHost 的 10 参构造与身份 getter 委托,与文件名声明的被测类不一致;后续维护者按类名检索 WorkerHost 的测试会误判其无覆盖。该用例只断言身份三元组与 engineIndex,未断言构造入参中的 workerStatusPort=18003 与 multiEngineNum=2——这两个值正是多引擎端口派生与 sibling 完整性判定(isHealthyPhysicalGroup:87-96)的输入,位置参数错位不会被发现;WorkerIdentity 由 Lombok 生成的 equals/hashCode(WorkerIdentity.java:22,被 CacheHitFeedback 相等性间接依赖)亦无覆盖。
建议: 将该用例移入新建的 WorkerHostTest 并补 getWorkerStatusPort() / getMultiEngineNum() 断言,WorkerIdentityTest 仅保留 WorkerIdentity 自身用例;如 equals/hashCode 用于集合去重或参数匹配,另补一条等值用例。
Checklist: [6.1] SRP:模块/类职责单一;[6.1] 边界 case 覆盖(空、单元素、最大值)
| public class ModelServiceConfiguration { | ||
|
|
||
| private static final String MODEL_SERVICE_CONFIG = "MODEL_SERVICE_CONFIG"; | ||
| private static final int MAX_TCP_PORT = 65_535; |
There was a problem hiding this comment.
[P3] 端口上界常量与既有字面量重复,端口范围校验错误语义不统一
新增私有常量 MAX_TCP_PORT = 65_535(:23),而 flexlb-sync/.../optimizer/OptimizerAddressResolver.java:105-111 对同一领域概念仍用裸字面量 port <= 0 || port > 65535 做等价判定,并以 log.warn 静默跳过。同一约束在两个模块各自编码,后续调整(如收窄到非特权端口)需多点同步;且同为 MODEL_SERVICE_CONFIG 字段,worker_status_port 越界会启动失败,optimizer 端口越界只在运行期告警丢弃,错误语义不一致。
建议: 将 TCP 端口上界提升为共享常量(本 PR 已在 CommonConstants 增补 LOGICAL_WORKER_ENGINE_INDEX_SEPARATOR,可顺带放入,或提供 isValidTcpPort 小工具)供两处共用;并考虑在 validateOptimizer 中对端口补一致的范围校验,让配置面统一走启动期 fail-fast。
Checklist: [6.1] DRY:重复非平凡逻辑被抽取或显式复用;[I] 同一功能用统一工具函数
| private void validateEngineEndpointConfiguration(Endpoint endpoint) { | ||
| int multiEngineNum = endpoint.getMultiEngineNum(); | ||
| if (multiEngineNum < 1) { | ||
| throw new IllegalArgumentException( |
There was a problem hiding this comment.
[P3] address 为空时端点级校验的错误信息退化为 null
validateServiceRoute 的循环先调用 validateEngineEndpointConfiguration(endpoint)(:54)再调用 serviceDiscovery.validate(endpoint)(:55),而 address 非空校验位于 RoutingServiceDiscovery:70-71。因此当某 endpoint 同时缺失 address 且 multi_engine_num 非法时,抛出的是 ... multi_engine_num must be greater than zero: null(:65-67),运维在多 endpoint 配置里无法据此定位是哪个 endpoint 出错;worker_status_port 的两条错误信息(:71-74、:85-88)同样以 address 结尾,存在相同问题。
建议: 在 validateEngineEndpointConfiguration 开头先校验 endpoint.getAddress() 非空,或在拼接消息时用 group/role 等可辨识信息兜底,保证任何端点级错误信息都能指向唯一 endpoint。
Checklist: [6.1] 错误语义:fail-fast/retry/fallback/silent 行为显式
Background
DashLLM can start multiple independently routable engines behind one frontend when
DS_LLM_MULTI_ENGINE_NUM > 1. FlexLB must preserve that engine index from service discovery through worker-status/cache polling, cache matching, scheduling, feedback, rollback, and observability.Implementation
ip:portip:port@engineIndexip@engineIndexmulti_engine_numlogical workers. Worker-status and cache-status gRPC useworker_status_port + engineIndex; for N=1, an omittedworker_status_portfalls back to the engine gRPC port.ip:port; passWorkerHostto runners so LOCAL_SYNC never parses a logical address as a TCP endpoint.engine_indexonly for N>1. N=1 responses omit it, while all internal cache/subscriber/KVCM identities remain@0.multi_engine_num >= 1,worker_status_portis in[1, 65535], and N>1 requires a configured worker-status port whose complete offset range fits in the TCP range.ip:port@0wins; only KVCM results may fall back to legacy physicalip:portfor a single-engine worker. Physical keys remain ignored for N>1 to avoid assigning a multi-engine cache match to the wrong index.Compatibility and rollout
multi_engine_numand a baseworker_status_port; engines expose one worker-control gRPC endpoint per index.worker_status_port; old worker-status responses are represented as index 0.host_ip_portvalues through the N=1 fallback. Multi-engine cache routing requires KVCM/subscriber to returnip:port@engineIndex.@indexis an internal logical identity, not a network endpoint suffix.Verification
./mvnw verify./mvnw -pl flexlb-sync -am -Dtest=LoadBalanceStrategyFactoryTest,DefaultRouterTest -Dsurefire.failIfNoSpecifiedTests=false test(20 tests passed)Explicitly deferred
multi_engine_num; the configured worker-status TCP port range remains the enforced bound.Rollback
Set
multi_engine_numto 1 and restart the deployment. FlexLB resumes one logical worker per frontend and omitsengine_indexfrom schedule responses.