从“能工作”到“可扩展、可靠、高性能”的系统演进
本页基于 Mooncake 项目历史讨论整理。重点不是罗列提交,而是解释每项工作的系统问题、 数据路径变化和架构价值:UB 的 TENT-native 化、连接预热、LOCAL_DISK 可靠删除、 HotCache 一致性修复、RDMA 数据路径缩短,以及 SSD 分配的反馈式负载均衡。
Mooncake 总体架构:关键模块与本次修改落点
flowchart TB
subgraph APP["应用与接入层"]
direction LR
A1["LLM / Inference Service"]
A2["Python Binding / Client SDK"]
A3["Mooncake Store Client API
Put / Get / Remove / BatchGet"]
A1 --> A2 --> A3
end
subgraph CACHE["客户端缓存与访问路径"]
direction LR
HC["HotCache"]
DC["DummyClient / RealClient"]
WU["Warmup Manager"]
A3 --> HC
A3 --> DC
DC --> WU
end
subgraph CTRL["控制面 / 元数据面"]
direction TB
MC["MasterClient"]
WS["WrappedMasterService / MasterService"]
META["Object Metadata / Replica Metadata"]
ALLOC["Allocation Strategy
SSD_FREE_RATIO_FIRST"]
RPC["ListWarmupTargets / Remove / FetchRemoveTasks / AckRemoveTasks"]
PLUGIN["Metadata Store Plugin
etcd / P2P Handshake"]
MC --> WS
WS --> META
WS --> ALLOC
WS --> RPC
WS --> PLUGIN
end
subgraph TE["数据搬运层"]
direction LR
FAC["TransferEngine Facade"]
subgraph OLDTE["Classic Transfer Engine"]
direction TB
MT["MultiTransport"]
RDMA0["RdmaTransport"]
TCP0["TcpTransport"]
UB0["UbTransport"]
META0["TransferMetadata / Handshake"]
MT --> RDMA0
MT --> TCP0
MT --> UB0
UB0 --> META0
end
subgraph TENT["TENT Runtime"]
direction TB
TI["tent::TransferEngineImpl"]
TS["TransportSelector"]
TL["TransportLoader"]
SM["SegmentManager"]
PFM["Platform / Topology"]
TUB["UB Transport
Non-native / Native"]
TRDMA["RdmaTransport"]
TTCP["TcpTransport"]
TI --> TS
TI --> TL
TI --> SM
TI --> PFM
TL --> TUB
TL --> TRDMA
TL --> TTCP
end
FAC --> OLDTE
FAC --> TENT
end
subgraph STORAGE["存储与后台任务"]
direction LR
MEM["Memory Segment / KV Buffer"]
FILE["FileStorage"]
BUCKET["BucketStorageBackend"]
GC["Bucket GC Worker"]
DFS["DFS / Remote Storage"]
FILE --> BUCKET --> GC
end
subgraph HW["网络与硬件"]
direction LR
RDMAH["RDMA NIC / verbs"]
UBH["UB NIC / URMA / liburma"]
TCPH["TCP Socket"]
SSD["LOCAL_DISK / SSD"]
DRAM["Host DRAM"]
end
HC --> FAC
DC --> MC
DC --> FAC
WU --> MC
WU --> FAC
META --> FAC
ALLOC --> FILE
RPC --> FILE
FAC --> MEM
FAC --> FILE
FAC --> DFS
RDMA0 --> RDMAH
TCP0 --> TCPH
UB0 --> UBH
TRDMA --> RDMAH
TTCP --> TCPH
TUB --> UBH
MEM --> DRAM
FILE --> SSD
BUCKET --> SSD
classDef core fill:#dfecff,stroke:#3973c6,stroke-width:2px,color:#123;
classDef mine fill:#ddf5e3,stroke:#32a852,stroke-width:3px,color:#132;
classDef key fill:#fff0cc,stroke:#e39a20,stroke-width:3px,color:#321;
classDef prob fill:#ffe1e1,stroke:#d64545,stroke-width:2px,color:#311;
class A1,A2,A3,MC,WS,META,PLUGIN,FAC,MT,RDMA0,TCP0,UB0,META0,TI,TS,TL,SM,PFM,TRDMA,TTCP,MEM,FILE,DFS,RDMAH,UBH,TCPH,SSD,DRAM core;
class HC,WU,ALLOC,RPC,TUB,BUCKET,GC mine;
class DC key;
class FILE,HC,RPC,ALLOC prob;
HotCache、Warmup Manager、Allocation Strategy、删除相关 RPC、TENT 中的 UB Transport、BucketStorageBackend 与 Bucket GC Worker。| 工作项 | 主要落点 | 在总体架构中的位置 | 关联模块 |
|---|---|---|---|
| UB 接入 TENT(Non-native / Native) | TENT Runtime / UB Transport | 数据搬运层 | TransportSelector、TransportLoader、SegmentManager、UB data plane |
| Warmup 预建链 | Store Client + Master + TE handshake | 客户端访问路径 + 控制面 + 数据搬运层 | Client::warmup、MasterClient、ListWarmupTargets、openSegment |
| Bucket 删除与 GC | MasterService + FileStorage + BucketStorageBackend + GC | 控制面 + 存储层 | FetchRemoveTasks、AckRemoveTasks、MarkRemoved、BucketGcWorker |
| HotCache 修复 | Store Client HotCache path | 客户端缓存与访问路径 | HotCache、Get/Remove/Upsert 路径 |
| RDMA 少一跳优化 | Store 远端读流程 + RDMA transport 编排 | 数据搬运层 + 远端存储读取路径 | Requester / Target staging / RDMA READ→WRITE 编排 |
| SSD 负载均衡 | Master allocation strategy | 控制面 / 资源调度 | SSD_FREE_RATIO_FIRST、ssd_used_bytes accounting |
flowchart LR
REQ["应用发起 Put / Get / Remove"]
--> API["Store Client API"]
--> PATH{"操作类型"}
PATH -->|"Get / BatchGet"| HOT["HotCache / Metadata Lookup"]
PATH -->|"Put / Write"| PUT["Store metadata + allocate replica"]
PATH -->|"Remove"| DEL["Master Remove / Delete Task"]
HOT --> TE1["TransferEngine / TENT"]
PUT --> TE1
DEL --> CTRL1["Master / Holder / GC 协议"]
TE1 --> XPORT{"传输后端"}
XPORT --> RD["RDMA"]
XPORT --> UB["UB / URMA"]
XPORT --> TCP["TCP"]
RD --> DATA["远端 Memory / SSD / DFS 数据"]
UB --> DATA
TCP --> DATA
CTRL1 --> FILE1["LOCAL_DISK Holder / BucketStorageBackend"]
FILE1 --> GC1["GC 回收"]
classDef core fill:#dfecff,stroke:#3973c6,stroke-width:2px,color:#123;
classDef mine fill:#ddf5e3,stroke:#32a852,stroke-width:3px,color:#132;
classDef key fill:#fff0cc,stroke:#e39a20,stroke-width:3px,color:#321;
class REQ,API,PATH,TE1,XPORT,RD,UB,TCP,DATA core;
class HOT,DEL,CTRL1,FILE1,GC1,PUT mine;
class PUT,HOT key;
UB 接入 TENT:从兼容层到 Native 数据面
非 Native 版本
- TENT 负责 Transport 加载、选择、控制面和 UB bootstrap。
UbTentTransport将 TENT Request / Metadata 转换为 Classic TE 模型。- 内存注册、切片、设备选择、Endpoint 与 URMA 数据面仍由 Classic
UbTransport执行。
Native 版本
- TENT Request 直接生成
UbTask / UbSlice,不再转换成 Classic TransferRequest。 UbWorkers负责切片执行、路径选择、提交与完成处理。- 通过
UrmaAdapter直接对接 liburma,UB 数据面完整归入 TENT。
flowchart TB
APP["Application / Mooncake Store"]
subgraph TENT["TENT Layer"]
direction LR
RT["TENT Runtime"] --> SEL["Transport Selector"] -->|"选择 UB"| UBT["UbTentTransport"] --> BR["Request / Metadata Bridge"]
end
subgraph CLASSIC["Classic TE UB Data Plane"]
direction LR
OLD["Classic UbTransport"] --> CTX["UbContext / UbEndpoint"] --> W["UB Worker"]
OLD --> TOPO["Legacy Topology"]
end
subgraph URMA["URMA / UMDK"]
direction LR
UC["UrmaContext / UrmaEndpoint"] --> LIB["liburma"] --> HW["UB NIC"]
end
APP --> RT
BR -->|"TENT Request → TransferRequest"| OLD
W --> UC
classDef new fill:#ddf5e3,stroke:#32a852,stroke-width:3px,color:#132;
classDef bad fill:#ffe1e1,stroke:#d64545,stroke-width:3px,color:#311;
classDef core fill:#dfecff,stroke:#3973c6,stroke-width:2px,color:#123;
classDef key fill:#fff0cc,stroke:#e39a20,stroke-width:3px,color:#321;
class UBT new;
class BR key;
class OLD,TOPO bad;
class RT,SEL,CTX,W,UC,LIB,HW core;
classDiagram
direction LR
class Transport {
<<interface>>
+install()
+uninstall()
+allocateSubBatch()
+freeSubBatch()
+submitTransferTasks()
+getTransferStatus()
+addMemoryBuffer()
+removeMemoryBuffer()
+getEstimatedBandwidth()
+cancelTransferTask()
}
class UbTentTransport {
<<TENT adapter>>
-legacyUbTransport
-metadataBridge
+install()
+addMemoryBuffer()
+allocateSubBatch()
+submitTransferTasks()
+getTransferStatus()
+cancelTransferTask()
}
class UbTentMetadataBridge {
<<metadata adapter>>
+getSegmentDescByName()
+getSegmentDescByID()
+translateHandshake()
}
class TransferMetadata {
<<Classic TE interface>>
+SegmentDesc
+HandShakeDesc
}
class SegmentManager {
+BufferDesc
+SegmentDesc
+lookupSegment()
}
class ControlService {
+BootstrapUb()
}
class UbTransport {
<<Classic TE UbTransport>>
+install()
+registerLocalMemory()
+allocateBatchID()
+submitTransfer()
+getTransferStatus()
}
class BatchID {
<<Classic TE task handle>>
}
class TransferRequest {
+opcode
+source
+target
+length
}
class TransferTask {
+status
+transferredBytes
}
Transport <|-- UbTentTransport
UbTentTransport --> UbTransport : delegates data path
UbTentTransport --> UbTentMetadataBridge : metadata / handshake
UbTentMetadataBridge ..|> TransferMetadata
UbTentMetadataBridge --> SegmentManager : TENT segment → legacy view
UbTentMetadataBridge --> ControlService : BootstrapUb
UbTransport --> TransferMetadata
UbTransport --> BatchID
UbTransport --> TransferRequest
UbTransport --> TransferTask
| TENT 侧 | 适配层动作 | Classic TE 侧 |
|---|---|---|
| Transport / SubBatch / Request | UbTentTransport 将 TENT 的提交与状态接口转换为旧 TE 的 batch/task 调用 |
BatchID / TransferRequest / TransferTask |
| SegmentManager / BufferDesc | UbTentMetadataBridge 提供旧 TransferMetadata 所需的 segment 视图 |
TransferMetadata::SegmentDesc |
| TENT ControlService | 通过 BootstrapUb 完成 UB 控制面握手并转换为旧握手描述 |
HandShakeDesc / Endpoint 初始化信息 |
| addMemoryBuffer() | 调用旧 UB 注册本地内存,并把生成的 UB segment 信息回写到 TENT BufferDesc | URMA tseg / registered segment |
sequenceDiagram
participant APP as Store / App
participant TENT as TENT Runtime
participant ADP as UbTentTransport
participant TE as Classic UbTransport
participant EP as UbEndpoint
participant URMA as liburma / NIC
APP->>TENT: submit(Request)
TENT->>TENT: TransportSelector 选择 UB
rect rgb(221,245,227)
TENT->>ADP: TENT Request
end
rect rgb(255,240,204)
ADP->>ADP: Request → TransferRequest
ADP->>TE: submitTransfer()
end
rect rgb(255,225,225)
TE->>TE: Classic topology / slice / batch
TE->>EP: 获取或建立 Endpoint
end
EP->>URMA: post READ / WRITE
URMA-->>EP: CQE
EP-->>TE: complete
TE-->>ADP: Classic status
ADP-->>TENT: TENT status
TENT-->>APP: Complete
flowchart TB
APP["Application / Store"]
subgraph TENT["TENT Runtime"]
direction LR
RT["TENT Runtime"] --> SEL["Transport Selector"] --> UB["Native UB Transport"]
end
subgraph EXEC["TENT-native UB Execution"]
direction LR
TASK["UbTask"] --> SL["UbSlice[]"] --> W["UbWorkers"] --> PATH["Path / Rail Selection"]
end
subgraph DATA["UB Data Plane"]
direction LR
CTX["UbContext"] --> EP["UbEndpoint"] --> UA["UrmaAdapter"]
end
subgraph UMDK["URMA / UMDK"]
direction LR
LIB["liburma"] --> HW["UB NIC"] --> JFC["UbJfc / Completion"]
end
APP --> RT
UB --> TASK
PATH --> CTX
UA --> LIB
JFC --> W
W --> TASK
classDef new fill:#ddf5e3,stroke:#32a852,stroke-width:3px,color:#132;
classDef key fill:#fff0cc,stroke:#e39a20,stroke-width:3px,color:#321;
classDef core fill:#dfecff,stroke:#3973c6,stroke-width:2px,color:#123;
class UB,TASK,SL,W new;
class PATH key;
class RT,SEL,CTX,EP,UA,LIB,HW,JFC core;
classDiagram
direction TB
class Transport {
<<TENT transport contract>>
+install()
+uninstall()
+allocateSubBatch()
+freeSubBatch()
+submitTransferTasks()
+getTransferStatus()
+addMemoryBuffer()
+removeMemoryBuffer()
+getEstimatedBandwidth()
+supportsCancellation()
+cancelTransferTask()
}
class UbTransport {
<<TENT-native UB>>
+allocateSubBatch()
+submitTransferTasks()
+getTransferStatus()
+addMemoryBuffer()
+getEstimatedBandwidth()
+cancelTransferTask()
}
class RdmaTransport {
<<TENT RDMA>>
+allocateSubBatch()
+submitTransferTasks()
+getTransferStatus()
+addMemoryBuffer()
+getEstimatedBandwidth()
+cancelTransferTask()
}
class TransportSelector {
+selectTransport()
}
class SegmentManager {
+registerBuffer()
+lookupSegment()
}
class Platform {
<<shared runtime abstraction>>
+probe()
+memory / topology info
}
class UbTask {
+status
+totalBytes
+transferredBytes
+cancelRequested
}
class UbSlice {
+offset
+length
+retryCount
+endpointGeneration
}
class UbWorkers {
+submitBatch()
+selectPath()
+post()
+poll()
}
class UbContext
class UbEndpoint
class UbJfc
class UrmaAdapter
Transport <|-- UbTransport
Transport <|-- RdmaTransport
TransportSelector --> Transport : only depends on common API
SegmentManager --> UbTransport
SegmentManager --> RdmaTransport
Platform --> UbTransport
Platform --> RdmaTransport
UbTransport --> UbTask
UbTask --> UbSlice
UbSlice --> UbWorkers
UbWorkers --> UbContext
UbWorkers --> UbEndpoint
UbWorkers --> UbJfc
UbEndpoint --> UrmaAdapter
| 共同数据面接口 | UB | RDMA | TENT 看到的语义 |
|---|---|---|---|
| install / uninstall | 初始化/释放 URMA 资源 | 初始化/释放 RDMA 资源 | Transport 生命周期 |
| allocateSubBatch / freeSubBatch | UB 子批次 | RDMA 子批次 | 统一 batch 生命周期 |
| submitTransferTasks | Request → UbTask / UbSlice | 进入 RDMA transport 的任务执行 | 统一 READ / WRITE 提交入口 |
| getTransferStatus | 汇总 UB task/slice 状态 | 汇总 RDMA 任务状态 | 统一完成/失败查询 |
| add/removeMemoryBuffer | 注册/释放 URMA segment | 注册/释放 RDMA MR | 统一内存注册接口 |
| getEstimatedBandwidth | UB path / rail 估计 | RDMA path 估计 | 供 selector / scheduler 使用 |
| supportsCancellation / cancelTransferTask | UB task 取消语义 | RDMA 取消语义 | 统一任务控制能力 |
sequenceDiagram
participant APP as Store
participant TENT as TENT Runtime
participant UB as Native UB Transport
participant W as UbWorkers
participant EP as UbEndpoint
participant UA as UrmaAdapter
participant HW as UB NIC / JFC
APP->>TENT: submit(Request)
TENT->>UB: submitTransferTasks()
rect rgb(221,245,227)
UB->>UB: create UbTask
UB->>UB: slice → UbSlice[]
UB->>W: submitBatch(slices)
end
rect rgb(255,240,204)
W->>W: selectPath()
W->>EP: get endpoint
end
EP->>UA: post(slice)
UA->>HW: URMA post
HW-->>UA: CQE
UA-->>W: completion
W->>UB: update UbTask
UB-->>TENT: Task Complete
TENT-->>APP: Complete
结构收益
去掉 Classic TE 中间层,数据路径归属清晰,TENT 真正成为统一传输框架。
调度收益
路径选择、slice、重试、带宽估计等能力可以在同一套 TENT 语义下扩展。
生态收益
UB 接入方式与 RDMA 等 transport 更接近,为统一抽象和后续开源维护降低成本。
Warmup:把首次建链成本移出业务热路径
flowchart LR
subgraph CTRL["控制面"]
direction TB
S["Store Client 初始化"] --> W["Warmup Manager"]
W --> L["ListWarmupTargets"]
L --> M["Master"]
M -->|"有界 targets"| W
end
subgraph INIT["连接初始化"]
direction TB
O["openSegment()"] --> META["Remote Metadata"] --> HS["Socket Handshake"] --> EP["RDMA / UB / TCP Endpoint"]
end
subgraph PROBE["Warmup Probe"]
direction TB
P["Small READ-only Probe"] --> EP2["复用已建立 Endpoint"]
P --> C["Cleanup Worker"]
end
W --> O
W --> P
EP --> EP2
B["后续真实 GET / WRITE"] -->|"直接复用"| EP2
classDef new fill:#ddf5e3,stroke:#32a852,stroke-width:3px,color:#132;
classDef key fill:#fff0cc,stroke:#e39a20,stroke-width:3px,color:#321;
classDef core fill:#dfecff,stroke:#3973c6,stroke-width:2px,color:#123;
class W,L,P,C new;
class HS,EP key;
class M,O,META core;
classDiagram
direction LR
class Client {
+setup()
+warmup()
}
class MasterClient {
+ListWarmupTargets()
}
class WrappedMasterService {
+ListWarmupTargets(clientId, maxTargets, preferredProtocols)
}
class WarmupTarget {
+segmentName
+segmentId
+clientId
+protocol
+isLocal
+allowWarmup
+priority
}
class TransferEngine {
+openSegment()
+submitTransfer(READ)
+getTransferStatus()
+freeBatchID()
}
class ProbeResource {
+buffer
+batchId
+status
}
class CleanupWorker {
+poll()
+requeueOnStatusError()
+releaseOnTerminal()
}
Client --> MasterClient : query targets
MasterClient --> WrappedMasterService
WrappedMasterService --> WarmupTarget : returns bounded list
Client --> TransferEngine : open + READ probe
Client --> ProbeResource : create
CleanupWorker --> ProbeResource
CleanupWorker --> TransferEngine : query/release
没有 Warmup
- 第一批业务请求同时触发 openSegment / metadata 查询。
- 多个请求并发进入 socket handshake。
- listener / worker / 状态轮询被短时间突发放大。
- 结果体现为尾延迟、超时,极端情况下出现 hang / fail。
加入 Warmup
- Client 初始化阶段先向 Master 获取有限数量目标。
- 对目标执行 READ-only probe,建立并缓存 transport endpoint。
- 资源清理与状态查询异步完成,避免影响真实请求。
- 业务流量到来后直接复用已建连接。
sequenceDiagram
participant R as 大量业务请求
participant C as Store Client
participant TE as Transfer Engine
participant H as Handshake Listener
participant P as Peer
rect rgb(255,225,225)
R->>C: GET / WRITE × N
par Request 1
C->>TE: openSegment()
TE->>H: handshake
and Request 2
C->>TE: openSegment()
TE->>H: handshake
and Request N
C->>TE: openSegment()
TE->>H: handshake
end
H->>P: 集中握手
Note over H,P: listener / metadata path 瞬时拥塞
P-->>H: response
end
H-->>TE: endpoint ready / timeout
TE-->>C: transfer / fail
sequenceDiagram
participant C as Store Client
participant M as Master
participant TE as Transfer Engine
participant P as Peer
participant R as 真实业务请求
rect rgb(221,245,227)
C->>M: ListWarmupTargets(max_targets)
M-->>C: target segments
loop bounded targets
C->>TE: openSegment(target)
TE->>P: metadata / handshake
P-->>TE: endpoint ready
C->>TE: small READ-only probe
TE->>P: READ
P-->>TE: complete
end
end
R->>C: GET / WRITE × N
C->>TE: submit transfer
Note over C,TE: endpoint 已存在
TE->>P: direct data transfer
Bucket 数据删除:从 metadata 删除到可靠的物理回收协议
flowchart LR
subgraph MASTER["Master"]
direction TB
U["RealClient
Remove / BatchRemove"] --> M["Object Metadata"]
M --> G["生成 Holder-specific Tasks"]
G --> Q["Bounded Remove Queue"]
G -->|"全部成功入队后"| DEL["删除 Object Metadata"]
end
subgraph HOLDER["LOCAL_DISK Holder"]
direction TB
H["Heartbeat"] --> F["FetchRemoveTasks"]
F --> MR["BucketStorageBackend::MarkRemoved"]
MR --> CK["校验 tenant / key / incarnation"]
CK --> D["tmp → fsync → rename → fsync(parent)"]
D --> A["AckRemoveTasks"]
end
subgraph GCBOX["Background GC"]
direction TB
GC["GC Worker"] --> CP["Bucket Compaction"] --> R["回收 SSD 物理空间"]
end
Q --> F
A --> Q
D --> GC
classDef new fill:#ddf5e3,stroke:#32a852,stroke-width:3px,color:#132;
classDef key fill:#fff0cc,stroke:#e39a20,stroke-width:3px,color:#321;
classDef core fill:#dfecff,stroke:#3973c6,stroke-width:2px,color:#123;
class G,Q,F,MR,A,GC new;
class CK,D key;
class U,M,H,CP,R,DEL core;
classDiagram
direction LR
class MasterService {
+Remove()
+BatchRemove()
+FetchRemoveTasks()
+AckRemoveTasks()
}
class RemoveTaskQueue {
+enqueue()
+fetchWithoutDelete()
+ackAndErase()
+capacity
}
class RemoveTask {
+taskId
+clientId
+tenantId
+key
+objectIncarnation
+replicaIdentity
}
class FileStorage {
+heartbeat()
+fetchRemoveTasks()
+ackRemoveTasks()
}
class BucketStorageBackend {
+MarkRemoved()
+persistMetadata()
+compactBucket()
}
class BucketGcWorker {
+start()
+scanCandidates()
+compact()
+stopAndJoin()
}
MasterService --> RemoveTaskQueue
RemoveTaskQueue o-- RemoveTask
FileStorage --> MasterService : heartbeat RPC
FileStorage --> BucketStorageBackend : execute task
BucketStorageBackend --> RemoveTask : validate incarnation
BucketGcWorker --> BucketStorageBackend : background GC
任务不丢
Fetch 只把任务交给 Holder,不立即从 Master 队列删除;只有 Holder 完成持久化并 ACK 后才移除。
重复投递安全
MarkRemoved 必须幂等;RPC/ACK 丢失、Holder 重启时允许重新处理同一任务。
不会删错新对象
任务携带 object incarnation / replica identity;旧任务遇到同 key 新对象时视为 stale task。
sequenceDiagram
participant U as RealClient
participant M as Master
participant Q as Remove Queue
participant H as Holder
participant B as BucketStorageBackend
participant GC as GC Worker
U->>M: Remove(key)
M->>M: 查 LOCAL_DISK replicas
rect rgb(221,245,227)
loop each holder replica
M->>Q: enqueue(holder-specific task)
end
end
Note over M,Q: 所有任务成功入队
M->>M: erase object metadata
H->>M: FetchRemoveTasks(max_tasks)
M-->>H: bounded task batch
rect rgb(255,240,204)
H->>B: MarkRemoved(task)
B->>B: validate tenant/key/incarnation
B->>B: write tmp + fsync
B->>B: atomic rename + fsync parent
end
B-->>H: durable
rect rgb(221,245,227)
H->>M: ACK(task_id)
M->>Q: remove task
end
GC->>B: scan tombstones
GC->>B: compact bucket
B-->>GC: reclaim old file
HotCache:修复异步 Fill 与 Remove/Upsert 的 stale race
sequenceDiagram
participant G as GET Thread
participant HC as HotCache
participant S as Storage / Metadata
participant R as Remove / Upsert Thread
G->>HC: Get(key)
HC-->>G: MISS
G->>S: read old object V1
G->>G: async Fill(V1) pending
rect rgb(255,225,225)
R->>HC: invalidate(key)
R->>S: Remove / Upsert
Note over G,R: 并发窗口
G->>HC: delayed Fill(V1)
Note over HC: 旧地址重新进入 HotCache
end
S-->>R: mutation success
G->>HC: next Get
HC-->>G: stale V1
sequenceDiagram
participant R as Remove / Upsert
participant HC as HotCache
participant S as Storage / Metadata
participant F as Racing Async Fill
rect rgb(221,245,227)
R->>HC: ① pre-mutation invalidate
end
R->>S: Remove / Upsert
rect rgb(255,225,225)
F->>HC: racing Fill(old V1)
end
S-->>R: mutation success
rect rgb(221,245,227)
R->>HC: ② success-side invalidate
Note over HC: 清掉窗口内重新插入的 stale entry
end
RDMA 少一跳:从 Requester Pull 改为 Target Push
flowchart LR
subgraph OLD["原方案 · RDMA READ Pull"]
R1["Requester"] -->|"① Read RPC"| T1["Target"]
SSD1["SSD"] -->|"② pread"| ST1["Target Staging"]
T1 -->|"③ staging addr/rkey"| R1
R1 -->|"④ RDMA READ"| ST1
R1 -->|"⑤ release"| T1
end
subgraph NEW["优化 · RDMA WRITE Push"]
R2["Requester Destination"] -->|"① RPC + dst addr/rkey"| T2["Target"]
SSD2["SSD"] -->|"② pread"| ST2["Target Staging"]
ST2 -->|"🟢 ③ RDMA WRITE"| R2
T2 -->|"④ complete"| R2
end
classDef bad fill:#ffe1e1,stroke:#d64545,stroke-width:3px,color:#311;
classDef good fill:#ddf5e3,stroke:#32a852,stroke-width:3px,color:#132;
class R1,T1,ST1 bad;
class R2,T2,ST2 good;
Requester Pull
Requester 需要等待 Target 准备 staging,再进行第二阶段 RDMA READ,并在结束后显式 release。
Target Push
Requester 的目标地址随初始 RPC 一次性传给 Target;Target SSD Read 完成后直接 WRITE 到最终 buffer。
sequenceDiagram
participant R as Requester
participant T as Target
participant SSD as Target SSD
participant S as Staging
rect rgb(255,225,225)
R->>T: Read RPC
T->>SSD: pread()
SSD-->>S: data
T-->>R: staging addr / rkey
Note over R,T: wait #1
R->>S: RDMA READ
S-->>R: data
Note over R,T: wait #2
R->>T: Release staging
end
sequenceDiagram
participant R as Requester Buffer
participant T as Target
participant SSD as Target SSD
participant S as Staging
rect rgb(221,245,227)
R->>T: Read RPC(dst addr + rkey)
T->>SSD: pread()
SSD-->>S: data
S->>R: RDMA WRITE
T-->>R: Complete
end
Note over R,T: Requester 不再额外发起 RDMA READ
SSD 负载均衡:从随机选择到资源反馈闭环
flowchart LR
subgraph MASTER["Master Allocation"]
direction TB
R["Object Offload Request"] --> A["Allocation Strategy"]
A --> S["SSD Usage Statistics"]
S --> C["计算 free_ratio"]
C --> P["SSD_FREE_RATIO_FIRST"]
end
subgraph HOLDERS["LOCAL_DISK Holders"]
direction TB
H1["Holder A · Used 80%"]
H2["Holder B · Used 45%"]
H3["Holder C · Used 20%"]
end
H1 -->|"usage"| S
H2 -->|"usage"| S
H3 -->|"usage"| S
P -. "skip" .-> H1
P -. "skip" .-> H2
P -->|"🟢 prefer"| H3
H3 -->|"allocation 后更新 used_bytes"| S
classDef good fill:#ddf5e3,stroke:#32a852,stroke-width:3px,color:#132;
classDef bad fill:#ffe1e1,stroke:#d64545,stroke-width:3px,color:#311;
classDef key fill:#fff0cc,stroke:#e39a20,stroke-width:3px,color:#321;
class H3,P good;
class H1 bad;
class S,C key;
flowchart LR
AL["Allocate LOCAL_DISK"] -->|"+"| U["ssd_used_bytes"]
U --> F["free_ratio"]
F --> S["Next Allocation"]
D["Remove / GC"] -->|"🟢 -"| U
classDef good fill:#ddf5e3,stroke:#32a852,stroke-width:3px,color:#132;
classDef key fill:#fff0cc,stroke:#e39a20,stroke-width:3px,color:#321;
class D good;
class U,F key;
只有“加”没有“减”会怎样
对象已经删除、物理空间也可能回收,但 Master 的统计不回退,节点会被永久误判为更满,后续策略持续避开它。
正确的资源模型
allocation、remove、metadata cleanup、GC 对 SSD 使用量的影响必须形成一致的 accounting 边界,否则调度算法再合理也会被错误输入破坏。
社区运作:从发现问题到 PR 合入的完整闭环
flowchart LR
A["查已有 Issue / PR"] --> B{"已有同类工作?"}
B -->|"有"| C["加入讨论 / 补充方案
避免重复实现"]
B -->|"没有"| D["创建 Issue / RFC"]
D --> E["基于 upstream/main
建立功能分支"]
E --> F["实现 + 回归测试"]
F --> G["format / pre-commit
本地测试"]
G --> H["同步 upstream/main
处理冲突"]
H --> I["提交 PR
Fixes / Closes #Issue"]
I --> J["GitHub CI"]
J --> K["Bot Review"]
K --> L["Maintainer / CODEOWNER Review"]
L --> M["修复问题 + 回复原因
Resolve threads"]
M --> N{"Required checks
+ approvals 完成?"}
N -->|"否"| F
N -->|"是"| O["Merge"]
classDef good fill:#ddf5e3,stroke:#32a852,stroke-width:3px,color:#132;
classDef key fill:#fff0cc,stroke:#e39a20,stroke-width:3px,color:#321;
classDef bad fill:#ffe1e1,stroke:#d64545,stroke-width:2px,color:#311;
class D,I,O good;
class J,L,N key;
class B,M bad;
Fixes #xxx / Closes #xxx。Issue 先固定问题归属和背景,PR 立即表明已有实现,避免只提 Issue 后长期空置或被其他人直接接走。先查社区,再开始实现
功能较大时怎么拆
flowchart LR U["upstream/main
kvcache-ai/Mooncake"] -->|"fetch"| L["本地仓库"] L -->|"checkout -b feature"| B["功能分支"] B -->|"实现 / 测试 / commit"| B U -->|"rebase / merge latest main"| B B -->|"push"| F["个人 / 公司 Fork"] F -->|"Pull Request"| U classDef good fill:#ddf5e3,stroke:#32a852,stroke-width:3px,color:#132; classDef core fill:#dfecff,stroke:#3973c6,stroke-width:2px,color:#123; class B,F good; class U,L core;
推荐的提交前命令链
如果 PR 已经公开,再改写历史要谨慎;个人 fork 必须改写时使用 --force-with-lease,不要无保护地 force push。
提交前代码状态
code_format.sh、pre-commit、codespell 等本地检查通过。CI:什么时候触发,哪些状态必须关注
触发方式
pull_request / workflow 条件启动。.github/workflows/*,除了 CI 自身通过,还可能额外触发 .github 的 CODEOWNER 审批要求。我们实际遇到的 Native UB CI
这证明的是 mock UB 构建、链接、TENT 单测和 CI matrix 正常,不等价于真实 UB 网卡、liburma 双节点数据传输已经验证。
| Check / 状态 | 应该如何判断 | 处理原则 |
|---|---|---|
| Build & Test / relevant matrix | 需要通过 | 与本次改动相关的构建和测试 job 是最基本的合入门槛。 |
| CI Gate / Required checks | 需要通过 | 最终以 PR merge box / Required 标记为准;Required 失败不能当作“环境问题”直接忽略。 |
| format / pre-commit / codespell | 应通过 | 尽量本地提前解决,避免把纯格式问题留给 CI。 |
| Qoder Review = skipped | 不需要变成 success | 我们遇到的 fork / cross-repo PR 中 Qoder 会被设计为 skipped;skipped 本身不是 CI failure。 |
| 条件不匹配的 matrix job = skipped | 正常 | 平台、路径或 feature 条件未命中时可以跳过,不能把 skipped 与 failed 混为一谈。 |
| Codecov / 外部分析 | 看是否 Required | 覆盖率是重要质量信号,但是否阻塞合并取决于仓库当时的 Required 配置;不能只看颜色判断。 |
CODEOWNERS:谁应该 Review
CODEOWNERS 是什么
仓库用路径规则把代码目录映射到负责维护的人。PR 修改某个受保护目录后,GitHub 可以自动请求相应 owner review,并在 branch protection 下要求 owner approval。
怎么决定 @ 谁
.github/CODEOWNERS,按实际修改目录定位 owner。.github/workflows/ci.yml,即使代码 reviewer 已经 Approve,也可能因为 .github 路径本身需要额外 CODEOWNER approval 而暂时无法合并。| 路径 | 当前 main 中的 CODEOWNERS 示例 | 使用方式 |
|---|---|---|
| /mooncake-store | @ykwd @stmatengss @XucSh @YiXR |
Store、Warmup、HotCache、SSD allocation / deletion 相关改动优先看这里。 |
| /mooncake-transfer-engine | @alogfans @doujiang24 @chestnut-Q |
Classic TE、transport 数据面、RDMA/UB 基础能力。 |
| /mooncake-transfer-engine/tent | @alogfans @doujiang24 @chestnut-Q @staryxchen @00fish0 @dtcccc |
TENT Runtime、Transport 接口、Native UB 等优先从这里选 reviewer。 |
| .github | @stmatengss @ykwd @Ann-1024 @luketong777 |
修改 CI workflow / GitHub 配置时会额外进入这一组 owner 的审批范围。 |
具体 reviewer 以提交 PR 当时的 .github/CODEOWNERS 和 GitHub 自动请求结果为准;上表用于展示我们这套流程里实际涉及的模块映射。
Review:Bot 和 Maintainer 的问题怎么处理
flowchart LR
A["Bot / Maintainer Comment"] --> B{"问题类型"}
B -->|"正确性 / 生命周期"| C["先验证代码路径
补 regression test"]
B -->|"架构 / scope"| D["重新确认模块边界
必要时拆 PR"]
B -->|"风格 / 文档"| E["直接修复并跑 format"]
B -->|"判断不成立"| F["用代码事实 + 测试结果解释"]
C --> G["提交修复"]
D --> G
E --> G
F --> H["回复 rationale"]
G --> H
H --> I["Resolve conversation"]
I --> J["等待 re-review / approval"]
Bot Review
- 适合发现格式、边界条件、潜在生命周期、测试覆盖等问题。
- 不能因为是 bot 就全部接受,也不能因为是 bot 就忽略。
- 先验证它指出的问题是否真的能沿代码路径发生,再决定修还是解释。
- 如果修复,最好用 regression test 证明问题确实被封住。
Maintainer Review
- 更关注模块边界、公共接口、长期维护成本和 PR scope。
- 设计方向被质疑时,不应只局部改代码;需要先确认 maintainer 想要的架构边界。
- 每个 comment 修复后回复“改了什么 + 为什么这样改 + 测试结果”。
- 过时 thread 即使代码已更新,也应确认是否需要手动 Resolve,避免 merge rule 阻塞。
测试要求:从单元正确性到真实硬件
| 层级 | 目的 | 适用示例 | 提交要求 |
|---|---|---|---|
| Unit Test | 验证类/函数局部语义 | allocation strategy、task queue、状态机 | 核心分支必须覆盖 |
| Regression Test | 稳定复现曾发生的 bug | HotCache stale、warmup status error | bugfix 最有价值的测试 |
| Integration Test | 验证跨模块路径 | Store ↔ Master ↔ TE、删除 Fetch/ACK | 接口改变时应补 |
| Concurrency / Repeat | 暴露 race、flaky、资源生命周期 | Warmup burst、Remove vs Fill | 并发问题不能只跑一次 |
| Mock Backend | 让 CI 无硬件也能验证数据面逻辑 | TENT native UB ub-mock | 明确写“mock”,不夸大 |
| Real HW / E2E | 验证驱动、设备、跨节点、真实吞吐 | UB/URMA、RDMA、LOCAL_DISK | 若尚未完成,PR 中明确列为未验证项 |
Issue / PR 模板与提交前检查
Issue / RFC 模板
PR 模板
./scripts/code_format.sh、运行 pre-commit run --all-files、必要时更新文档、补充能够证明修改有效的测试;如果改动超过 500 LOC,模板要求先提交 RFC Issue。提交前 · 代码
无无关 diff;接口命名、错误处理、资源释放与社区已有风格一致;大块临时代码不直接进入 PR。
提交前 · 文档
Issue、PR、代码行为三者一致;不要让 PR 描述停留在旧设计;Module 只勾实际改到的模块。
提交前 · Reviewability
Reviewer 能从 PR 描述快速回答:为什么改、数据路径怎么变、风险在哪里、如何验证、哪些还没验证。
技术实现与社区协作的共同主线
| 工作 | 原问题 | 关键改变 | 系统设计思想 |
|---|---|---|---|
| UB → TENT | 双层 transport / 双模型 | Classic Adapter → Native Task/Slice | 减少架构耦合,统一抽象边界 |
| Warmup | 首次业务请求集中建链 | 初始化阶段有界预建链 | 把不可避免的冷启动成本前移 |
| Bucket Delete | 远端 SSD 删除缺乏可靠协议 | Task → durable tombstone → ACK → GC | 将逻辑正确性与物理回收解耦 |
| HotCache | 异步 Fill 可回写 stale address | 两阶段 invalidate | 封闭并发窗口,保证最终状态一致 |
| RDMA Direct | Requester 二次 Pull | Target RDMA WRITE Push | 减少同步点与中间网络交互 |
| SSD Balance | 随机分配不感知容量 | free-ratio feedback | 用可观测资源状态驱动调度 |