深入解析 AIBrix statesync:基于 Redis 的跨副本状态同步组件与接入实战 深入解析 AIBrix statesync基于 Redis 的跨副本状态同步组件与接入实战【免费下载链接】aibrixCost-efficient and pluggable Infrastructure components for GenAI inference项目地址: https://gitcode.com/GitHub_Trending/ai/aibrix导读AIBrix 网关插件以多副本方式部署时各副本内存中的本地状态如前缀缓存 PrefixHashTable彼此隔离会造成同一请求命中不同副本、路由决策不一致的问题。pkg/plugins/gateway/statesync提供了一套通用的、基于 Redis 的跨副本状态同步层以每个实体一个 KeySETEX 每记录 TTL、MGET批量拉取、SCAN枚举为技术底座采用周期性的 Pull-first Push同步模型不依赖 PUB/SUB。本文将以 statesync README 为主线结合其底层实现 redissync.go、接口定义 syncable.go 以及网关默认接入代码 cmd/plugins/main.go完整讲解它的设计原理、配置项、接入步骤与实战示例帮助你在自己的组件中快速落地跨副本状态同步。一、statesync 是什么设计目标与核心机制statesync 的目标非常聚焦让多个网关或其他副本之间以 Redis 为共享存储周期性地同步各自持有的本地状态。它被设计为通用组件——statesync包本身不依赖任何组件类型只依赖抽象接口syncable.Syncable具体状态由拥有该状态的包自己实现适配器然后注册进来。其核心机制可以概括为以下几点均可在 redissync.go 的包注释与实现中验证每实体一个 Keyper-entity keys状态以namespace → entityId → serialized bytes的键值形态建模每个实体对应一个独立的 Redis Key使用SETEX写入并携带每记录独立 TTL拉取用MGET批量读取枚举用SCAN。Key 命名空间与哈希标签Key 格式为aibrix:{namespace}:e:{entityId}其中{namespace}是 Redis Cluster 的 hash tag用于保证同一 namespace 的实体 Key 落在同一分片便于批量操作namespace 必须非空。周期性同步每个副本的同步循环采用Pull first先从 Redis 加载再 Push推送增量或全量的顺序每syncPeriod执行一轮。可选的写穿write-through本地状态变更后可以直接调用Put/Delete立即写入 Redis但网关默认接线只使用周期性增量推送不在每次AddPrefix时调用Put。删除与墓碑tombstone如果 Syncable 实现了TombstoneSupport删除操作会以墓碑载荷写入对端在下一次 Pull 时据此删除本地条目否则直接使用DEL。抖动与退避启动时有一个最大为syncPeriod/2的随机延迟打散多副本同时启动的负载每轮周期叠加约 ±10% 的抖动出错时指数退避上限 1 分钟见syncLoop的实现。生命周期约束所有 Syncable 必须在Start()之前注册Start()之后的注册会被静默忽略源码中会打印 warning 并直接返回。每记录 TTL 的语义每个实体独立存储为 Key默认 TTL 2 分钟。一个记录如果在 TTL 窗口内未被重写就会过期消失单个记录的 TTL 不影响其他记录。这意味着陈旧数据有界——最坏情况下一个副本写入的状态会在 TTL 之后从 Redis 中自然消失而不是无限期残留在共享存储中。二、技术流程一轮同步循环到底做了什么README 中给出了完整的时序图它精确对应 syncLoop 与 runOneSyncCycle 的实现把图中的每一步对应到代码Pull 阶段Pull 函数先检查 Syncable 是否实现了OptionalHooks若是则在 Pull 前后调用OnSyncStart/OnSyncEnd注意钩子只包住 Pull不包住 Push。用SCANCOUNT 提示为 500按aibrix:{namespace}:e:*模式枚举该 namespace 下所有实体 Key累积到mgetBatchSize默认 200一批调用MGET批量取值。对每个值依次处理若实现了TombstoneSupport且载荷是墓碑则调用DeleteLocal删除本地条目并跳过若实现了StalePolicy且IsStale返回 true则删除 Redis 中该 Key并跳过应用否则调用ApplyRemote(id, bytes)合入本地状态。Redis 会过滤已过期的 Key因此拉取到的都是 TTL 窗口内的有效记录。Push 阶段Push 函数若 Syncable 实现了DeltaSyncable走pushDelta只推送自上次ClearDirty以来变更的实体updated 用SETEXdeleted 用DEL或墓碑推送成功后调用ClearDirty清空脏标记。否则走pushFull调用GetSnapshot拿到全量快照全部用SETEX写入快照超过 1 万字段或平均每字段超过 512KB 时会打印告警。写入使用 pipeline 分块setexChunkSize默认 25降低往返开销。循环调度syncLoop启动先随机延迟最多syncPeriod/2→ 立即执行一轮Pull-first Push→ 之后按syncPeriod ± 10%抖动间隔循环任何一轮出错下一轮间隔按 2 倍指数退避封顶 1 分钟成功后恢复为syncPeriod。为什么要 Pull-first 再 Push因为先拉取对端最新状态、再推送本地状态可以减少本地方案覆盖对端更新的概率见 runOneSyncCycle 注释。三、网关默认接线前缀缓存Prefix Cache如何同步前缀缓存是 statesync 在网关插件中的第一个落地场景。启用开关为环境变量AIBRIX_STATESYNC_ENABLED默认false接线逻辑位于 cmd/plugins/main.gostateSyncEnabled : utils.LoadEnvBool(AIBRIX_STATESYNC_ENABLED, false) var syncManager *statesync.RedisSync if stateSyncEnabled { klog.InfoS(statesync enabled; starting cross-replica state sync) table : prefixcacheindexer.GetSharedPrefixHashTable() table.EnableDeltaSync() syncManager statesync.New(redisClient) syncManager.Register(prefixcacheindexer.NewPrefixHashTableSyncable(table)) syncManager.Start() } else { klog.InfoS(statesync disabled; set AIBRIX_STATESYNC_ENABLEDtrue to enable cross-replica state sync) }启用后的三步与 README 描述完全一致调用table.EnableDeltaSync()激活PrefixHashTable的脏标记追踪。该方法的实现位于 hash.go它只是初始化dirtyIds集合在AddPrefix写入新块时会顺手把块哈希记入dirtyIds见 hash.go。未启用时该追踪是零成本的 no-op。通过NewPrefixHashTableSyncable(table)把共享表包装成syncable.Syncable注册进 manager。依赖周期性 delta push pull完成同步不会在每次AddPrefix时调用Put做写穿。由此得到的是一致性模型为最终一致多副本间状态传播的滞后上界由同步周期 每 Key TTL共同决定。PrefixHashTableSyncable的完整实现namespace 为prefixcache在 pkg/utils/prefixcacheindexer/sync.go并显式声明var _ syncable.DeltaSyncable (*PrefixHashTableSyncable)(nil)来保证接口契约。如果你想更快地让其他副本看到本副本新缓存的前缀块README 给出了替代方案在每次AddPrefix之后手动调用syncManager.Put(ctx, prefixcache, blockIDStr, data)data 取自已实现的EncodeBlockForSync(blockHash)见 sync.go。代价是 Redis 写入负载上升——这是时效性与写入成本的权衡。关于 LRU 淘汰与删除传播的边界README 特别强调了一个重要边界EnableDeltaSync()必须在把表注册给*statesync.RedisSync之前调用否则脏追踪不生效GetDeltaForSync永远返回空 delta。此外如果一个块被标记为脏但在下一次成功推送前被本地 LRU 淘汰GetDeltaForSync会静默跳过它见 sync.go 注释 与 sync.go 中store.Get未命中的分支。对应的 Redis Key 不会被删除只能等每实体 TTL 到期自然消失。因此对已淘汰条目不保证强删除传播。这是 LRU 缓存与同步层组合下的固有取舍需要在设计自己的 Syncable 时留意。相关的跨副本行为都有测试覆盖例如TestCrossReplicaPrefixCacheLookupAfterSync验证了副本 A 写入 → 取 delta → 副本 B ApplyRemote → B 能命中 A 写入的前缀前提是两端使用相同 hash seedTestDeltaLifecycle验证了 delta 在AddPrefix后非空、ClearDirtyForSync后清空见 pkg/utils/prefixcacheindexer/sync_test.go。四、核心接口syncable.Syncable 与可选扩展statesync 不感知具体状态类型全部通过 pkg/utils/syncable/syncable.go 中定义的接口交互。实现的归属原则是接口实现放在拥有该状态的包里statesync 只依赖接口绝不反向依赖组件类型。必选接口Syncable方法用途Namespace() string该状态的稳定名称如mytracker会进入 Redis Keyaibrix:{namespace}:e:{id}中的{namespace}即 hash tag。GetSnapshot(ctx) (map[string][]byte, error)返回当前本地状态形态为实体 id → 序列化字节。调用方可能修改返回的 map实现不应保留引用。ApplyRemote(ctx, id string, data []byte) error将 Redis 中的一个实体应用到本地状态合并或覆盖。可选扩展接口接口用途DeltaSyncable额外提供GetDelta/ClearDirty推送时只发变更实体而不是全量快照。注意GetDelta内部不要清脏ClearDirty由同步层在推送成功后调用。OptionalHooksOnSyncStart/OnSyncEnd在每个 Syncable 的每次 Pull 前后调用不包住 Push。StalePolicyIsStale——返回 true 时同步层会删除该远程 Key 并跳过应用用于在 Redis TTL 之外再做应用层级的过期判断例如基于 lastUpdated 元数据。TombstoneSupportMakeTombstone/IsTombstone/DeleteLocal——把删除以墓碑载荷传播而不是DEL适合删除也需要被对端感知的场景。这些接口的语义注释如GetDelta不要清脏、ClearDirty由同步层调用等都原样写在 syncable.go是接入时最容易踩坑的约定。五、在新组件中接入 statesync 的四个步骤README 明确接入改动只发生在拥有该状态的包内statesync包保持通用。四个步骤如下。1. 实现syncable.Syncable依赖github.com/vllm-project/aibrix/pkg/utils/syncable实现必选三方法Namespace / GetSnapshot / ApplyRemote按需叠加第 4 节中的可选接口。2. 添加序列化编码Encode把结构体转成[]byte例如 JSON必须使用所有副本都能解码的稳定格式。解码Decode在ApplyRemote中反序列化data合并进或覆盖本地状态。3. 暴露返回syncable.Syncable的构造函数返回一个持有本地状态指针、实现上述接口的适配器调用方把它传给(*statesync.RedisSync).Register(...)。参考范例就是prefixcacheindexer.NewPrefixHashTableSyncable(table)。4. 在进程入口接线rs : statesync.New(redisClient, statesync.WithSyncPeriod(30*time.Second), // optional; default 10s statesync.WithOpTimeout(15*time.Second), // optional; default 30s statesync.WithKeyPrefix(myapp), // optional; default aibrix ) // Register all Syncables before Start. rs.Register(mypkg.NewMyStateSyncable(myState)) rs.Start() // On shutdown: rs.Stop() // blocks until the sync loop exits or stop-wait timeout (default 90s)可用配置项statesync.New的 OptionsOption默认值说明WithSyncPeriod(d)10s同步周期会叠加抖动。WithOpTimeout(d)30s每轮 Pull/Push 的操作级 context 超时。WithKeyPrefix(s)aibrix所有 Redis Key 的前缀。WithRecordTTL(d)2m每实体SETEX的 TTLd 0会被忽略。WithStopWaitTimeout(d)90sStop()等待后台同步循环退出的最大时间。WithSetexChunkSize(n)25SETEX 写入的 pipeline 分块大小。WithMGetBatchSize(n)200Pull 阶段 MGET 的批量大小。这些默认值在 redissync.go 常量区 中有精确对应defaultSyncPeriod、defaultRecordTTL、setexChunkSize、mgetBatchSize、opTimeout、maxBackoff等。此外同步周期还支持通过环境变量AIBRIX_STATESYNC_SYNC_PERIOD覆盖默认值见 redissync.go这比代码内传参更便于运维调整。Redis 客户端从哪来Redis 客户端由调用方提供例如使用utils.GetRedisClient()生产环境应在客户端选项中配置 TLS 与鉴权参见 pkg/utils/redis.go。六、完整实战为组件 TokenTracker 接入 statesync下面以 README 中的完整示例为主线——一个维护tokenID → lastUsedTime映射、需要跨网关副本同步的组件。Step 1 – 在组件包内实现 Syncable// pkg/plugins/gateway/algorithms/vtc/token_tracker_sync.go import ( context encoding/json strconv time github.com/vllm-project/aibrix/pkg/utils/syncable ) const tokenTrackerNamespace token_tracker // TokenTrackerSyncable adapts TokenTracker for statesync. type TokenTrackerSyncable struct { Tracker *TokenTracker } func NewTokenTrackerSyncable(t *TokenTracker) syncable.Syncable { return TokenTrackerSyncable{Tracker: t} } func (s *TokenTrackerSyncable) Namespace() string { return tokenTrackerNamespace } func (s *TokenTrackerSyncable) GetSnapshot(ctx context.Context) (map[string][]byte, error) { s.Tracker.mu.RLock() defer s.Tracker.mu.RUnlock() out : make(map[string][]byte, len(s.Tracker.entries)) for id, t : range s.Tracker.entries { b, _ : json.Marshal(t.UnixNano()) out[strconv.FormatInt(id, 10)] b } return out, nil } func (s *TokenTrackerSyncable) ApplyRemote(ctx context.Context, id string, data []byte) error { var nano int64 if err : json.Unmarshal(data, nano); err ! nil { return err } parsed, _ : strconv.ParseInt(id, 10, 64) remote : time.Unix(0, nano) s.Tracker.mu.Lock() defer s.Tracker.mu.Unlock() if existing, ok : s.Tracker.entries[parsed]; !ok || remote.After(existing) { s.Tracker.entries[parsed] remote } return nil }这个例子体现了三个关键设计点① 序列化格式稳定——时间统一编码为 UnixNano 的 JSON 数字所有副本可互相解码② 合并语义明确——ApplyRemote采用取较新的 lastUsedTime的合并策略而不是盲目覆盖③ 锁粒度安全——GetSnapshot用 RLock 快照拷贝ApplyRemote用 Lock 写回避免并发读写竞态。仓库中TokenTracker的真实实现位于 pkg/plugins/gateway/algorithms/vtc/token_tracker.go并有对应的单元测试 token_tracker_test.go 可以对照。Step 2 – 在入口接线rs : statesync.New(redisClient, statesync.WithSyncPeriod(30*time.Second)) rs.Register(prefixcacheindexer.NewPrefixHashTableSyncable(prefixHashTable)) rs.Register(vtc.NewTokenTrackerSyncable(tokenTracker)) rs.Start() defer rs.Stop()多个 Syncable 可以同时注册到同一个 manager共享一个同步循环和同一份 Redis 连接注意所有注册必须在Start()之前完成。Step 3 – 可选的写穿write-through如果希望其他副本在下一次周期推送之前就看到本次变更func (t *TokenTracker) RecordUse(ctx context.Context, id int64, rs *statesync.RedisSync) { t.mu.Lock() t.entries[id] time.Now() t.mu.Unlock() if rs ! nil { data, _ : json.Marshal(time.Now().UnixNano()) _ rs.Put(ctx, tokenTrackerNamespace, strconv.FormatInt(id, 10), data) } }Put内部就是SETEX带每实体 TTL见 redissync.go因此写穿的记录同样享受独立过期语义。是否使用写穿是一个明确的权衡要更快的跨副本可见性就承担更高的 Redis 写入负载默认网关接线前缀缓存选择了周期性增量推送 不写穿的低负载方案。七、参考实现与测试验证statesync 本身的单元测试位于 pkg/plugins/gateway/statesync/redissync_test.go使用miniredis模拟 Redis覆盖了 fake Syncable / fake DeltaSyncable 的 Pull、Push、TTL、墓碑、抖动与退避等行为是理解各选项实际效果的最佳活文档。前缀缓存的同步参考实现则集中在 pkg/utils/prefixcacheindexer/sync.goPrefixHashTableSyncable与NewPrefixHashTableSyncable(table)—— 完整的DeltaSyncable实现GetSnapshotForSync/GetDeltaForSync/ClearDirtyForSync/ApplyRemoteForSync/EncodeBlockForSync—— 挂在PrefixHashTable上的同步方法族接入前记得调用table.EnableDeltaSync()需要 delta push 时如果只用全量快照同步则不需要。对应的行为验证在 pkg/utils/prefixcacheindexer/sync_test.goTestDeltaLifecycle验证 dirty 标记的Add 后出现、Clear 后清空生命周期TestCrossReplicaPrefixCacheLookupAfterSync端到端验证副本 A 的增量被副本 B 应用后B 能命中 A 缓存的 prefixTestCrossReplicaPrefixCacheLookupRequiresSharedSeed则从反面证明两端 hash seed 不一致时即使块载荷同步成功也无法命中——seed 一致性是跨副本前缀缓存可用的前提。结语statesync 以每实体 Key 独立 TTL Pull-first/Push的朴素设计为 AIBrix 网关插件提供了一套无 PUB/SUB 依赖、易于理解和接入的跨副本状态同步方案。它的通用性来自syncable接口的抽象任何id → bytes形态的本地状态都可以在四个步骤内接入而它的边界也同样清晰——最终一致性、LRU 淘汰不保证强删除传播、delta 推送依赖EnableDeltaSync预先开启。理解了这些机制你就能在自己的网关组件中按需选择周期性 delta 推送与写穿 Put两种同步策略在时效性与 Redis 写入成本之间做出正确取舍。【免费下载链接】aibrixCost-efficient and pluggable Infrastructure components for GenAI inference项目地址: https://gitcode.com/GitHub_Trending/ai/aibrix创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考