深入解析 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 的跨副本状态同步层:以"每个实体一个 Key"(SETEX+ 每记录 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 的包注释与实现中验证):
- 每实体一个 Key(per-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)。 - 用
SCAN(COUNT 提示为 500)按aibrix:{namespace}:e:*模式枚举该 namespace 下所有实体 Key,累积到mgetBatchSize(默认 200)一批,调用MGET批量取值。 - 对每个值依次处理:若实现了
TombstoneSupport且载荷是墓碑,则调用DeleteLocal删除本地条目并跳过;若实现了StalePolicy且IsStale返回 true,则删除 Redis 中该 Key并跳过应用;否则调用ApplyRemote(id, bytes)合入本地状态。 - Redis 会过滤已过期的 Key,因此拉取到的都是 TTL 窗口内的有效记录。
- 先检查 Syncable 是否实现了
- Push 阶段(Push 函数):
- 若 Syncable 实现了
DeltaSyncable,走pushDelta:只推送自上次ClearDirty以来变更的实体(updated 用SETEX,deleted 用DEL或墓碑),推送成功后调用ClearDirty清空脏标记。 - 否则走
pushFull:调用GetSnapshot拿到全量快照,全部用SETEX写入;快照超过 1 万字段或平均每字段超过 512KB 时会打印告警。 - 写入使用 pipeline 分块(
setexChunkSize默认 25),降低往返开销。
- 若 Syncable 实现了
- 循环调度(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.go:
stateSyncEnabled := 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_ENABLED=true 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 seed;TestDeltaLifecycle验证了 delta 在AddPrefix后非空、ClearDirtyForSync后清空(见 pkg/utils/prefixcacheindexer/sync_test.go)。
四、核心接口:syncable.Syncable 与可选扩展
statesync 不感知具体状态类型,全部通过 pkg/utils/syncable/syncable.go 中定义的接口交互。实现的归属原则是:接口实现放在"拥有该状态的包"里,statesync 只依赖接口,绝不反向依赖组件类型。
必选接口Syncable
| 方法 | 用途 |
|---|---|
Namespace() string | 该状态的稳定名称(如"mytracker"),会进入 Redis Key,aibrix:{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由同步层在推送成功后调用。 |
OptionalHooks | OnSyncStart/OnSyncEnd,在每个 Syncable 的每次 Pull 前后调用(不包住 Push)。 |
StalePolicy | IsStale——返回 true 时,同步层会删除该远程 Key 并跳过应用,用于在 Redis TTL 之外再做应用层级的过期判断(例如基于 lastUpdated 元数据)。 |
TombstoneSupport | MakeTombstone/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的 Options)
| Option | 默认值 | 说明 |
|---|---|---|
WithSyncPeriod(d) | 10s | 同步周期(会叠加抖动)。 |
WithOpTimeout(d) | 30s | 每轮 Pull/Push 的操作级 context 超时。 |
WithKeyPrefix(s) | "aibrix" | 所有 Redis Key 的前缀。 |
WithRecordTTL(d) | 2m | 每实体SETEX的 TTL;d <= 0会被忽略。 |
WithStopWaitTimeout(d) | 90s | Stop()等待后台同步循环退出的最大时间。 |
WithSetexChunkSize(n) | 25 | SETEX 写入的 pipeline 分块大小。 |
WithMGetBatchSize(n) | 200 | Pull 阶段 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.go:
PrefixHashTableSyncable与NewPrefixHashTableSyncable(table)—— 完整的DeltaSyncable实现;GetSnapshotForSync/GetDeltaForSync/ClearDirtyForSync/ApplyRemoteForSync/EncodeBlockForSync—— 挂在PrefixHashTable上的同步方法族;- 接入前记得调用
table.EnableDeltaSync()(需要 delta push 时;如果只用全量快照同步则不需要)。
对应的行为验证在 pkg/utils/prefixcacheindexer/sync_test.go:TestDeltaLifecycle验证 dirty 标记的"Add 后出现、Clear 后清空"生命周期;TestCrossReplicaPrefixCacheLookupAfterSync端到端验证"副本 A 的增量被副本 B 应用后,B 能命中 A 缓存的 prefix";TestCrossReplicaPrefixCacheLookupRequiresSharedSeed则从反面证明:两端 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),仅供参考