分布式基础:RPC 和远程调用
系列入口:分布式阅读路径。 前置:操作系统、网络、并发与存储。 下一篇:状态机、分片、路由与元数据。
本文要解决的问题
一个请求超时后,客户端往往不知道服务端是否已经修改了状态。本文围绕这个不确定性,解释 deadline、重试预算、幂等与操作 ID,并用一个可运行案例观察“已提交但响应丢失”之后该怎样重试。
RPC 是什么?
RPC 是 Remote Procedure Call,远程过程调用。
它的目标是让调用远程服务像调用本地函数一样方便。例如本地代码看起来像:
metadata = get_block_group_metadata(block_group_id)
但实际执行时,可能发生了:
1. 客户端把 block_group_id 序列化成请求。
2. 请求通过网络发送到 metadata server。
3. metadata server 查找元数据。
4. metadata server 把结果序列化成响应。
5. 响应通过网络返回客户端。
6. 客户端反序列化得到 metadata。
RPC 框架屏蔽了很多细节,但不能消除远程调用的本质问题:
网络可能超时。
请求可能丢失。
响应可能丢失。
服务端可能崩溃。
服务端可能执行成功但客户端不知道。
服务端可能很慢。
客户端可能重试。
重试可能导致重复执行。
所以 RPC 的重点不是“怎么像本地调用一样写代码”,而是“知道它绝对不是本地调用”。
请求/响应模型
最常见的 RPC 是请求/响应模型:
client -> request -> server
client <- response <- server
例如:
get_block_group(block_group_id, epoch)
正常路径是:
1. client 发送请求。
2. server 收到请求。
3. server 执行读取。
4. server 返回响应。
5. client 得到结果。
但真实系统里,异常路径更多:
client 发不出去。
request 在网络中丢失。
server 收到了,但还没执行就崩溃。
server 执行成功了,但 response 丢了。
server 执行太慢,client 超时。
client 超时后重试,server 收到两次请求。
所以设计 RPC 接口时,不能只考虑正常返回,还要考虑调用方不知道服务端到底执行到哪一步。
序列化和反序列化
远程调用要跨进程、跨机器传输数据,所以内存里的对象不能直接发送,必须先编码成字节流。
这个过程叫序列化:
object -> bytes
接收方再把字节流恢复成对象:
bytes -> object
常见序列化方式包括:
JSON
Protocol Buffers
Thrift
FlatBuffers
MessagePack
不同方式有不同取舍:
JSON:可读性好,但体积大、解析慢。
Protobuf:体积小、速度快,适合服务间通信。
Thrift:常用于跨语言 RPC。
FlatBuffers:可减少反序列化开销,适合性能敏感场景。
在 KV Cache 场景中,一般不会把巨大的 K/V tensor 直接用普通 JSON 传输。RPC 请求里更常传:
block_group_id
location
epoch
token_range
size_bytes
checksum
metadata
真正的大块数据可能通过专门的数据通道传输,例如:
RDMA
对象存储
共享内存
文件
GPU Direct
专门的 streaming channel
因此要区分两类数据:
控制面数据:元数据、状态、命令,适合 RPC。
数据面数据:大块 KV Cache 内容,通常需要更高效的数据传输方式。
这个区分很重要。不要把所有东西都塞进普通 RPC,否则很容易让控制面被大数据传输拖垮。
远程调用为什么会失败?
本地函数调用失败时,通常能比较明确地知道发生了什么。
远程调用失败时,最大的问题是:客户端经常不知道服务端是否执行过。
例如:
evict_block_group(bg_1) 超时
可能发生了四种情况:
1. 请求没发到服务端。
2. 服务端执行了,但响应丢了。
3. 服务端正在执行,还没返回。
4. 服务端执行失败。
从客户端视角看,这些情况可能都表现为:
timeout
但它们的后果完全不同。
如果客户端简单重试:
evict_block_group(bg_1)
evict_block_group(bg_1)
服务端可能执行两次。对于查询类接口,这通常没问题;但对于修改状态的接口,重复执行可能带来错误。
所以分布式接口设计必须考虑:
这个接口能不能重试?
重复调用是否安全?
如果服务端已经执行成功,客户端重试应该返回什么?
如果请求过期了,服务端是否应该拒绝?
可运行案例:已提交,但客户端没有收到结果
下面是为解释失败语义构造的最小示例。它用 SQLite 保存一个 BlockGroup 的预留计数,使用子进程退出来制造“没有响应”;没有建立真实网络连接,也不测量 RPC 性能。
先冻结请求与不变量
客户端希望为 bg_1 预留一个单位,第一次请求和所有重试都使用同一个操作 ID:
{"operation_id": "reserve-1", "block_group": "bg_1", "epoch": 7, "slots": 1}
初始状态是 epoch=7, reserved=0。希望保持的不变量是:同一逻辑操作即使重试,也只增加一次计数。 不同操作使用不同 ID;相同 ID 携带不同参数时必须拒绝,不能默默当成重试。
第一版为什么会重复执行?
第一版只在一个事务里增加计数,然后返回成功:
| 步骤 | 服务端持久状态 | 客户端观察 |
|---|---|---|
| 收到首次请求 | reserved=0 | 等待 |
| 增加计数并提交事务 | reserved=1 | 仍未收到响应 |
| 服务端在返回前退出 | reserved=1 | 没有结果,无法确定是否提交 |
| 新进程收到相同请求 | reserved=2 | 收到成功,但副作用已重复 |
这个反例中的每次数据库更新都是原子的。问题在于接口没有识别“两次调用其实是同一个逻辑操作”,所以单次更新正确仍不足以保证重试正确。
第二版:把操作结果与副作用一起提交
增加一张操作表,保存 operation_id、规范化请求参数和原始结果。一次请求按下面的事务执行:
BEGIN IMMEDIATE
查询 operation_id
已存在:参数相同则返回保存的结果;不同则拒绝
不存在:校验当前 epoch,更新 reserved
插入 operation_id、请求参数和本次结果
COMMIT
返回结果
操作结果与计数在同一个事务里提交:
- 提交前退出:计数修改与操作记录一起回滚,重试可以真正执行一次。
- 提交后、返回前退出:计数和操作记录都保留,新进程查到记录后返回原始结果。
- 并发重试:示例的
BEGIN IMMEDIATE串行化写事务,后续请求在同一检查路径看到已保存的 ID。
如果先增加计数并提交,再单独记录 ID,那么两次提交之间退出仍会重复执行。只在 Python 字典里记住 ID 也不够:进程重启会丢失去重信息。
下载、运行与实际结果
完整脚本:rpc-idempotency.py 只依赖 Python 3.10+ 标准库中的 SQLite 支持。以下命令将它下载到新建临时目录,脚本自己创建并清理临时数据库:
rpc_demo_dir=$(mktemp -d)
curl --fail --location https://zqwiki.cn/examples/rpc-idempotency.py \
--output "$rpc_demo_dir/rpc-idempotency.py"
python3 "$rpc_demo_dir/rpc-idempotency.py"
2026-09-07 使用 Python 3.14.7 / SQLite 3.53.4 在本地运行该脚本,得到以下结果。每次请求都由新的子进程执行;两种故障点使用 os._exit 直接退出,避免正常清理替我们完成事务处理:
unsafe / response lost + retry: reserved=2 (duplicate effect)
safe / restart + same ID: reserved=1 (saved result replayed)
safe / crash before commit: reserved=0; retry -> reserved=1
safe / same ID + changed parameters: rejected; reserved=1
safe / new ID + stale epoch: rejected; reserved=1
safe / 4 concurrent retries: reserved=1
PASS: 6 failure/contract scenarios
这些检查支持的是单机数据库边界内的副作用去重。它们没有验证磁盘损坏、断电、跨数据库事务或分布式共识,也没有证明网络能提供“恰好一次投递”。SQLite 原子提交依赖的环境条件见 Atomic Commit。
还需要明确三个接口边界:操作表的保存时间必须覆盖允许的重试窗口;多租户服务的去重键需要包含相应作用域;返回缓存结果表示该操作当时成功,不代表 BlockGroup 的当前状态仍与当时相同。真正的 GPU 资源释放或跨节点迁移若无法纳入这一个事务,就需要额外的状态机、补偿或查询机制。
Timeout:超时
远程调用必须设置超时。
没有超时会导致:
请求一直等待。
线程或协程被占用。
连接池被耗尽。
上游请求堆积。
故障扩散到其他模块。
例如 worker 读取远端 block group:
get_block_group(block_group_id, epoch)
如果远端节点卡住,而本地没有 timeout,那么 decode 请求会一直等待,进而拖慢整个 batch。
超时时间怎么设置?
超时时间不能随便写一个固定值,要结合调用路径和业务目标。
需要考虑:
这个调用在不在在线请求关键路径上?
调用失败后是否可以重试?
重试是否会造成更大延迟?
下游服务正常 P99 是多少?
上游请求整体超时时间是多少?
一个基本原则是:
下游 timeout 不能超过上游剩余时间。
例如一次推理请求整体最多允许 2 秒,而某个远程 cache 读取已经消耗了 1.8 秒,那么这个 RPC 不应该再设置 1 秒超时。
超时不是取消
客户端 timeout 只代表客户端不等了,不保证服务端应用工作已经停止。gRPC Deadlines 会传播取消状态,但服务端仍需让自己启动的工作检查取消并退出;已经提交的副作用也不会因此自动回滚。
例如:
migrate(block_group_id=bg_1, src=A, dst=B, epoch=10)
客户端 500ms 后超时,但服务端可能还在迁移。客户端如果立刻发起新的迁移任务,就可能和旧任务冲突。
因此服务端需要能识别同一个操作,客户端也需要能查询操作状态。
Retry:重试
重试可以提高成功率,但也可能放大故障。
适合重试的情况:
临时网络抖动。
连接被对端关闭。
服务端短暂过载。
请求没有明显副作用。
接口是幂等的。
不适合盲目重试的情况:
非幂等写操作。
下游已经严重过载。
上游剩余时间不足。
请求体很大,重试成本高。
重试会触发重复迁移或重复释放。
重试风暴
如果大量客户端同时发现下游变慢,然后一起重试,就会形成重试风暴。
这会让本来已经变慢的服务更慢。
常见缓解方式:
限制最大重试次数。
指数退避。
增加随机抖动 jitter。
只对部分错误重试。
使用重试预算 retry budget。
下游过载时快速失败。
例如:
第一次失败后等待 10ms。
第二次失败后等待 20ms。
第三次失败后等待 40ms。
每次等待时间加一点随机抖动。
重试要带请求 ID
重试时最好带上唯一请求 ID 或操作 ID:
operation_id = "migrate-bg_1-epoch_10-A-to-B"
服务端可以用它识别重复请求:
如果操作已经成功,直接返回成功结果。
如果操作正在执行,返回 in_progress 或等待。
如果操作失败,返回失败原因。
这比每次都创建一个新操作安全得多。
幂等
幂等的意思是:同一个操作执行一次和执行多次,最终效果相同。
例如:
set_state(bg_1, EVICTED)
执行多次结果仍然是 EVICTED。
但下面这种操作不是天然幂等:
ref_count += 1
如果客户端重试两次,ref_count 就可能多加一次。
查询类接口
查询类接口通常天然幂等:
get_metadata(block_group_id)
get_location(block_group_id)
get_state(block_group_id)
重复查询不会改变系统状态。
写入类接口
写入类接口需要特别设计。
例如淘汰接口不应该只写成:
evict(block_group_id)
更好的形式是:
evict(block_group_id, epoch, operation_id)
这样服务端可以判断:
epoch 是否仍然是当前版本?
operation_id 是否已经执行过?
block group 当前状态是否允许 evict?
如果已经是 EVICTED,是否可以直接返回成功?
pin/unpin 的幂等设计
pin 和 unpin 很容易出错。
如果接口是:
pin(block_group_id)
unpin(block_group_id)
那么重试可能导致计数错误。
更好的方式是带上 request_id:
pin(block_group_id, request_id)
unpin(block_group_id, request_id)
服务端维护:
block_group_id -> pinned_request_set
这样:
同一个 request_id pin 多次,只算一次。
同一个 request_id unpin 多次,只释放一次。
这种设计比简单的 pin_count += 1、pin_count -= 1 更适合分布式重试场景。
migrate 的幂等设计
迁移接口也需要幂等:
migrate(block_group_id, src, dst, epoch, operation_id)
服务端处理时应该检查:
block_group_id 是否存在?
epoch 是否匹配?
当前 location 是否仍然是 src?
目标是否已经是 dst?
operation_id 是否已经执行过?
当前状态是否允许迁移?
如果客户端重试同一个 operation_id,服务端不应该启动两次迁移,而应该返回同一个操作的状态。
限流
限流是为了保护系统,避免请求量超过服务能力。
常见限流维度:
按服务限流。
按用户或租户限流。
按接口限流。
按 block group 迁移流量限流。
按远端节点限流。
在 KV Cache 系统里,限流尤其重要,因为某些操作很重:
远端读取大 block group。
跨节点迁移 cache。
从 SSD 拉取冷数据。
批量淘汰和释放。
如果不限制,后台迁移任务可能抢占在线 decode 的网络带宽和存储带宽。
常见限流算法
常见算法包括:
计数器:固定窗口内限制请求数。
滑动窗口:比固定窗口更平滑。
漏桶:以固定速率处理请求。
令牌桶:允许一定突发,但整体速率受限。
对于在线服务,令牌桶很常见:
系统按固定速率生成 token。
请求要先拿到 token 才能执行。
token 桶允许短时间突发。
桶空了就等待或拒绝。
KV Cache 迁移可以设计独立的 token:
migration_bytes_token
remote_read_qps_token
ssd_read_iops_token
这样可以限制后台任务,不让它们影响在线请求。
熔断
熔断是为了避免持续调用已经异常的下游。
如果某个服务持续超时,客户端继续打请求只会浪费资源,并加重下游压力。
熔断器通常有三种状态:
CLOSED:正常调用。
OPEN:熔断,直接失败。
HALF_OPEN:半开,放少量请求探测恢复情况。
例如 remote cache manager 连续超时,客户端可以进入 OPEN 状态:
后续请求不再访问这个远端节点。
调度器把请求转移到其他节点。
必要时走重新计算或降级路径。
一段时间后进入 HALF_OPEN:
只放少量请求访问远端。
如果成功率恢复,再切回 CLOSED。
如果仍然失败,继续 OPEN。
熔断的目标不是解决下游故障,而是限制故障扩散。
连接池
RPC 通常会复用连接,而不是每次请求都新建连接。
连接池的作用是:
减少 TCP 建连成本。
减少 TLS 握手成本。
控制并发连接数量。
复用已有连接提高吞吐。
但连接池也可能成为瓶颈。
需要关注:
连接池大小。
每条连接上的并发请求数。
连接是否健康。
空闲连接是否回收。
连接上的请求是否出现队头阻塞。
例如某个 worker 到 metadata server 的连接池太小,所有请求都排队等连接,即使 metadata server 本身很空,P99 也会变差。
队头阻塞
队头阻塞是指前面的慢请求挡住后面的请求。
如果多个 RPC 共用一条连接,而协议或实现不能很好地并发处理响应,就可能出现:
请求 A 很慢。
请求 B 本来很快。
但 B 排在 A 后面,必须等待。
这会显著影响尾延迟。
所以性能敏感服务要关注连接池、协议多路复用和请求排队情况。
控制面和数据面
在 KV Cache 系统里,需要区分控制面和数据面。
控制面负责:
元数据查询。
状态更新。
调度命令。
迁移任务创建。
pin/unpin。
evict。
数据面负责:
真实 KV Cache 数据传输。
大块 tensor 复制。
GPU HBM 和 CPU Memory 之间搬运。
跨节点传输 block group。
SSD 读写。
RPC 更适合控制面。数据面如果也完全依赖普通 RPC,很容易出现:
大响应占满连接。
控制请求被大数据传输阻塞。
序列化开销过高。
内存拷贝过多。
P99 变差。
所以一种常见设计是:
RPC 负责发命令和返回元数据。
数据传输走专门的数据通道。
RPC response 里返回数据位置、token、offset、checksum。
客户端再通过数据通道拉取或写入。
例如:
prepare_get_block_group(block_group_id, epoch)
-> 返回 remote_addr、size、checksum、transfer_token
transfer_data(remote_addr, size, transfer_token)
-> 传输真实 KV 数据
这样可以减少控制面被大数据传输拖慢的风险。
KV Cache 接口设计示例
下面用几个接口把前面的概念串起来。
查询元数据
get_metadata(block_group_id)
特点:
查询类接口,天然幂等。
需要 timeout。
可以短重试。
客户端可以缓存结果,但要关注 epoch。
返回值可以包括:
block_group_id
location
state
epoch
size_bytes
token_range
读取 block group
get_block_group(block_group_id, epoch, request_id)
需要考虑:
epoch 不匹配时拒绝或返回最新元数据。
block group 正在 MIGRATING 时如何处理。
远端读取是否超时。
是否允许重试。
大数据是否走数据面传输。
如果 block group 当前状态是 MIGRATING_OUT,服务端可以选择:
等待迁移完成。
返回 RETRY_LATER。
返回新 location。
拒绝读取并要求客户端刷新元数据。
淘汰 block group
evict(block_group_id, epoch, operation_id)
服务端应该检查:
epoch 是否匹配。
pin_count 是否为 0。
ref_count 是否为 0。
当前 state 是否允许 evict。
operation_id 是否已经执行过。
如果已经被淘汰:
state == EVICTED
那么重复调用可以直接返回成功。这就是幂等。
迁移 block group
migrate(block_group_id, src, dst, epoch, operation_id)
这个接口要特别小心,因为迁移是长操作。
可以把它拆成两类接口:
start_migrate(...)
get_migrate_status(operation_id)
这样客户端超时后,不必盲目重试迁移,而是先查询这个 operation_id 的状态。
可能的状态:
PENDING
RUNNING
SUCCEEDED
FAILED
CANCELLED
pin/unpin
pin(block_group_id, request_id, epoch)
unpin(block_group_id, request_id, epoch)
服务端可以维护:
pinned_request_set
这样 pin 和 unpin 都可以做到幂等:
同一个 request_id 重复 pin,不重复增加 pin_count。
同一个 request_id 重复 unpin,不重复减少 pin_count。
这对超时重试非常重要。
常见错误设计
把远程调用当成本地调用
错误做法:
result = remote_call()
use(result)
但没有考虑:
timeout 怎么办?
失败是否重试?
重试是否幂等?
下游过载怎么办?
调用是否会卡住线程?
写接口不带版本
错误做法:
evict(block_group_id)
问题是服务端不知道客户端看到的是不是旧状态。
更好的方式:
evict(block_group_id, epoch, operation_id)
pin/unpin 只操作计数
错误做法:
pin_count += 1
pin_count -= 1
问题是请求超时重试会导致计数不准。
更好的方式:
pin(block_group_id, request_id)
unpin(block_group_id, request_id)
由服务端维护请求集合,再计算 pin_count。
大数据走普通 RPC
错误做法:
get_block_group(...)
-> response 里直接塞几百 MB KV Cache
问题是:
序列化开销大。
内存拷贝多。
控制面连接被占满。
小请求被大响应阻塞。
P99 变差。
更好的方式是控制面和数据面分离。
如何排查 RPC 问题?
RPC 问题通常需要看指标和日志。
指标
常见指标包括:
rpc_qps
rpc_error_rate
rpc_timeout_count
rpc_retry_count
rpc_latency_p50
rpc_latency_p99
rpc_inflight_requests
rpc_queue_length
connection_pool_in_use
connection_pool_wait_time
circuit_breaker_state
rate_limited_count
KV Cache 相关 RPC 还可以看:
get_metadata_latency
get_block_group_latency
migrate_start_latency
migrate_status_latency
evict_latency
pin_latency
unpin_latency
remote_read_timeout_count
日志
日志里最好带上:
request_id
operation_id
block_group_id
src_node
dst_node
epoch
rpc_method
timeout_ms
retry_count
error_code
latency_ms
如果没有这些字段,排查问题时很难判断:
这是第一次请求还是重试?
服务端有没有执行过?
客户端看到的是哪个 epoch?
迁移操作是否重复提交?
超时发生在哪个节点?
面向 KV Cache 的 RPC 检查清单
看组内代码或设计文档时,可以用下面的问题检查:
1. 这个远程调用是否在在线请求关键路径上?
2. 是否设置了 timeout?
3. timeout 是否小于上游剩余时间?
4. 失败后是否 retry?
5. retry 是否有最大次数、退避和 jitter?
6. 被 retry 的接口是否幂等?
7. 写接口是否带 epoch/version?
8. 长操作是否有 operation_id?
9. pin/unpin 是否按 request_id 幂等?
10. 大块 KV 数据是否和控制面 RPC 分离?
11. 是否有限流保护下游?
12. 下游连续失败时是否熔断?
13. 连接池是否可能成为瓶颈?
14. 日志里是否包含 request_id、operation_id、block_group_id 和 epoch?
把接口契约写完整
选一个 evict、migrate 或 pin 接口,写出相同操作 ID 的重试结果、参数变化时的拒绝方式、旧 epoch 的处理,以及服务端重启后的查询入口。用前面的响应丢失案例检查每个提交点,而不是只验证一次正常返回。
参考与对照
- gRPC Deadlines:区分客户端停止等待、RPC 取消通知和应用主动停止后台工作。
- gRPC Retry:核对可重试状态、退避与调用提交点;框架重试不替代业务副作用去重。
- SQLite Atomic Commit:对照示例中状态修改与操作结果的原子提交,以及它对文件系统的假设。