本节摘要:Reactive Streams 之后,R2DBC、RSocket、MQTT 5 等协议把「非阻塞 + 背压」推到数据库与消息层。选型时要问:端到端是否仍是无界 push?
JPA/Hibernate 的 Session 与懒加载本质是阻塞线程。R2DBC(Reactive Relational Database Connectivity)定义与 JDBC 平行的 API,但 Connection.createStatement() 返回 Publisher<Result>,行以 onNext 到达,适合 WebFlux 全链路非阻塞。
| 维度 | JDBC | R2DBC |
|---|---|---|
| 线程模型 | 一连接阻塞一等 | 事件驱动 |
| 结果集 | 一次加载或分页 | 流式 Flux<Row> |
| ORM | Hibernate 成熟 | 以 SQL 为主,无 JPA 级 ORM |
| 背压 | 无(整结果集) | request 驱动逐行拉取 |
SOURCE 指出:PostgreSQL R2DBC 驱动较稳,MySQL 仍 Beta,Oracle 近乎空白——迁移前须做驱动与 SQL 方言评估。
ConnectionFactory factory = ConnectionFactories.get("r2dbc:postgresql://localhost/db"); Flux.from(factory.create()) .flatMapMany(conn -> conn.createStatement("SELECT id, name FROM users").execute()) .flatMap(result -> result.map((row, meta) -> row.get("name", String.class)));
R2DBC 与 JDBC 的差异不只是「API 长什么样」,而是执行模型的根本不同。JDBC 的 SELECT 执行后,驱动把结果行一次性(或按 fetch size)同步读进内存,连接在读取期间被独占——大量慢查询 = 大量挂起的连接。R2DBC 的 execute() 返回 Publisher,行以 onNext 流式到达,连接在等待期间可以被复用给其他请求,且下游可以 request(n) 按需拉取。同一个 PostgreSQL,线程占用从「查询期间独占」变成「查询期间空闲」,这就是非阻塞数据库驱动的意义。
RSocket 在 TCP/WebSocket 上提供四种交互模型,其中 REQUEST_STREAM 原生带 Reactive Streams 语义,适合微服务间长连接流式 RPC。
RSocket 的四种交互模型对应四种通信模式:
| 模型 | 请求/响应 | 语义 | 场景 |
|---|---|---|---|
| REQUEST_RESPONSE | 1:1 | 单请求单响应 | 普通 RPC |
| REQUEST_STREAM | 1:N | 单请求流响应 | 日志流、实时通知 |
| FIRE_AND_FORGET | 1:0 | 单向发射 | 事件上报、埋点 |
| CHANNEL | 1:1 双向流 | 双向对流 | 交互式会话、流式 AI |
其中 REQUEST_STREAM 与 CHANNEL 都携带 Reactive Streams 的请求-响应语义:接收方可以按需请求帧,实现跨进程的背压传导。这是背压第一次离开 JVM 内部、成为网络协议的组成部分——也是 4.2 选型时关注 RSocket 的根本原因。
MQTT 5 增加流控与共享订阅,边缘 IoT 设备向云端上报时,Broker 可与消费者协商 in-flight 窗口,避免 SOURCE 所述「推模型压垮消费者」的老问题。
// RSocket 客户端发起流式请求(概念示意) RSocket socket = ...; socket.requestStream(DefaultPayload.create(route, data)) // REQUEST_STREAM .map(Payload::getDataUtf8) .take(100) // 限制消费量,触发背压 .subscribe(System.out::println);
MQTT 5 的流控与 Reactive Streams 的 request 是同构的:Broker 与消费者协商一个「最大在途消息数」(Receive Maximum),超过后 Broker 停止发送,等消费者 ack 后再续发。这避免了经典推模型的恶性循环——消费者处理不过来,Broker 仍全速推送,导致客户端内存溢出、重启、再溢出。任何长连接的推送系统,最终都会重新发明 request(n)。
block() 调 JDBC)会截断背压链request 原子性要求不一,跨语言桥接需显式测试// 反模式:响应式链中混入 block(),背压在此断裂 public Flux<User> getUsers() { List<User> all = jdbcTemplate.query("SELECT * FROM users", ROW_MAPPER); // ↑ 一次阻塞拉全部,request(n) 在此失效 return Flux.fromIterable(all); }
这段代码看起来「返回了 Flux」,但阻塞 JDBC 查询已经把所有数据一次性拉进内存——下游无论怎么 request(n),内存里都已经躺着全量结果。背压的价值前提是「数据可以按需产生」,一旦数据已被急切物化,背压就失去意义。 正确做法要么用 R2DBC 让查询本身流式,要么明确告诉团队「这一跳是阻塞的,用 boundedElastic 隔离」,而不是伪装成响应式。
| 协议/驱动 | 成熟度 | 采用建议 |
|---|---|---|
| PostgreSQL R2DBC | 较成熟 | 可试点 |
| MySQL R2DBC | Beta | 谨慎评估 |
| Oracle R2DBC | 近乎空白 | 暂缓 |
| RSocket | 生产可用 | 微服务流式 RPC |
| MQTT 5 | 成熟 | IoT 首选 |
面对 R2DBC、RSocket、MQTT 5 等衍生规范,选型时建议逐项核对:
| 检查项 | 说明 | 达标标准 |
|---|---|---|
| 协议是否流式 | 数据能否按需到达 | 支持 request/窗口/ack |
| 背压是否跨进程 | 消费方能否约束发送方 | 有 in-flight 窗口或请求帧 |
| 驱动成熟度 | 生产可用等级 | 官方文档 + 社区案例 |
| 错误语义 | 错误是帧还是中断 | 有 error frame / onError |
| 取消语义 | 能否中途终止 | 有 cancel/断连机制 |
// 检查链:WebFlux → R2DBC → PostgreSQL // 1. WebFlux 层:Flux 原生 RS,OK // 2. R2DBC 层:execute() 返回 Publisher,OK // 3. PostgreSQL 驱动:较成熟,OK // 4. 若换成 Oracle R2DBC → 暂无驱动 → 链条断裂,需回退 JDBC + boundedElastic
检查的价值在于提前发现断点。很多团队把应用层换成 WebFlux 后才发现 Oracle 没有 R2DBC 驱动,只能回头给所有 DAO 加 subscribeOn(boundedElastic) 隔离——这就是 6.4 反模式「伪异步」的变体。规范的生态扩展不是均匀的,选型必须从最薄弱的环节往回推。
request一句话总结:R2DBC、RSocket、MQTT 5 是把「背压契约」从进程内扩展为进程间、从内存扩展为磁盘与网络的三条主线——但每一跳都要单独验证。
第 4 章进入 Rx 家族与 WebFlux/Vert.x 等具体实现选型。