起点:TCP 字节流上的三个问题

远程过程调用(RPC)的目标朴素:让一个进程像调用本地函数一样,调用另一个进程的函数。传输层通常复用一条 TCP 连接,但 TCP 提供的是一条没有边界的字节流,直接把函数调用架在上面会立刻撞上三个问题。

  1. 没有消息边界。TCP 不区分”这是一条完整消息”。发送方写两次、接收方可能一次读完,也可能分三次读完,边界信息在字节流里丢失。
  2. 一条连接串行排队。若每个请求都等前一个的响应才能发下一个,前一个慢请求会卡住后面所有请求(这种前一个未完成、后一个被阻塞的现象称为队头阻塞),连接的吞吐被最慢的请求拖垮。
  3. 一次调用只有一次往返。普通调用是”请求—响应”两个帧,但有些场景需要一端持续推送(例如导入镜像时不断上报进度),单次往返表达不了。

这三个问题分别对应三个机制:帧化切消息边界,多路复用让一条连接并发跑多个请求,流式让一次调用承载连续多帧。本文先逐一拆解这三个机制,再用一个四层模型说明它们如何协作,并解释为什么流式传输本质上与框架本身无关。

核心概念

后续章节会频繁引用以下术语,先做统一定义。术语保留英文形式(如 StreamID、readLoop),因为它们直接对应代码里的标识符。

术语定义
字节流(Byte stream)TCP 提供的语义:一串连续字节,没有”一条消息”的边界。发送顺序保证,但读写次数不保证对齐。
帧(Frame)协议自行切出的最小单位,一条消息对应一个帧。帧由长度前缀、流标识、类型、负载组成。
StreamID帧头里的流标识,同一逻辑流的所有帧共享一个 StreamID,是多路复用的钥匙。
逻辑流(Stream)一条 TCP 上复用出的独立通道,由 StreamID 区分。一条逻辑流可以反复收发多个帧。
半关闭(Half-close)只关闭一个方向:发送方声明”我发完了”,但接收方向继续有效。区分本地收尾与对端结束两类。
readLoop一个长期运行的 goroutine,不停从 TCP 读帧并按 StreamID 投递,是多路复用的分派核心。

字节流、帧、StreamID 三者的关系是本文的基线:TCP 提供字节流,帧化把它切成消息,StreamID 把消息归到各自的逻辑流

机制总览

一次完整的 RPC 由三个机制叠加而成,下方的总览图给出它们的协作位置。

   Sender                                 Receiver
   ──────                                 ────────

   app calls RPC ──► [Frame] ──┐
   (encode args)               │  length-prefix framing  (framing)

   ═════════════════ one TCP byte stream ═════════════════
   ...[len][id][type][payload][len][id][type][payload]...   (byte stream, no boundary)
                               │ ReadFrame() + readLoop  (multiplexing)

                    dispatch by StreamID into per-stream queues


                    per logical stream: Send / Recv repeatedly  (streaming)
                    + half-close (CloseSend / Close)

机制一在发送侧把消息打成带长度前缀的帧,机制二在接收侧用 readLoop 按 StreamID 分发,机制三在已建立的逻辑流上支撑连续收发。下面逐层展开。

机制一:帧化(长度前缀切消息边界)

既然 TCP 不提供消息边界(见上文起点的第一个问题),协议必须自己切。解决办法是经典的长度前缀:发送方先写一个长度字段,再写负载;接收方先读长度,再按长度精确读够负载。既不丢字节,也不会越界读到下一条消息的字节。

下面的帧格式来自一个教学协议 minirpc——作者维护的从零写的极简 RPC 框架(零外部依赖、负载用 JSON),灵感来自 containerd 的内部 RPC 协议 ttrpc。需要先点明的是,这是 minirpc 自定义的格式,并非某种标准答案:真实的 RPC 协议各有各的线上格式(gRPC 走 HTTP/2,ttrpc 有自己的长度前缀结构),但下文讲的三个机制——帧化、多路复用、流式——是共通的。读懂 minirpc 这一种,换到别的协议只是换一层工程外壳。minirpc 的长度字段用大端序(网络协议惯例,收发双方约定同一字节序即可)。

┌────────────┬───────────┬────────┬──────────────┐
│ Length     │ StreamID  │ Type   │ Payload      │
│ 4 bytes BE │ 4 bytes BE│ 1 byte │ Length-5 B   │
└────────────┴───────────┴────────┴──────────────┘
 Length = bytes after itself (StreamID + Type + Payload), excludes Length

读侧的关键是用”读够 N 字节”的语义读取,而不是裸用 Read。因为 TCP 一次 Read 可能只返回一两个字节(partial read 很常见),必须循环读满。Go 标准库的 io.ReadFull 正好提供这个语义:

func ReadFrame(r io.Reader) (*Frame, error) {
    var lenBuf [4]byte
    if _, err := io.ReadFull(r, lenBuf[:]); err != nil {  // 读够 4 字节长度
        return nil, err
    }
    length := binary.BigEndian.Uint32(lenBuf[:])

    rest := make([]byte, length)                          // 按长度分配
    if _, err := io.ReadFull(r, rest); err != nil {       // 读够剩余部分
        return nil, err
    }
    return &Frame{
        StreamID: binary.BigEndian.Uint32(rest[0:4]),
        Type:     FrameType(rest[4]),
        Payload:  rest[5:],
    }, nil
}

Type 字段标明这条消息的用途:REQUEST(方法调用请求)、RESPONSE(一问一答式调用的最终返回,这类单次请求—响应的调用称为 unary)、DATA(流中的一条数据)、CLOSE(流结束信号,对应读到结尾)。读侧拿到帧后先看 Type 才决定如何处理,这是机制二分发的前提。帧化只保证”一整条消息完整到达”,它不关心负载里装的是什么——这一点会在四层模型里再次出现。

机制二:多路复用(按 StreamID 分发)

帧化解决了边界,但若每条消息都等响应才能发下一条,连接依然是串行的。多路复用的目标是:一条 TCP 上并发跑多个请求,谁先完成谁先回,互不阻塞。

实现它的全部秘密是 1 个读循环加 N 个队列:每个逻辑流开一条自己的入站队列(一个按 StreamID 索引的 channel 映射);一个 readLoop goroutine 不停读帧,看帧的 StreamID,把帧投递进对应流的队列;每个流的 Recv 阻塞读自己的队列,各流互不阻塞。

                  one TCP byte stream

                      ReadFrame()
                           │   readLoop: 1 goroutine, runs forever

          dispatch each frame by its StreamID

        ┌──────────────────┼──────────────────┐
        │                  │                  │
    StreamID=1         StreamID=2         StreamID=3
        │                  │                  │
        ▼                  ▼                  ▼
    [ queue ]         [ queue ]         [ queue ]   (buffered channels)
        │                  │                  │
    waiter A          waiter B          waiter C    (each Recv blocks on its own)

readLoop 的循环体只有几行:一个读循环不断读帧,按 StreamID 投递进对应流的队列。

// readLoop: 1 goroutine, runs forever, dispatches by StreamID
for {
    f := ReadFrame(conn)
    recv := streams[f.StreamID]   // map[StreamID] -> buffered queue
    recv.ch <- f                  // deliver to that stream's queue
}

为什么 0.5 秒的请求能比 1 秒的先回来?因为不同 StreamID 的帧去不同队列,各自被各自的等待者取走,跨流没有先来后到。下面是一次实际运行的时序,三个请求几乎同时发出,但完成顺序由耗时决定,而非发起顺序。

T+0.0s   client issues A(1s), B(2s), C(0.5s)  ──►  server: 3 goroutines sleep
T+0.5s                       ◄──  stream#C RESPONSE   first back, not issue order
T+1.0s                       ◄──  stream#A RESPONSE
T+2.0s                       ◄──  stream#B RESPONSE   last

这里有一个容易混淆的边界:业务标签(如上方的 A/B/C)只是客户端本地的展示标签,并不进入协议负载;真正把帧归位到正确调用的是 header 里的 StreamID。所以哪个 StreamID 对应哪个耗时请求,每次运行都可能不同——这正是多路复用下”完成顺序由耗时决定”的体现。

分发还牵涉一个读写不对称的并发约束。读侧无需加锁,因为只有一个 readLoop 在读;写侧必须加锁串行化,因为 TCP 字节流上两个 goroutine 同时写,两帧字节会交错,对端解出来全是乱码。所以约定”一帧等于一次完整的连续写”,用一把锁串起来。readLoop 把帧投递进队列时还要处理一个边界:流的本地收尾可能在投递前发生,此时不能往已关闭方向写,需用 select 同时监听流的结束信号,已结束就丢弃这一帧。

这里需要厘清多路复用消除的是哪一层的串行。它消除的是连接级串行:客户端在 TCP 连接上阻塞等响应,导致别的请求发不出去;换成在逻辑流上等待后,各流并行推进。它不消除单条逻辑流内部的串行:一条流上若采用”发一帧等一回包”的一问一答模式,该流内部依旧是请求—响应串行,只是不再连累其它流。

multiplexed: N streams share 1 TCP

  stream A ┤ A.req ──wait──► A.resp
  stream B ┤ B.req ──wait──► B.resp
  stream C ┤ C.req ──wait──► C.resp

换句话说,多路复用给单车道加了 N 条车道,而非把每条车道改成不限速。某条车道堵车,车在那条车道里照样排队,但堵车不会蔓延到别的车道。要打破单条流内部的往返串行,需要机制三:让一条逻辑流承载连续多帧,而非一问一答。

机制三:流式(逻辑流上多帧加半关闭)

一旦一条逻辑流不止于一问一答,而是可以反复 SendRecv,流式就出现了:服务端在一整段时间内连续发多个 DATA 帧,客户端边收边处理。这对应镜像导入这类场景——一个操作持续产出进度,而非一次性返回。

流式之所以能成立,关键在于区分两种”结束”。这两种结束不是同一件事,混在一起就会出错。

  • 对端结束发送(远端 EOF):对端发一个 CLOSE 帧,表示”这个方向不会再有数据”。该帧与它前面的 DATA 在同一队列里排队,所以接收方会先按先进先出读完所有 DATA,再在 CLOSE 上收到读到结尾的信号(Go 标准库约定为 io.EOF,表示”正常读到了末尾”)。
  • 本地收尾(本地取消):本端声明”我不再听了”,注销自己的接收队列。它不产生任何协议帧,只是本地拆掉接收通道。

对应到接口上,CloseSend 发一个 CLOSE 帧(协议动作,通知对端),Close 注销本地队列(本地动作,不通知对端)。这种”只关一个方向”的设计称为半关闭:一端可以先把数据推完并声明”我说完了”,同时继续读另一端的回包。这与 HTTP/2 用单独标志位声明某方向流结束(END_STREAM)是同一思想。

半关闭的存在,引出一个更根本的问题:框架究竟负责到哪一层,剩下的交给谁?答案就是下一节的四层模型。

四层模型:为什么流式与框架无关

前面三个机制分开看都很直观,但合起来会出现一个看似矛盾的现象:流式传输时,客户端和服务端是在各自的方法里自己处理 DATA 帧的——框架并没有一个”流式模式”的开关。这是为什么?

把同一时刻一条逻辑流上的活动拆成四层,答案就清晰了。以服务端在一整条流上连续推送进度为例:

time -->

┌─ Layer 1: TCP byte stream ────────────────────────────────────┐
│ one pipe, no message boundary                                 │
│ wire bytes: [len][id=B][DATA][33%][len][id=B][DATA]...        │
│ framework asks only: is this a complete frame?                │
└───────────────────────────────────────────────────────────────┘
                 ReadFrame() slices into frames
                                |
┌─ Layer 2: logical-stream dispatch  (framework's job) ─────────┐
│ readLoop delivers each frame by StreamID into queues          │
│   stream#A queue: [REQUEST] ... [RESPONSE]  main call         │
│   stream#B queue: [DATA][DATA][DATA][CLOSE]  progress         │
│ framework does NOT inspect DATA contents                      │
└───────────────────────────────────────────────────────────────┘
                   stream.Recv() takes a frame
                                |
┌─ Layer 3: framework-recognized signals ───────────────────────┐
│ REQUEST  -> dispatch to handler                               │
│ CLOSE    -> translate to io.EOF                               │
│ identical handling for unary and streaming                    │
└───────────────────────────────────────────────────────────────┘
               DATA frames carried along untouched
                                |
┌─ Layer 4: function-negotiated protocol  (framework-agnostic) ─┐
│ what DATA carries, how many frames, which direction           │
│ = contract between client and server for this method          │
└───────────────────────────────────────────────────────────────┘

关键在第 2 层与第 4 层的分工。第 2 层是框架的职责:它只看 header 里的 StreamID,把帧投递进对应队列,完全不解析 DATA 的内容。第 4 层与框架无关:DATA 里装什么、传几帧、往哪个方向传,全靠客户端和服务端在这个方法上的接口契约。

这就是文章开头那个看似矛盾现象的根源:流式传输不是 RPC 框架的某种”模式”,而是框架抽象天然支持的副作用。框架只要保证”帧完整、有序地送到对应的逻辑流”,剩下的就是双方自己协商的协议——客户端接下来发的 DATA 你按某种格式解,这段协议长什么样,和框架无关。

由此得到本文最核心的判断:逻辑流的建立由框架负责,流上要传的内容由 function 协商,两者是解耦的。一个更直观的说法是——框架修的是路(逻辑流),流式传输是上面跑的车,车的设计是车主(function)的事。用不用流式、流上格式怎么定,都与框架无关。

四层模型套到 gRPC 上

既然”机制共通、外壳不同”,拿这套四层模型去套最主流的生产 RPC——gRPC,能验证这个判断,也顺带澄清一个常见误解:gRPC 的切帧并不是 protobuf 干的

先说约定。gRPC 客户端和服务端的约定由两部分组成:proto 文件声明数据结构和 service 方法(方法名、参数类型、返回类型),HTTP/2 提供承载它的传输。proto 文件只描述”调什么、传什么结构”,不规定字节怎么切。

再看四层各是谁的职责。gRPC 并非重新发明帧格式,而是复用 HTTP/2 现成的两层:

  • 第 1、2 层(帧化 + 多路复用)由 HTTP/2 承担。HTTP/2 自己定义了 9 字节帧头(含长度、stream id、类型字段),并按 stream id 做多路复用——这正是上文机制一、机制二的内容,gRPC 直接复用。所以 gRPC 的”StreamID”就是 HTTP/2 的 stream id,多路复用由 HTTP/2 完成。
  • 长度前缀层(嵌在 DATA 帧里)由 gRPC 协议本身定义。HTTP/2 帧的 payload 里,gRPC 又套了一层 Length-Prefixed-Message:1 字节压缩标志 + 4 字节消息长度(大端序)+ 消息体。这层是 gRPC 定义的,对应 minirpc 的 Length 字段那层语义。
  • 第 3、4 层(内容)交给 protobuf。DATA 帧里那个消息体,才是 protobuf 序列化的产物。protobuf 只负责把结构化对象编码成字节,对切帧、多路复用、流式都不知情。

一个 HTTP/2 DATA 帧的字节构成如下图:外层是 HTTP/2 帧头,内层嵌套 gRPC 的 Length-Prefixed-Message,最里层才是 protobuf 序列化的消息体。

图 1:一个 HTTP/2 DATA 帧的字节构成
HTTP/2 frame header(9 B,由 HTTP/2 规范定义)Layer 1 + 2
lengthtype = DATAflagsstream id
gRPC Length-Prefixed-Message(DATA payload)
compressed-flag (1 B)message-length (4 B, BE)
protobuf-serialized message bytesLayer 3 + 4 内容
protobuf 只负责把结构化对象编码成字节;切帧与多路复用与它无关。
stream id = 多路复用钥匙 gRPC 自己的长度前缀

对照 minirpc 的单层帧(Length | StreamID | Type | Payload,JSON 负载)就能看清分工:gRPC 把”帧化”拆给 HTTP/2 做了一部分(帧头 + stream-id 复用),自己只补了一层 Length-Prefixed-Message 来框住 protobuf 消息;minirpc 则把这几件事打在一条帧里,负载用 JSON。两者的第 4 层”内容”分别由 protobuf 和 JSON 承担,但切帧与多路复用的机制是一模一样的。

这也印证了四层模型第 4 层与框架无关的判断:换掉负载格式(JSON 换 protobuf),换掉帧的外壳(自定帧换 HTTP/2),RPC 依然成立,因为前两层”切帧 + 按 id 分发”的机制没变。所谓应用层协议,本质上都是在 TCP 这条字节流之上,叠加上”切帧、复用、内容约定”这三件事。

亲手跑一遍:minirpc

前面用 minirpc 的协议讲完了三个机制,若想把它们拼成一个可运行的整体跑一遍,这个项目本身就是一个直接可用的样本:零外部依赖,只用 Go 标准库,配三个 demo 把三个机制分别可视化。

  • Demo 1 unaryAdd(2,3)=5,展示一次远程函数调用从序列化、帧化到等待响应的完整链路。
  • Demo 2 multiplex:一条 TCP 上并发三个慢请求(耗时 1 秒、2 秒、0.5 秒),完成顺序由耗时决定而非发起顺序,即上文机制二的可视化。
  • Demo 3 stream:服务端持续推送进度,复现一次真实的 containerd 竞态排查(PR #13625)——多路复用下”RPC 响应流”与”进度流”异步到达,漏掉最后一条进度。demo 提供 buggy 与 fixed 双模式,两者唯一差异是是否等待进度消费 goroutine 读完后再取计数。

minirpc 的帧对应 containerd 内部 RPC 协议 ttrpc 的消息帧,Conn 按 StreamID 分发对应 ttrpc 的流多路复用。demo 里学到的机制与生产级 RPC 一致,生产版只是多了 protobuf、TLS、完善的错误处理等工程外壳。如果想在读完本文后把这些机制亲手跑一遍、改一改参数看时序变化,minirpc 是一个直接可用的样本。仓库地址:gzb1128/minirpc