grpc-reactor 实战(一):为什么直接在 Reactor Netty 上实现 gRPC

grpc-reactor 实战(一):为什么直接在 Reactor Netty 上实现 gRPC

Elvis Lv2

最近我开始开发 grpc-reactor:一个直接构建在 Reactor Netty HTTP/2 之上的实验性 gRPC 实现。项目使用 Protobuf,但运行时不依赖 grpc-java transport、ClientCallServerCallStreamObserver

这不是为了证明 grpc-java 不好。grpc-java 已经成熟、稳定,并覆盖负载均衡、NameResolver、重试和丰富的可观测性。这个项目要回答的是另一个问题:如果应用的编程模型本来就是 Reactor,能否让 MonoFlux 从生成 API 一直贯穿到 HTTP/2 stream,而不是在两套异步抽象之间适配?

本文基于 JDK 25、Gradle 9.2.1、Reactor 3.8.6、Reactor Netty 1.3.6 和 Netty 4.2.15.Final。main 已完成 Stage 0 到 Stage 10:在协议、四种 RPC cardinality、生产传输、DNS、codegen、标准服务和 Stage 9 加固之外,Stage 10 又加入了可选的负载均衡扩展制品,支持 grpclb、RLS、负载报告和 ORCA。

API 目标不是包装 StreamObserver

Protobuf 方法有四种 cardinality,目标 API 直接映射为:

Protobuf 方法 Reactor 签名
unary Mono<Resp> method(Mono<Req>)
server streaming Flux<Resp> method(Mono<Req>)
client streaming Mono<Resp> method(Flux<Req>)
bidirectional streaming Flux<Resp> method(Flux<Req>)

为什么 gRPC 要定义四种而不是只做 unary?本质上这是请求和响应各自独立选择”单值还是流”的 2x2 组合:

单值响应 流式响应
单值请求 unary server streaming
流式请求 client streaming bidirectional

它们解决不同场景:unary 覆盖经典 request-response;server streaming 用于服务端推送(事件订阅、大数据分页拉取);client streaming 用于客户端批量上传(文件分片、批量写入);bidirectional streaming 用于实时双向通信(聊天、协同编辑)。这不是凭空发明的四种模式——HTTP/2 stream 本身就是全双工的,单个连接可以多路复用数百条并发 stream,请求和响应都可以独立发送多帧数据。gRPC 把这种传输能力提升为一等 API 语义:不是 chunked transfer encoding 的变体,而是让编译器在类型层面检查调用方和实现方是否匹配。

低层传输只保留一个统一模型:

1
Flux<Req> -> HTTP/2 stream -> Flux<Resp>

生成代码负责在边界应用 single() 等 cardinality 检查。这样不需要为四种 RPC 分别实现四套网络逻辑,也能明确拒绝 unary 空请求、重复请求或重复响应。

兼容的是 wire protocol

项目的兼容目标是公开的 gRPC over HTTP/2 协议规范,不是 grpc-java 内部 API:

1
2
3
4
5
6
7
8
9
generated Reactor API
|
call dispatcher / service registry
|
marshaller + message framer/deframer
|
metadata + status + deadline
|
Reactor Netty HTTP/2 stream

每个 RPC 对应一个 HTTP/2 stream。请求以 HEADERS frame 开始,至少包含:

1
2
3
4
5
:method      POST
:path /<package.Service>/<Method>
content-type application/grpc+proto
te trailers
grpc-timeout <相对超时,如 100m 表示 100ms>

消息体使用 Length-Prefixed-Message 格式封装在 DATA frame 中——每条 Protobuf 消息前有五字节 envelope(1 字节压缩标志 + 4 字节大端长度)。一个 DATA frame 可能包含多条 gRPC 消息,一条大消息也可能跨越多个 DATA frame。

为什么 gRPC 必须依赖 HTTP trailers? 这是协议最反直觉的设计之一。HTTP status code 对 gRPC 来说几乎无用——协议要求 HTTP 层始终返回 200,真正的调用结果(grpc-statusgrpc-message)放在 trailers 中。原因在于:streaming 场景下,server 可能已经发送了数千条消息,最终成功还是失败只有处理完最后一条数据后才能确定。HTTP headers 在 body 之前发送,无法承载这个后验结果。trailers 是 HTTP/2 中唯一可以在 body 之后附加 metadata 的机制。

协议还定义了 trailers-only 模式:当 server 在读取 body 之前就能判定失败(比如路径不存在、认证失败),grpc-status 直接放在 response headers 中返回,不发送 body,节省一次 round-trip。Client 必须同时检查 headers 和 trailers 两个位置才能正确提取最终状态。

按协议边界拆模块

项目目前有九个模块:

1
2
3
4
5
6
7
8
9
grpc-reactor-protocol      → 与传输无关的协议原语
grpc-reactor-transport → Reactor Netty HTTP/2 映射
grpc-reactor-codegen → Protoc 插件,生成 Reactor stub
grpc-reactor-gradle-plugin → Gradle 集成
grpc-reactor-maven-plugin → Maven generate-sources 集成
grpc-reactor-services → 可选 Health、Reflection 与 Channelz 标准服务
grpc-reactor-binlog → 可选、有界的 canonical binary logging
grpc-reactor-lb-extensions → 可选 grpclb、RLS、负载报告与 ORCA
grpc-reactor-interop-test → grpc-java 兼容性测试

protocol 只处理与传输无关的值和编解码:消息帧、metadata、status、timeout、compression 和 protobuf marshaller。它依赖 Reactor Core 与 Netty Buffer,但不依赖 Reactor Netty——这意味着协议层可以独立测试而不需要真实的 HTTP/2 连接。

transport 把协议映射到 Reactor Netty HTTP/2,负责 client、server、service registry 和 per-call context。

codegen 消费 Protobuf CodeGeneratorRequest,生成类型安全的 Reactor client 和 service binder。

services 使用同一套 codegen 生成 canonical gRPC service binding,并在 transport 的 immutable descriptor/diagnostics snapshot 之上实现 Health v1、Reflection v1 和 Channelz v1。该模块需要显式注册,不会因为加入依赖就自动暴露管理端点。

binlog 同样是显式启用的独立模块。它以 transport interceptor 捕获 canonical binary-log v1 event,通过 metadata/message 截断、敏感 key 脱敏和固定容量 sink 控制信息暴露与内存上限。

interop-test 在测试作用域引入 grpc-java。grpc-java 在这里是兼容性判定器,不进入项目运行时。

依赖选型与版本锁定

运行时依赖链刻意保持短小:

依赖 版本 用途
Reactor Core 3.8.6 Mono/Flux 编程模型
Reactor Netty 1.3.6 HTTP/2 client/server
Netty 4.2.15 ByteBuf、HTTP/2 codec
Protobuf-java 4.35.1 消息序列化

项目不引入 Spring、Micrometer 或任何 DI 框架。测试使用 JUnit 6.1.2 与 Reactor Test 的 StepVerifier。grpc-java 1.82.2 仅出现在 interop-test 的 test classpath,运行时零侵入。

版本通过 gradle/libs.versions.toml 集中管理,配合 Gradle dependency locking 生成锁文件,确保不同机器的构建完全可复现。

为什么协议层必须先完成

直接从”启动一个 HTTP/2 Server”开始很容易得到一个能 echo 的演示,却会把真正困难的问题推迟:

  • DATA 可能在五字节 gRPC header 中间分片;
  • 一个 DATA buffer 可能包含多条消息;
  • metadata 允许重复键,二进制值需要 Base64;
  • timeout wire value 最多八位并带单位;
  • 最终状态来自 trailers;
  • ByteBuf 在成功、失败和取消路径都必须释放;
  • Reactive Streams demand 以消息计数,HTTP/2 flow control 以字节计数。

因此项目按 Stage 推进:先固定构建和互操作 fixture,再完成协议层,之后才实现 unary transport、streaming、生产特性和 codegen。每个 Stage 都有可执行退出条件,不以”类已经创建”作为完成标准。

Stage 0:构建基线与互操作 fixture

在写任何协议代码之前,Stage 0 解决的是”如何证明代码正确”:

JDK 25 编译配置——项目使用 -Xlint:all -parameters -encoding UTF-8 严格编译。所有警告都是编译错误,不允许静默忽略。选择 JDK 25 是为了提前验证 Netty 和 Protobuf 在最新 JVM 上的兼容性。

Spotless 格式化——统一使用 Eclipse formatter 配置加 ktlint,spotlessCheck 是 CI 门禁的第一道关。这消除了所有关于格式的 code review 讨论。

interop.proto 测试 fixture——定义了包含全部四种 RPC cardinality 的测试 service:

1
2
3
4
5
6
service InteropTestService {
rpc Unary (TestRequest) returns (TestResponse);
rpc ServerStreaming (TestRequest) returns (stream TestResponse);
rpc ClientStreaming (stream TestRequest) returns (TestResponse);
rpc BidirectionalStreaming (stream TestRequest) returns (stream TestResponse);
}

GrpcJavaFixture——在 interop-test 模块中,一个测试工具类同时启动 grpc-java server 和 client,提供 start() / close() 生命周期。通过它可以双向验证:Reactor client 调用 grpc-java server,以及 grpc-java client 调用 Reactor server。两者使用相同的 .proto 生成代码,确保 wire 兼容。

辅助工具

  • FreePorts:为每个测试分配独立端口,避免并行测试冲突;
  • TlsTestCertificates:为后续 TLS 测试预生成自签名证书;
  • LeakDetection:集成 Netty 的 ResourceLeakDetector,确保 ByteBuf 泄漏在测试中立即暴露。

Stage 0 退出条件./gradlew clean test 从 fresh checkout 通过,CI 在 Linux + Java 25 绿灯,protoc 生成确定性可重现,grpc-java fixture 双向通信正常。

分阶段验证策略

项目按 12 个 Stage 递进,每个 Stage 有明确的目标、可执行退出条件和回归保障:

Stage 目标 关键交付
0 构建基线 CI、格式、interop fixture
1 协议基础 帧编解码、metadata、status、timeout、compression
2 Unary transport 端到端 h2c unary 调用
3 Server streaming 多消息响应流
4 全 cardinality client streaming + bidirectional
5 生产传输 TLS、gzip、deadline、GOAWAY、连接池、keepalive
6 Name resolution DNS、subchannel、pick_first、round_robin
7 Codegen 与构建集成 protoc 插件、描述符注册表、Gradle/Maven 插件
8 标准服务 Health v1、Reflection v1、Channelz v1
9 运行面与加固 interceptor、observer、binlog、canonical smoke、fuzz/churn
10 负载均衡扩展 grpclb、RLS、load reporter、ORCA
11 诊断(计划中) 在有界 diagnostics registry 之上增加 Channelz v2

每个 Stage 完成时必须满足:./gradlew clean spotlessCheck test --no-daemon 全量通过,且前序 Stage 的测试作为回归继续运行。这意味着 Stage 4 完成时,Stage 2 的 unary 测试仍然是绿的。

Reactor 与 grpc-java 的语义差异

两套实现共享 wire protocol,但应用层契约并不相同。互操作测试验证消息、状态和 metadata 能跨实现传输;运行时测试则分别验证各自的生命周期规则,不能把两者简单视为同一种背压或取消机制:

关注点 Reactor 语义 grpc-java 语义 测试需要证明
Demand Subscription.request(n) 以解码后的消息计数,传输层还要同时遵守 HTTP/2 的字节窗口 入站流通过 ClientCall.request(n) 和 readiness callback 控制 下游没有 demand 时不能交付 response DATA,合并在一个 frame 中的消息也不能绕过 demand
Cancellation 取消 Mono/Flux 会取消两侧 application publisher,并在 stream 仍打开时发送 reset ClientCall.cancelContext cancellation 通过 StreamObserver callback 传播 两侧都停止,只产生一个 terminal signal,迟到的 DATA/trailers 被忽略
Trailers 成功 trailers 保存在 call result 中;非 OK trailers 转换为 GrpcException.trailers() 通过 ClientCall.Listener#onCloseStatusRuntimeException#getTrailers() 暴露 grpc-statusgrpc-message 和自定义 trailers 在双向调用中都能保留,包括 trailers-only 错误
Buffer ownership Netty ByteBuf 有引用计数,成功、失败和取消路径都必须释放;只有解码后的 Protobuf 值越过 API 边界 生成的 Protobuf message 对应用隐藏了传输 buffer 被拒绝、分片、压缩和取消的 frame 最终都达到 refCnt() == 0

例如,Reactor 的双向流测试从零 demand 开始,逐批请求响应:

1
2
3
4
5
6
7
StepVerifier.create(client.bidirectionalStreaming(BIDI, requests), 0)
.thenRequest(1)
.expectNext(first)
.thenRequest(2)
.expectNext(second, third)
.expectComplete()
.verify();

对应的 grpc-java 测试使用 StreamObserver callback 和完成信号。消息与 trailers 必须 wire-compatible,但两套测试断言的背压机制并不相同;这也是项目同时保留跨实现互操作测试和运行时故障测试的原因。

仓库现在还在 ServerStreamingInteroperabilityTest 中加入了两个 server streaming 取消互操作案例:Reactor client 从 grpc-java server 收到一条响应后取消,以及 grpc-java client 在收到 Reactor server 的第一条响应后取消。两侧都在 paranoid leak detection 下等待对端的取消 callback;更底层的 GrpcFrameCodecTest 则直接断言取消和 partial frame 输入最终达到 refCnt() == 0

核心取消路径其实很短:

1
2
3
4
5
6
7
8
9
10
11
12
// Reactor client:收到第一条响应后取消 HTTP/2 调用
TestResponse first = reactorStub
.serverStreaming(Mono.just(request(5)))
.take(1)
.single()
.block(Duration.ofSeconds(5));

// grpc-java client:在响应 callback 中取消
@Override
public void onNext(TestResponse response) {
requestStream.cancel("cancel after first response", null);
}

grpc-java server 通过 setOnCancelHandler 观察取消,Reactor server 则在响应 Flux 上使用 doOnCancel。包含 latch、对端状态断言、资源关闭和 leak detection 的完整代码见 ServerStreamingInteroperabilityTest.java

当前验证结果

项目统一使用下面的完整门禁,同时检查格式、编译、协议测试、transport 测试和 grpc-java 互操作测试:

1
./gradlew clean spotlessCheck test --no-daemon

本文更新时该命令已经通过。当前门禁除了协议、transport、codegen、标准服务、binary-log 和 grpc-java 互操作测试外,也包含 Stage 10 负载均衡扩展测试。依赖基线是 Protobuf 4.35.1、grpc-java 1.82.2 和 JUnit 6.1.2。JDK 25 下仍会看到 Protobuf Unsafe、Netty/Gradle native access 与 Gradle deprecated feature 提示,它们不影响本次测试结果,但需要在后续 JDK 和 Gradle 升级时持续跟踪。

后续文章将继续讨论五字节 gRPC message envelope、ByteBuf 所有权、四种 RPC cardinality、生产传输、DNS、代码生成、标准管理服务,以及如何在不破坏 Reactive Streams 语义的前提下接入 interceptor、观测与 binary log。

  • 标题: grpc-reactor 实战(一):为什么直接在 Reactor Netty 上实现 gRPC
  • 作者: Elvis
  • 创建于 : 2026-06-01 15:30:00
  • 更新于 : 2026-08-07 10:00:00
  • 链接: https://qianwj.github.io/2026/06/01/grpc-reactor-series-1/
  • 版权声明: 本文章采用 CC BY-NC-SA 4.0 进行许可。