Storforge
编码预计 50 分钟

网络与 RPC:tonic/gRPC

节点之间怎么说话?用 protobuf 定义 forge 的内部协议,tonic 生成服务骨架,顺便搞懂连接池、超时与重试的正确姿势。

学完这节你能做到

  • 用 protobuf 定义 put/get/gossip 协议并保持向后兼容
  • 实现带超时与重试的客户端,理解幂等性前提
  • 解释 HTTP/2 多路复用对存储流量的利弊

节点之间怎么说话:从你 tcpdump 过的流量说起

运维 Ceph 时你抓过 messenger 的包,调过心跳超时,也一定处理过「OSD 之间通信异常」的告警。 分布式存储的第一块积木不是什么高深算法,是节点间通信:消息怎么编码、连接怎么管、 超时了算成功还是失败。这节课把 forge 的内部协议立起来 —— 从 L0 就占着位的 forge-proto 终于开工,同时新增 forge-node:节点守护进程,你以后 systemctl start 的那个东西。

选 gRPC/tonic 而不是自己搞 TCP 私有协议,理由和 L2 选 io_uring 而不是 SPDK 一样: 同一个问题的低一档解法,机制完整,工程量可控。Weka 用 DPDK 自研 RDMA 风格协议, 我们用 tonic 走内核 TCP —— 慢一档,但连接管理、多路复用、流控这些课一节不少。

protobuf 与协议演进:字段号就是盘上格式的兄弟

L1 设计盘上格式时你体会过「格式定了就改不动」的敬畏感。网络协议是同一个问题的镜像: 盘上格式要兼容过去的自己(老数据还在盘上),网络协议要兼容别的节点—— 滚动升级时,新老版本的 forge-node 会同时在线,谁也不能拒收谁的消息。

protobuf 的答案是字段号:每个字段绑定一个永不复用的编号,编码时只写编号不写名字。 新版本加字段,老节点解码时跳过不认识的编号;老版本缺字段,新节点读到默认值。

syntax = "proto3";
package forge.v1;

message PutBlobRequest {
  string key = 1;
  bytes  data = 2;
  uint32 crc32c = 3;   // 端到端校验:客户端算,服务端验,复用 forge-util
}

message PutBlobResponse {
  uint64 cluster_epoch = 1;  // 服务端当时的集群配置版本,placement 课兑现
}

message GetBlobRequest {
  string key = 1;
}

message GetBlobResponse {
  bytes  data = 1;
  uint32 crc32c = 2;
}

service BlobService {
  rpc PutBlob(PutBlobRequest) returns (PutBlobResponse);
  rpc GetBlob(GetBlobRequest) returns (GetBlobResponse);
}

// 成员管理协议,下一课实现;先把字段号占住
service GossipService {
  rpc Exchange(GossipRequest) returns (GossipResponse);
}

协议演进的三条铁律,和盘上格式的版本号规则一一对应:

  1. 字段号只增不改不复用。删掉的字段用 reserved 3; 立墓碑,防止后人误用
  2. 只加字段,不改类型。把 uint32 改成 uint64 在 proto3 里有些组合恰好兼容,但别赌
  3. 语义变更就加新字段,老字段留着,让新代码同时认两种 —— 和盘上格式升级一个套路

tonic 服务端与客户端骨架

forge-proto 用 build script 在编译期把 .proto 生成 Rust 代码:

// crates/forge-proto/build.rs
fn main() -> Result<(), Box<dyn std::error::Error>> {
    tonic_build::configure()
        .compile_protos(&["proto/forge.proto"], &["proto"])?;
    Ok(())
}

服务端在新 crate forge-node 里,把 L1 的 BlobStore trait 挂到网络上 —— L2 做过的「同步核心 + 异步外壳」改造在这里兑现:

use forge_proto::forge::v1::blob_service_server::BlobService;
use forge_proto::forge::v1::{PutBlobRequest, PutBlobResponse, GetBlobRequest, GetBlobResponse};
use tonic::{Request, Response, Status};

pub struct BlobSvc<S> {
    store: std::sync::Arc<tokio::sync::Mutex<S>>,
}

#[tonic::async_trait]
impl<S: forge_store::BlobStore + Send + 'static> BlobService for BlobSvc<S> {
    async fn put_blob(
        &self,
        req: Request<PutBlobRequest>,
    ) -> Result<Response<PutBlobResponse>, Status> {
        let msg = req.into_inner();
        // 红线代码:先验校验和再落盘,网络上翻转的位不许进引擎
        if forge_util::crc32c_of(&msg.data) != msg.crc32c {
            return Err(Status::data_loss("crc32c mismatch, refusing to write"));
        }
        let mut store = self.store.lock().await;
        store
            .put(msg.key.as_bytes(), &msg.data)
            .map_err(|e| Status::internal(e.to_string()))?;
        Ok(Response::new(PutBlobResponse { cluster_epoch: 0 }))
    }

    async fn get_blob(
        &self,
        req: Request<GetBlobRequest>,
    ) -> Result<Response<GetBlobResponse>, Status> {
        // 读路径对称:从引擎读出后重算 crc 再回包
        // ...
        todo!()
    }
}

注意错误的翻译:引擎的 StoreError 到网络边界要映射成 gRPC 的 Status。 crc 不匹配用 data_loss,key 不存在用 not_found,盘写失败用 internal —— 客户端靠这个分类决定重试还是报警,别全都 internal 糊过去。

超时、重试、幂等:分布式第一课

单机上函数调用只有两种结果:成功或失败。跨网络多了第三种:不知道。 请求超时,可能是没送到,也可能是服务端已经写成功了、只是响应丢了 —— 你在运维里早见过这个:命令超时,重跑之前先查一眼上次到底执行没执行。

代码里的对策是一套组合拳,每个数字都要有出处:

  • 连接超时 1s:同机房 TCP 握手是亚毫秒级,1s 还连不上就是节点或网络有病,别等
  • 请求超时按数据量算:1MiB 对象在 10GbE 上传输约 1ms,加上服务端 fsync 几毫秒, 给 500ms 已经宽到荒唐 —— 超过它,等待的代价高于重试
  • 只重试可重试的错误:超时、unavailable 重试;data_loss、参数错误重试一万次也是错
  • 指数退避:50ms 起步,每次翻倍,上限 1s,最多 3 次 —— 防止全体客户端同步重试打死刚恢复的节点(你见过重试风暴)

而这一切的前提是幂等:重试等于同一个请求执行两次,结果必须一样。 PutBlob 按 key 覆盖写,天然幂等;将来要是加「追加写」这种非幂等接口, 就得带上客户端请求 ID 做去重 —— 现在协议里没有,是因为我们刻意让每个 RPC 都幂等。

×超时不等于失败

客户端超时后,那个请求可能仍在服务端排队,几百毫秒后成功落盘。 如果你的重试逻辑假设「超时 = 没写进去」,后面叠加删除或改写操作时就会出现乱序覆盖。 幂等 put 扛得住这种重放,这是我们把幂等当协议设计前提、而不是优化项的原因。

Checkpoint单选

forge 集群滚动升级到一半,新版本节点把 PutBlobRequest 里废弃的字段号 3 复用成了别的含义(bytes 改成了 flags),会发生什么?

大对象传输:流式接口与 zero-copy 的取舍

gRPC 单条消息默认上限 4MiB,而且整条消息要在内存里攒齐才能解码 —— put 一个 1GiB 对象显然不能塞一条消息。答案是流式 RPC: rpc PutBlobStream(stream PutChunk) returns (PutBlobResponse), 客户端切成 1MiB 一片连续发,服务端边收边写,内存占用恒定。

顺便把 HTTP/2 多路复用对存储流量的利弊说清,这是面试和设计评审都躲不开的问题:

  • :一条 TCP 连接跑多个并发流,连接数从「客户端数 × 节点数 × 并发」降到「客户端数 × 节点数」, 连接管理和 keepalive 成本大降
  • :TCP 层的队头阻塞 —— 丢一个包,这条连接上所有流一起等重传; 大对象数据流还会挤占同连接上的小控制消息
  • 对策:控制面和数据面用不同的连接(甚至不同端口)。你在 Ceph 里配 public network 和 cluster network 分离,是同一个思想的物理版

至于 zero-copy:prost 解码时会把 bytes 字段拷贝进新分配的 Vec,一次 1MiB 拷贝 在现代 CPU 上约 20~50µs,而网络往返是毫秒级 —— 这份拷贝暂时不值得优化。 Weka 用 DPDK 把这层也省了;我们记下这笔账,L5 调优时再看值不值得追。

AI 结对:让 forge 第一次开口说话pair with ai

这节课的产出是后面五节课的地基,按三道关认真 review:重点读错误映射 (每个 StoreError 变成了哪种 Status?有没有全部糊成 internal?)和重试循环 (哪些错误码会触发重试?data_loss 重试了就是 bug)。 最后问 AI 一个问题:如果 put 成功但响应丢了,客户端重试会发生什么 —— 答案应该落在幂等上。

在 forge workspace 里完成节点间 RPC 的第一步,分四块:

1. forge-proto:新增 proto/forge.proto,package forge.v1,定义 BlobService(PutBlob/GetBlob,请求响应字段见下)与流式 PutBlobStream(stream 请求,1MiB 分片);PutBlobRequest 含 key(string)、data(bytes)、crc32c(uint32);加 build.rs 用 tonic-build 生成代码,依赖 tonic、prost;
2. forge-node:新建 bin crate,实现 BlobService:put 前用 forge-util 的 crc32c_of 校验,不匹配返回 Status::data_loss;存储后端接 forge-store 的 BlobStore trait,用 tokio Mutex 包装;启动参数 --listen 与 --data 用 clap 解析;
3. forge-cli:新增子命令 blob put <ADDR> <KEY> <FILE> 和 blob get <ADDR> <KEY>,客户端设置连接超时 1s、请求超时 500ms,对 unavailable 与超时做最多 3 次指数退避重试(50ms 起,上限 1s),其他错误直接报出;
4. 集成测试:tokio 测试里起服务端,put 后 get 校验内容一致;再测 crc 故意写错时服务端拒绝。

生产路径不许 unwrap。完成后跑 cargo clippy --workspace 和 cargo test --workspace,贴结果。

小结

  • 字段号是网络协议的「盘上格式」:只增不改不复用,滚动升级时新老节点才能互认
  • tonic 骨架把 L1 的 BlobStore 挂上网络;错误在边界处翻译成可分类的 Status,别糊成一团
  • 分布式第一课:超时是第三种结果(不知道),重试的前提是幂等,退避防止重试风暴
  • 大对象走流式分片;控制面和数据面分连接,对付 HTTP/2 的队头阻塞
  • 节点能说话了,下一课解决更难的问题:怎么知道对方还活着