grpc-reactor 实战(一):为什么直接在 Reactor Netty 上实现 gRPC
最近我开始开发 grpc-reactor:一个直接构建在 Reactor Netty HTTP/2 之上的实验性 gRPC 实现。项目使用 Protobuf,但运行时不依赖 grpc-java transport、ClientCall、ServerCall 或 StreamObserver。
这不是为了证明 grpc-java 不好。grpc-java 已经成熟、稳定,并覆盖负载均衡、NameResolver、重试和丰富的可观测性。这个项目要回答的是另一个问题:如果应用的编程模型本来就是 Reactor,能否让 Mono 和 Flux 从生成 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 | generated Reactor API |
每个 RPC 对应一个 HTTP/2 stream。请求以 HEADERS frame 开始,至少包含:
1 | :method POST |
消息体使用 Length-Prefixed-Message 格式封装在 DATA frame 中——每条 Protobuf 消息前有五字节 envelope(1 字节压缩标志 + 4 字节大端长度)。一个 DATA frame 可能包含多条 gRPC 消息,一条大消息也可能跨越多个 DATA frame。
为什么 gRPC 必须依赖 HTTP trailers? 这是协议最反直觉的设计之一。HTTP status code 对 gRPC 来说几乎无用——协议要求 HTTP 层始终返回 200,真正的调用结果(grpc-status 和 grpc-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 | grpc-reactor-protocol → 与传输无关的协议原语 |
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 | service InteropTestService { |
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.cancel 或 Context cancellation 通过 StreamObserver callback 传播 |
两侧都停止,只产生一个 terminal signal,迟到的 DATA/trailers 被忽略 |
| Trailers | 成功 trailers 保存在 call result 中;非 OK trailers 转换为 GrpcException.trailers() |
通过 ClientCall.Listener#onClose 或 StatusRuntimeException#getTrailers() 暴露 |
grpc-status、grpc-message 和自定义 trailers 在双向调用中都能保留,包括 trailers-only 错误 |
| Buffer ownership | Netty ByteBuf 有引用计数,成功、失败和取消路径都必须释放;只有解码后的 Protobuf 值越过 API 边界 |
生成的 Protobuf message 对应用隐藏了传输 buffer | 被拒绝、分片、压缩和取消的 frame 最终都达到 refCnt() == 0 |
例如,Reactor 的双向流测试从零 demand 开始,逐批请求响应:
1 | StepVerifier.create(client.bidirectionalStreaming(BIDI, requests), 0) |
对应的 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 | // Reactor client:收到第一条响应后取消 HTTP/2 调用 |
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 进行许可。