Skip to content

HTTP/2、RPC 与通信协议设计 ​

#网络 · #网络协议 · #HTTP2 · #gRPC · #RPC · #Protobuf · #序列化 · #Thrift

Web 服务从 HTTP/1.1 演进到 HTTP/2,再到 gRPC(基于 HTTP/2),通信协议从文本走向二进制。本文涵盖:HTTP/2 帧格式、gRPC 的 Protobuf 序列化、自定义二进制协议的打包/解包、以及 JSON/Protobuf/Thrift/MessagePack 等序列化方案对比。


1. HTTP/2 二进制帧深入 ​

1.1 帧格式 ​

帧结构(9 字节固定头部):
+-----------------------------------------------+
|           Length (24 bits)       |Type(8)|Flags|
+-----------------------------------------------+
|R|              Stream Identifier (31 bits)     |
+-----------------------------------------------+
|              Frame Payload (Length 字节)        |
+-----------------------------------------------+
字段位数说明
Length24负载长度,最大 2^24-1 = 16MB(默认 16KB)
Type8帧类型 (DATA=0x0, HEADERS=0x1, SETTINGS=0x4 等)
Flags8帧标志 (END_STREAM=0x1, END_HEADERS=0x4 等)
R1保留位
Stream ID31流标识(0 = 连接控制帧)

1.2 关键帧类型 ​

go
// HTTP/2 帧类型常量
const (
    FrameData         = 0x0 // 请求/响应体
    FrameHeaders      = 0x1 // 头部帧(HPACK 压缩)
    FramePriority     = 0x2 // 流优先级
    FrameRSTStream    = 0x3 // 终止流
    FrameSettings     = 0x4 // 连接参数协商
    FramePushPromise  = 0x5 // 服务器推送承诺
    FramePing         = 0x6 // 心跳/ RTT 测量
    FrameGoAway       = 0x7 // 优雅关闭连接
    FrameWindowUpdate = 0x8 // 流量控制窗口更新
    FrameContinuation = 0x9 // HEADERS/PUSH_PROMISE 续帧
)

1.3 DATA 帧与流结束 ​

客户端: HEADERS (Stream 1, END_HEADERS)
        DATA    (Stream 1, 前 8KB)
        DATA    (Stream 1, 最后 4KB, END_STREAM)

服务端: HEADERS (Stream 1, END_HEADERS, :status: 200)
        DATA    (Stream 1, 响应体, END_STREAM)

2. gRPC — HTTP/2 上的 RPC 框架 ​

2.1 gRPC 协议分层 ​

mermaid
flowchart TB
    subgraph gRPC["gRPC 协议栈"]
        Business["业务层: .proto 定义<br/>service Greeter {<br/>  rpc SayHello(Req) returns(Res);<br/>}"]
        Serialization["序列化: Protobuf 编解码"]
        Transport["传输: gRPC 帧 → HTTP/2 帧"]
        Network["网络: TCP + TLS"]
    end

    Business --> Serialization --> Transport --> Network

2.2 gRPC 帧格式(在 HTTP/2 DATA 帧中) ​

gRPC 消息帧 (Length-Prefixed Message):
+-----------------------------------------------+
|  Compressed (1 bit) |    Length (31 bits)     |
+-----------------------------------------------+
|              Protobuf 序列化数据                |
+-----------------------------------------------+

多个消息可以在同一个 HTTP/2 DATA 帧中连续排列:
[Len1][Msg1][Len2][Msg2][Len3][Msg3]...

2.3 Protobuf 序列化:怎么打桩、解包 ​

protobuf
// user.proto
syntax = "proto3";
package example;

message User {
  int64  id       = 1;  // field number = 1
  string name     = 2;  // field number = 2
  string email    = 3;
  int32  age      = 4;
}

Protobuf 的 Wire Format(二进制编码):

每个字段 = Tag + Value
Tag     = (field_number << 3) | wire_type

Wire Types:
  0 = Varint (int32, int64, uint32, bool, enum)
  1 = 64-bit (fixed64, double)
  2 = Length-delimited (string, bytes, embedded messages, packed arrays)
  5 = 32-bit (fixed32, float)

编码示例:User{id: 42, name: "Alice"}

序列化后的字节流:
08 2A                    ← field 1 (id), varint 42
12 05 41 6C 69 63 65     ← field 2 (name), length 5, "Alice"

解析过程 (解包):
08 = 00001 000 → field_number=1, wire_type=0 (varint) → 读取 varint → 42
12 = 00010 010 → field_number=2, wire_type=2 (length-delimited)
  → 读取 varint length → 5 → 读取 5 字节 → "Alice"
mermaid
flowchart LR
    subgraph Encode["编码 (打桩/Stub)"]
        OBJ["User{id:42, name:'Alice'}"] --> TAG["08 2A<br/>12 05 41 6C 69 63 65"]
    end

    subgraph Decode["解码 (解包/Unmarshal)"]
        BIN["08 2A 12 05 41 6C 69 63 65"] --> T1["field 1: varint → 42"]
        BIN --> T2["field 2: length-delimited → 'Alice'"]
        T1 --> RES["User{id:42, name:'Alice'}"]
        T2 --> RES
    end

    Encode -.->|"传输"| Decode

2.4 Varint 编码(Protobuf 核心算法) ​

go
// Varint: 变长整数编码,小数字省空间
// 每组 7 bit 数据 + 1 bit 标志(是否还有下一组)

// 编码 42 (0x2A):
//   42 < 128 → 1 字节: 00101010

// 编码 300 (0x12C):
//   300 = 10101100 00000010
//   → 低 7 位: 0101100, 高 7 位: 0000010
//   → 字节 1: 1 0101100 = 0xAC (MSB=1, 还有后续)
//   → 字节 2: 0 0000010 = 0x02 (MSB=0, 结束)
//   → varint: AC 02

func encodeVarint(v uint64) []byte {
    var buf []byte
    for v >= 0x80 {
        buf = append(buf, byte(v)|0x80)
        v >>= 7
    }
    buf = append(buf, byte(v))
    return buf
}

2.5 gRPC 四种通信模式 ​

模式客户端服务端场景
Unary一发一回普通 RPC
Server Streaming一发多回日志推送、大文件下载
Client Streaming多发一回上传大文件、批量写入
Bidirectional多发多回聊天、实时协作

3. RPC 框架设计要点 ​

3.1 RPC 的完整调用链 ​

mermaid
sequenceDiagram
    participant Client as 客户端
    participant Stub as 客户端 Stub
    participant Net as 网络
    participant Skel as 服务端 Skeleton
    participant Server as 服务实现

    Client->>Stub: SayHello("Alice")
    Stub->>Stub: 序列化为二进制
    Stub->>Net: 发送请求 (方法名 + 参数)
    Net->>Skel: 接收请求
    Skel->>Skel: 反序列化参数
    Skel->>Server: 调用 SayHello("Alice")
    Server->>Skel: 返回 "Hello Alice"
    Skel->>Net: 发送响应
    Net->>Stub: 接收响应
    Stub->>Client: 返回 "Hello Alice"

3.2 Stub(桩)的核心职责 ​

职责说明
序列化将参数编码为字节流(Protobuf/JSON/自定义)
协议封装添加方法名、请求 ID、超时、元数据
传输发送到网络,接收响应
反序列化将响应字节解码为返回类型
错误处理网络超时、服务不可用、业务错误码

3.3 RPC 协议设计 ​

通用 RPC 请求格式:
+---------+----------+--------+--------+------+----------+
| Magic   | HeaderLen| Header | BodyLen| Body | Checksum |
| 2B      | 2B       | N      | 4B     | M    | 4B (可选)|
+---------+----------+--------+--------+------+----------+

Header 中包含:
  - Request ID / Sequence Number (用于多路复用和匹配请求-响应)
  - Method Name / Service.Method
  - Timeout (毫秒)
  - Trace ID (链路追踪)
  - Flags (压缩/加密标志)

4. 通信协议对比 ​

4.1 序列化协议对比 ​

协议格式大小 (相对 JSON)速度Schema可读性适用场景
JSON文本1× (基准)慢❌✅Web API, 调试友好
Protobuf二进制0.2× ~ 0.5×快✅ (.proto)❌gRPC, 高性能微服务
Thrift二进制0.2× ~ 0.5×快✅ (.thrift)❌跨语言 RPC
MessagePack二进制0.5× ~ 0.8×中❌❌ (可转 JSON)Redis 替代序列化
Avro二进制0.2× ~ 0.5×快✅ (.avsc)❌Hadoop/Kafka
FlatBuffers二进制0.2× ~ 0.5×极快 (零拷贝)✅ (.fbs)❌游戏、移动端
Cap'n Proto二进制0.2× ~ 0.5×极快 (零拷贝)✅❌沙盒 IPC
BSON二进制0.8× ~ 1.2×中❌❌MongoDB
XML/SOAP文本1.5× ~ 3×慢✅ (XSD/WSDL)✅遗留企业系统

4.2 选型决策 ​

mermaid
flowchart TD
    Q0["选通信协议"] --> Q1{"需要浏览器支持?"}
    Q1 -->|"是"| Q2{"需要双向实时通信?"}
    Q2 -->|"是"| WS["WebSocket + JSON/Protobuf"]
    Q2 -->|"否"| Q3{"需要高性能?"}
    Q3 -->|"是"| GRPC_WEB["gRPC-Web"]
    Q3 -->|"否"| REST["REST + JSON"]

    Q1 -->|"否, 后端间通信"| Q4{"需要 IDL / Schema?"}
    Q4 -->|"是"| Q5{"零拷贝需求?"}
    Q5 -->|"是"| FB["FlatBuffers / Cap'n Proto"]
    Q5 -->|"否"| Q6{"语言生态?"}
    Q6 -->|"Go/多语言"| GRPC["gRPC + Protobuf"]
    Q6 -->|"Java/遗留"| THRIFT["Thrift"]
    Q4 -->|"否, 简单场景"| Q7{"格式要求?"}
    Q7 -->|"人类可读"| JSON["JSON (文本)"]
    Q7 -->|"紧凑高效"| MSGPACK["MessagePack (二进制)"]

    style GRPC fill:#4CAF50,color:#fff
    style REST fill:#2196F3,color:#fff
    style WS fill:#FF9800,color:#fff

4.3 自定义二进制协议(TLV 编码) ​

TLV = Type-Length-Value

+--------+--------+--------+
| Type   | Length | Value  |
| 2B     | 2B     | N      |
+--------+--------+--------+

Type 0x0001: 用户 ID   (uint64)
Type 0x0002: 用户名    (UTF-8 string)
Type 0x0003: 邮箱      (UTF-8 string)
Type 0x0004: 年龄      (uint8)

示例: {name: "Alice", age: 30}
  → 00 02  00 05  41 6C 69 63 65      (name)
  → 00 04  00 01  1E                    (age)

TLV 的优势:

  • 向后兼容:未知 Type 直接跳过
  • 扩展性好:新增字段不影响旧解析器
  • 变长支持:Length 指定长度

5. WebSocket ​

5.1 协议升级 ​

http
客户端请求 (HTTP Upgrade):
GET /chat HTTP/1.1
Host: server.example.com
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==
Sec-WebSocket-Version: 13

服务端响应:
HTTP/1.1 101 Switching Protocols
Upgrade: websocket
Connection: Upgrade
Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=

Key → Accept 的算法:

Accept = Base64(SHA-1(Key + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"))

这是为了证明服务器理解 WebSocket 协议,而非误打误撞。

5.2 WebSocket 帧格式 ​

 0                   1                   2                   3
 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1 2 3 4 5 6 7 8 9 0 1
+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
|F|R|R|R|  Opcode   |M|        Payload Length (7|7+16|7+64)     |
|I|S|S|S|   (4)     |A|                                         |
|N|V|V|V|           |S|                                         |
+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
| Extended Payload Length (0|16|64 bit, 可选)                    |
+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
| Masking-Key (0 or 4 bytes, 仅客户端→服务端消息)                |
+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
| Payload Data                                                  |
+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+-+
字段说明
FIN1 = 最后一帧
Opcode0x1=Text, 0x2=Binary, 0x8=Close, 0x9=Ping, 0xA=Pong
MASK客户端→服务端必须掩码(防缓存投毒攻击)
Payload Length7bit (≤125) → 7+16bit (126) → 7+64bit (127)
Masking-Key4 字节,对 Payload 做 XOR 运算

5.3 Ping/Pong 心跳检测 ​

Ping 帧 (Opcode=0x9): "are you alive?"
Pong 帧 (Opcode=0xA): "yes, I'm alive"

应用层定时发送 Ping → 如果超时未收到 Pong → 断开重连

6. gRPC 深入 ​

6.1 Protocol Buffers 编码 ​

Protobuf 使用**变长编码(Varint)和标签-值(Tag-Value)**格式:

消息 = 字段1(标签+值) + 字段2(标签+值) + ...

标签 = (field_number << 3) | wire_type
wire_type: 0=Varint, 1=64bit, 2=Length-delimited, 5=32bit
protobuf
message Person {
  string name = 1;   // tag = (1<<3)|2 = 10
  int32  age  = 2;   // tag = (2<<3)|0 = 16
}

编码后:0A 05 41 6C 69 63 65 10 1E

  • 0A = tag(10) → name 字段,长度 5,"Alice"
  • 10 = tag(16) → age 字段,Varint 30

6.2 Proto3 vs Proto2 ​

Proto2Proto3
默认值自定义零值(0/空串/false)
required有无(所有字段 optional)
未知字段丢弃保留(向前兼容)
JSON 映射无原生支持
推荐旧项目新项目 ✅

6.3 gRPC 四种通信模式 ​

protobuf
service ChatService {
  // 1. 一元 RPC(请求-响应)
  rpc SendMessage(Message) returns (Response);

  // 2. 服务端流式
  rpc Subscribe(Topic) returns (stream Message);

  // 3. 客户端流式
  rpc Upload(stream Chunk) returns (UploadStatus);

  // 4. 双向流式
  rpc Chat(stream Message) returns (stream Message);
}
go
// 双向流式示例
func (s *ChatServer) Chat(stream pb.ChatService_ChatServer) error {
    for {
        msg, err := stream.Recv()
        if err == io.EOF { return nil }
        if err != nil { return err }

        // 回复
        reply := &pb.Message{Text: "echo: " + msg.Text}
        if err := stream.Send(reply); err != nil {
            return err
        }
    }
}

6.4 gRPC 拦截器(Interceptor) ​

go
// 一元拦截器链
func UnaryServerInterceptor(ctx context.Context, req interface{},
    info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {

    // 前置逻辑:认证、日志、限流
    start := time.Now()
    log.Printf("call %s with %v", info.FullMethod, req)

    // 执行实际处理
    resp, err := handler(ctx, req)

    // 后置逻辑:耗时、错误记录
    log.Printf("call %s took %v err=%v", info.FullMethod, time.Since(start), err)
    return resp, err
}

// 注册拦截器
s := grpc.NewServer(
    grpc.UnaryInterceptor(UnaryServerInterceptor),
)

6.5 gRPC 负载均衡 ​

mermaid
flowchart LR
    subgraph Client["客户端"]
        R["Resolver<br/>(服务发现)"]
        LB["Load Balancer<br/>(Round Robin/Pick First)"]
        CONN["连接池"]
    end
    DNS["DNS/K8s Service/Etcd"] --> R
    R --> LB
    LB --> CONN
    CONN --> S1["Server 1"]
    CONN --> S2["Server 2"]

Go gRPC 支持:

  • pick_first:默认,选第一个可用
  • round_robin:轮询
  • xds:Envoy xDS 控制面(高级)
go
conn, err := grpc.Dial("dns:///my-service:50051",
    grpc.WithDefaultServiceConfig(`{"loadBalancingPolicy":"round_robin"}`),
    grpc.WithInsecure(),
)

6.6 gRPC 四种通信模式 — 时序图 ​

mermaid
sequenceDiagram
    participant C as Client
    participant S as Server

    rect rgb(240, 248, 255)
        Note over C,S: Unary RPC (一发一回)
        C->>S: Request
        S-->>C: Response
    end

    rect rgb(255, 248, 240)
        Note over C,S: Server Streaming (一发多回)
        C->>S: Request
        S-->>C: Response 1
        S-->>C: Response 2
        S-->>C: Response N (EOF)
    end

    rect rgb(240, 255, 240)
        Note over C,S: Client Streaming (多发一回)
        C->>S: Chunk 1
        C->>S: Chunk 2
        C->>S: Chunk N (EOF)
        S-->>C: Response
    end

    rect rgb(255, 240, 255)
        Note over C,S: Bidirectional Streaming (多发多回)
        C->>S: Msg 1
        S-->>C: Msg A
        C->>S: Msg 2
        S-->>C: Msg B
        C->>S: Msg 3 (EOF)
        S-->>C: Msg C
    end
模式客户端服务端HTTP/2 实现典型场景
Unary1次1次一个 Stream, HEADERS+单DATA普通 API 调用
Server Streaming1次N次单 Stream, 多个 DATA+END_STREAM日志推送、大文件下载、行情推送
Client StreamingN次1次单 Stream, 多个 DATA→HEADERS文件上传、批量写入、传感器数据上报
BidirectionalN次N次单 Stream, 双工交错聊天、实时协作、视频会议信令

6.7 gRPC vs REST — 完整对比 ​

维度gRPCREST (HTTP/1.1 + JSON)
协议HTTP/2HTTP/1.1
序列化Protobuf (二进制, 0.2-0.5× JSON)JSON (文本, 1×)
接口定义.proto 文件 → 类型安全OpenAPI/Swagger (非强制)
代码生成✅ 原生 protoc 多语言第三方工具
流式传输✅ 四种模式原生支持❌ 需 SSE / WebSocket
浏览器兼容❌ 需 grpc-web 代理✅ 原生支持
调试工具grpcurl/grpcui/BloomRPCcurl / Postman
连接复用✅ 多路复用 (天生)🟡 Keep-Alive (手动)
负载均衡需要 L7 LB (gRPC 连接不释放)L4/L7 均可
性能✅✅ 高 (二进制+多路复用)🟡 中 (文本解析)
适用场景微服务间通信Web API / 对外开放接口

一句话选型:内部微服务 → gRPC;对外开放 API、浏览器 → REST;需要实时双向通信 → gRPC Bidirectional 或 WebSocket。


7. 工程实践:gRPC 真正难的是流控、连接复用和超时传播 ​

7.0 一眼看懂:gRPC 慢通常慢在连接复用后的共享成本 ​

mermaid
flowchart LR
    A["长连接"] --> B["多 Stream 复用"]
    B --> C["流控 / 消息大小 / 丢包共享"]
    C --> D["deadline 传播"]
    D --> E["重试、排队、级联超时"]
现象优先看什么常见误判
deadline exceeded 多下游 RT、预算传播、重试只怪网络不稳
连接数少但 RT 高单连接负载、流控、丢包只看“连接少就是好”
Streaming 卡住send/recv 节奏、取消传播只怪 Protobuf 慢
负载不均长连接粘连、LB、连接刷新只怪服务发现

7.1 gRPC 快,不代表调用一定便宜 ​

gRPC 的优势来自:

  • Protobuf 编码更紧凑
  • HTTP/2 多路复用减少连接数量
  • 长连接减少握手成本
  • 代码生成降低协议使用错误

但它的工程复杂度也更高,因为它把很多问题集中到了长连接 + 多 Stream + 流控 + deadline 传播 这一层。

7.2 多路复用不是没有阻塞,而是阻塞形态变了 ​

HTTP/2/gRPC 解决了 HTTP/1.1 应用层队头阻塞,但并不代表没有等待:

  • TCP 丢包仍会让同一连接上的多个 Stream 一起受影响
  • 连接级窗口和流级窗口不足会限制发送速度
  • 某些大流量 Stream 会挤压同连接上的小请求
  • 单连接承载过多请求时,尾延迟会被放大

所以 gRPC 的“快”依赖前提是:连接复用合理、流量模型稳定、窗口和并发没有被极端请求打穿。

7.3 Streaming 的风险不只是代码复杂,而是生命周期变长 ​

流式 RPC 常见问题:

  • Stream 长时间不结束,资源长期占用
  • 读写双方节奏不一致,背压积累
  • 客户端取消不及时,服务端还在继续发送
  • 消息过大,导致单次发送和反序列化成本过高

这意味着 Streaming 不是简单的“多收几条消息”,而是把一次调用从短事务变成了长生命周期会话。

7.4 deadline 是预算,不是装饰字段 ​

在线上,一个 gRPC 调用如果没有合理 deadline,常见后果是:

  • 下游慢调用长期悬挂
  • 上游请求已经超时,但下游还在继续执行
  • 连接池、goroutine、线程池迟迟不释放
  • 重试叠加后形成放大风暴
mermaid
flowchart LR
    A["上游未设置或错误设置 deadline"] --> B["下游调用持续等待"]
    B --> C["资源长时间占用"]
    C --> D["连接池和 goroutine 堆积"]
    D --> E["超时、重试、雪崩放大"]

7.5 gRPC 的连接池问题,经常表现为“看起来连接很少,但延迟越来越高” ​

因为 gRPC 倾向长连接复用,所以线上常见现象是:

  • 连接数不多
  • 但单连接负载越来越重
  • 某些连接被热点实例、热点调用压满
  • P99/P999 抖动明显

因此 gRPC 不是“连接越少越好”,而是要看:

  • 每条连接上承载多少并发流
  • 服务发现变更后连接是否及时更新
  • 负载均衡策略是否把请求均匀打散
  • 是否存在单大流拖累小请求

7.6 拦截器是能力入口,也是最容易隐藏成本的地方 ​

拦截器里经常放:

  • 鉴权
  • 日志
  • 指标
  • trace
  • 限流
  • 重试

但如果这里做得过重,就会出现:

  • 每个请求都重复做昂贵逻辑
  • 日志和序列化开销变成公共成本
  • 重试和超时策略互相打架
  • 错误看起来都像“网络超时”,但根因藏在拦截器链里

7.7 一个典型故障链:grpc client -> grpc server -> mysql ​

mermaid
sequenceDiagram
    participant C as gRPC Client
    participant S as gRPC Server
    participant M as MySQL

    C->>S: Unary / Stream 请求
    S->>M: 查询数据库
    Note over M: 慢 SQL / 锁等待 / I/O 抖动
    M-->>S: 返回变慢
    Note over S: stream 占用更久 / handler 处理超时边缘
    S-->>C: deadline exceeded
    Note over C: 重试 / 上游级联超时

7.8 常见误判 ​

现象容易误判为实际可能是
deadline exceeded 很多网络不稳定下游慢、超时预算不合理、重试放大
连接数少但 RT 高连接管理做得很好单连接过载、流控或丢包放大
Streaming 慢Protobuf 编码差消息太大、背压、取消不及时
负载不均服务发现有问题长连接粘连、LB 策略或连接刷新不及时
CPU 不高但请求超时服务空闲等待下游、流控、I/O 或连接排队

7.9 排障顺序 ​

现象优先看什么
deadline exceededdeadline 预算、下游 RT、是否有重试
某实例特别慢连接分布、LB 策略、热点流量
Streaming 卡住send/recv 节奏、窗口、取消传播
RT 抖动丢包、单连接并发、消息大小、拦截器成本
调用链级联超时上下游 timeout 是否按预算递减

7.10 一个实战原则 ​

text
把 gRPC 看成"长连接上的多路复用调用系统",
而不是"更快的 HTTP JSON";
这样你更容易看懂延迟、流控和雪崩是怎么产生的。

8. gRPC 实战专题:Deadline、Cancellation 与 Stream 队头阻塞 ​

8.1 Deadline 传递与超时预算分配 ​

gRPC 的 deadline 会自动在调用链上传递,这是防止级联超时的核心机制:

mermaid
sequenceDiagram
    participant C as Client<br/>deadline=5s
    participant A as Service A<br/>剩余=4.2s
    participant B as Service B<br/>剩余=2.1s

    C->>A: gRPC call (deadline=5s)
    Note over A: 处理 800ms
    A->>B: gRPC call (deadline=4.2s)
    Note over B: 处理 2100ms
    B-->>A: response
    A-->>C: response

    Note over C,B: 如果 B 处理超 4.2s → DEADLINE_EXCEEDED<br/>A 收到后可以: 返回错误 OR 用降级数据

超时预算分配原则:

go
// ✅ 好的做法: 每层预留预算
func ServiceA(ctx context.Context, req *pb.Request) (*pb.Response, error) {
    // 获取当前剩余 deadline
    deadline, ok := ctx.Deadline()
    if !ok { deadline = time.Now().Add(5 * time.Second) }

    remaining := time.Until(deadline)

    // 预留 20% 给自己,80% 给下游
    downstreamCtx, cancel := context.WithTimeout(ctx, remaining*8/10)
    defer cancel()

    // 调用下游
    resp, err := client.Call(downstreamCtx, req)
    if err != nil {
        // 可以降级返回
        return fallbackResponse(), nil
    }
    return resp, err
}

// ❌ 坏的做法: 每层硬编码超时
ctx, _ := context.WithTimeout(ctx, 3*time.Second) // 不管上游 deadline

8.2 Cancellation 传播链 ​

text
Client cancel → gRPC RST_STREAM (CANCEL) → Server context.Done()
  → 如果 server 没有监听 ctx.Done() → goroutine 继续执行 → 泄漏!

正确写法:
  select {
  case <-ctx.Done():
      return nil, ctx.Err()  // 收到取消信号,立即停止
  case result := <-processAsync():
      return result, nil
  }

8.3 HTTP/2 Stream 队头阻塞 vs TCP 队头阻塞 ​

mermaid
flowchart TD
    subgraph "TCP 队头阻塞"
        T1["TCP 流<br/>[P1][P2][P3]...→"] --> T2["P2 丢包"]
        T2 --> T3["P3-P10 全部卡住<br/>必须等 P2 重传"]
    end
    subgraph "HTTP/2 Stream 队头阻塞"
        H1["Stream 1: GET /a<br/>Stream 2: GET /b<br/>Stream 3: GET /c"] --> H2{"Stream 2 丢包?"}
        H2 -->|"所有 Stream 共享 TCP 连接"| H3["Stream 1,3 也被阻塞!"]
    end
    subgraph "HTTP/3 解决方案"
        Q1["Stream 1: QUIC Stream"]
        Q2["Stream 2: QUIC Stream"]
        Q3["Stream 3: QUIC Stream"]
        Q2X["Stream 2 丢包?"]
        Q2X --- Q1N["Stream 1 不受影响 ✅"]
        Q2X --- Q2N["Stream 2 单独重传"]
        Q2X --- Q3N["Stream 3 不受影响 ✅"]
    end

8.4 gRPC Client 连接池与负载均衡 ​

go
// gRPC 默认: 长连接 + 不关闭
// 意味着连接数 = server 实例数

conn, _ := grpc.Dial("dns:///service:8080",
    grpc.WithDefaultServiceConfig(`{
        "loadBalancingPolicy": "round_robin"
    }`),
    grpc.WithKeepaliveParams(keepalive.ClientParameters{
        Time:    30 * time.Second, // 30s 发一次 keepalive
        Timeout: 10 * time.Second, // 10s 没 ack 就断开
    }),
)
问题表现修复
单点连接新 server 实例启动后流量不均衡启用 round_robin + DNS 刷新
死连接server 挂了但 client 连接还保留keepalive.Timeout + MaxConnectionAge
连接泄漏conn.Close() 被忘记确保 defer close

8.5 DEADLINE_EXCEEDED 排查 ​

bash
# 1. 看 error 统计
# grpc-status-code: [4] = DEADLINE_EXCEEDED

# 2. 用 grpc-trace 加详细日志
GRPC_GO_LOG_VERBOSITY_LEVEL=99 GRPC_GO_LOG_SEVERITY_LEVEL=info ./server

# 3. 检查 deadline 传递
# 服务端打印 ctx.Deadline() 确认是否正确传递
原因排查方向
上游 deadline 太紧缩短上游处理时间 或 延长 deadline
中间层没有传递 deadline检查 context.WithTimeout 是否正确使用
下游服务慢降级返回 或 异步处理
网络丢包导致重传抓包看重传统计


gRPC 流式通信实战 ​

Proto 定义 ​

protobuf
syntax = "proto3";
package stream;

service StreamService {
    // 服务端流: 客户端发一个请求, 服务端持续推送
    rpc ServerStream(Request) returns (stream Response);

    // 客户端流: 客户端持续发送, 服务端汇总后返回一个响应
    rpc ClientStream(stream Request) returns (Response);

    // 双向流: 双方独立收发
    rpc BidiStream(stream Request) returns (stream Response);
}

message Request  { string msg = 1; }
message Response { string msg = 1; int32 seq = 2; }

服务端 ​

go
func (s *streamServer) ServerStream(req *pb.Request, stream pb.StreamService_ServerStreamServer) error {
    for i := 0; i < 5; i++ {
        if err := stream.Send(&pb.Response{
            Msg: fmt.Sprintf("server msg %d for %s", i, req.Msg),
            Seq: int32(i),
        }); err != nil {
            return err
        }
    }
    return nil
}

func (s *streamServer) ClientStream(stream pb.StreamService_ClientStreamServer) error {
    var msgs []string
    for {
        req, err := stream.Recv()
        if err == io.EOF {
            return stream.SendAndClose(&pb.Response{
                Msg: fmt.Sprintf("received %d messages: %v", len(msgs), msgs),
            })
        }
        if err != nil {
            return err
        }
        msgs = append(msgs, req.Msg)
    }
}

func (s *streamServer) BidiStream(stream pb.StreamService_BidiStreamServer) error {
    for {
        req, err := stream.Recv()
        if err == io.EOF { return nil }
        if err != nil { return err }

        stream.Send(&pb.Response{
            Msg: fmt.Sprintf("echo: %s", req.Msg),
            Seq: int32(time.Now().Unix()),
        })
    }
}

客户端 ​

go
func callServerStream(client pb.StreamServiceClient) {
    stream, _ := client.ServerStream(ctx, &pb.Request{Msg: "hello"})
    for {
        resp, err := stream.Recv()
        if err == io.EOF { break }
        fmt.Printf("seq=%d msg=%s\n", resp.Seq, resp.Msg)
    }
}

func callBidiStream(client pb.StreamServiceClient) {
    stream, _ := client.BidiStream(ctx)

    // 收发独立 goroutine
    go func() {
        for i := 0; i < 3; i++ {
            stream.Send(&pb.Request{Msg: fmt.Sprintf("ping %d", i)})
        }
        stream.CloseSend()
    }()

    for {
        resp, err := stream.Recv()
        if err == io.EOF { break }
        fmt.Printf("recv: %s\n", resp.Msg)
    }
}
流模式客户端服务端典型场景
Unary1 请求1 响应普通 RPC
Server Stream1 请求N 响应订阅推送、大数据导出
Client StreamN 请求1 响应批量上传、日志采集
Bidi StreamN 请求M 响应聊天、实时协作、游戏同步

参考 ​

批注模式

💬 文章评论

暂无评论,来说点什么吧 👇

编程学习笔记