RSocket 通信
一、RSocket 概述
1.1 什么是 RSocket
RSocket 是一种二进制、异步、多路复用的应用层网络协议,由 Netflix、Facebook 等公司联合设计,2017 年进入 Reactive Foundation 规范。RSocket 专为响应式系统间通信而生,原生支持背压(Backpressure)、双向通信(Bidirectional Streaming)和多路复用(Multiplexing)。
RSocket 的核心设计理念:
- 异步二进制协议:基于 Reactive Streams 语义,天然适配 Project Reactor / RxJava
- 单连接双向多路复用:一个 TCP/WebSocket 连接可承载多个独立数据流,避免连接风暴
- 连接语义与路由解耦:客户端与服务端通过 route 寻址,而非固定的 HTTP 端点
- 内置背压:Subscriber 可控制上游生产速率,防止下游过载
- 会话恢复(Resumption):断线重连后可恢复断点数据流,不丢消息
1.2 RSocket 与 HTTP 的对比
| 维度 | RSocket | HTTP/1.1 | HTTP/2 |
|---|---|---|---|
| 协议格式 | 二进制(元数据 + 数据分离) | 文本 | 二进制帧 |
| 传输层 | TCP / WebSocket | TCP | TCP |
| 交互模式 | 4 种(RR/FF/RS/CH) | 1 种(请求-响应) | 1 种 + Server Push |
| 多路复用 | 原生支持 | 不支持(需多连接) | 支持(流 ID 帧) |
| 背压 | 原生支持 | 不适用 | 不适用 |
| 双向流 | 原生支持 | 不支持 | 需 SSE/WebSocket |
| 会话恢复 | 支持 | 不支持 | 不支持 |
1.3 RSocket 协议帧结构
RSocket 协议帧包含三个部分:
+----------------+----------------+------------------+
| 帧头 (Frame Header) | 元数据 (Metadata) | 数据 (Data) |
+----------------+----------------+------------------+- 帧头:流 ID(Stream ID)、帧类型(Frame Type)、标志位(Flags)
- 元数据:路由信息、安全 Token、MIME 类型等,长度前缀编码
- 数据:业务载荷,长度前缀编码
元数据与数据分离是 RSocket 的重要特性:路由、安全等信息作为元数据传递,业务数据作为 Data 传递,两者可配置不同的序列化格式(如元数据用 JSON,数据用 Protobuf)。
二、RSocket 四种交互模式
RSocket 定义了四种标准交互模式,覆盖从简单的请求-响应到复杂的双向流通信场景。
2.1 Request-Response(请求-响应)
单对单模型:客户端发送一个请求,服务端返回一个响应。语义上等价于 HTTP 请求,是最常用的模式。
Client ─── Request ──→ Server
Client ←── Response ── Server// 服务端处理
@MessageMapping("stock.price.query")
public Mono<StockPrice> queryStockPrice(StockQueryRequest request) {
return stockService.getLatestPrice(request.getSymbol());
}
// 客户端调用
Mono<StockPrice> price = requester
.route("stock.price.query")
.data(new StockQueryRequest("AAPL"))
.retrieveMono(StockPrice.class);2.2 Fire-and-Forget(即发即忘)
客户端发送请求后不等待响应。适用于日志上报、事件通知、指标采集等不需要确认的场景,吞吐量最高。
Client ─── Request ──→ Server// 服务端处理
@MessageMapping("trade.log")
public Mono<Void> logTrade(TradeEvent event) {
return tradeLogger.log(event); // 返回 Mono<Void>
}
// 客户端调用(不等待响应)
Mono<Void> result = requester
.route("trade.log")
.data(new TradeEvent("BUY", "AAPL", 100))
.send();Fire-and-Forget 在 RSocket 协议层面仍保证消息到达(依赖 TCP 可靠性),但服务端不会发送响应帧,从而节省带宽和 CPU。
2.3 Request-Stream(请求-流)
客户端发送一个请求,服务端返回有序的多个元素的 Flux。适用于订阅实时数据流、分页查询、日志流等场景。
Client ─── Request ──→ Server
Client ←─── Stream ──── Server
Client ←─── Stream ──── Server (持续推送)// 服务端处理
@MessageMapping("stock.price.stream")
public Flux<StockTick> streamStockPrices(String symbol) {
return stockService.subscribePrice(symbol)
.map(price -> new StockTick(symbol, price, Instant.now()));
}
// 客户端订阅
Flux<StockTick> stream = requester
.route("stock.price.stream")
.data("AAPL")
.retrieveFlux(StockTick.class);
stream.subscribe(tick ->
System.out.printf("价格快照: %.2f @ %s%n",
tick.getPrice(), tick.getTimestamp()));背压在 Request-Stream 模式中天然生效:客户端可通过 request(n) 控制接收速率,服务端根据 Demand 信号调整生产速率。
2.4 Channel(双向通道)
双向流模型:客户端和服务端彼此独立地发送双向数据,形成全双工通信。适用于实时协同编辑、金融交易撮合、物联网双向指令等场景。
Client ───→ Server (独立发送)
Client ←─── Server (独立发送)
(双向并发)// 服务端——接收客户端的股票订阅列表,返回对应行情
@MessageMapping("channel.stocks")
public Flux<StockTick> stockChannel(Flux<StockSubscription> subscriptions) {
return subscriptions
.flatMap(sub -> stockService.subscribePrice(sub.getSymbol())
.map(price -> new StockTick(sub.getSymbol(), price, Instant.now())));
}
// 客户端——发送订阅列表,接收多只股票推送
Flux<StockSubscription> subscriptions = Flux.just(
new StockSubscription("AAPL"),
new StockSubscription("GOOGL"),
new StockSubscription("MSFT")
);
Flux<StockTick> ticks = requester
.route("channel.stocks")
.data(subscriptions)
.retrieveFlux(StockTick.class);2.5 四种模式对比总结
| 模式 | 请求数 | 响应数 | 语义 | 典型场景 |
|---|---|---|---|---|
| Request-Response | 1 | 1 | 同步调用 | RPC、查询 |
| Fire-and-Forget | 1 | 0 | 异步通知 | 日志、事件 |
| Request-Stream | 1 | N | 流式推送 | 订阅、日志流 |
| Channel | N | N | 双向流 | 交互、协同 |
三、Spring RSocket 支持
Spring Framework 从 5.2 开始引入 RSocket 支持,Spring Boot 2.2+ 提供自动配置。核心模块包括:
- spring-messaging:
@MessageMapping注解、RSocketRequesterAPI - spring-boot-starter-rsocket:Spring Boot 自动配置
- spring-security-rsocket:RSocket 安全支持
3.1 @MessageMapping 注解
服务端通过 @MessageMapping 注解处理方法,支持四种交互模式的自动路由:
@Controller
public class StockRSocketController {
// Request-Response
@MessageMapping("stock.query")
public Mono<StockQuote> queryStock(StockRequest req) {
return stockService.findQuote(req.getSymbol());
}
// Fire-and-Forget
@MessageMapping("stock.log")
public Mono<Void> logOperation(TradeLog log) {
return tradeLogService.save(log);
}
// Request-Stream
@MessageMapping("stock.subscribe")
public Flux<StockQuote> subscribeStock(String symbol) {
return stockService.streamQuotes(symbol);
}
// Channel(双向流)
@MessageMapping("stock.batch")
public Flux<StockQuote> batchSubscription(Flux<StockRequest> requests) {
return requests.flatMap(req -> stockService.streamQuotes(req.getSymbol()));
}
}方法的返回类型决定交互模式:
| 输入类型 | 返回类型 | RSocket 交互模式 |
|---|---|---|
| 普通类型 / Mono | Mono | Request-Response |
| 普通类型 / Mono | Mono<Void> | Fire-and-Forget |
| 普通类型 / Mono | Flux | Request-Stream |
| Flux | Flux | Channel |
3.2 RSocketRequester
RSocketRequester 是客户端的核心 API,用于向服务端发起 RSocket 请求。它支持链式编程风格定义路由、数据和预期响应类型。
创建 RSocketRequester:
@Configuration
public class RSocketClientConfig {
@Bean
RSocketRequester rSocketRequester(RSocketRequester.Builder builder) {
return builder
.dataMimeType(MediaType.APPLICATION_JSON)
.metadataMimeType(parseMetadata(WellKnownMimeType.MESSAGE_RSOCKET_ROUTING))
.connectTcp("localhost", 7000)
.block();
}
}四种交互操作:
@Component
public class StockClient {
private final RSocketRequester requester;
public StockClient(RSocketRequester requester) {
this.requester = requester;
}
// 1. Request-Response
public Mono<StockQuote> queryPrice(String symbol) {
return requester
.route("stock.query")
.data(new StockRequest(symbol))
.retrieveMono(StockQuote.class);
}
// 2. Fire-and-Forget
public Mono<Void> sendLog(TradeLog log) {
return requester
.route("stock.log")
.data(log)
.send();
}
// 3. Request-Stream
public Flux<StockQuote> subscribePrice(String symbol) {
return requester
.route("stock.subscribe")
.data(symbol)
.retrieveFlux(StockQuote.class);
}
// 4. Channel
public Flux<StockQuote> batchSubscribe(Flux<StockRequest> requests) {
return requester
.route("stock.batch")
.data(requests)
.retrieveFlux(StockQuote.class);
}
}3.3 Spring Boot 自动配置
在 application.yaml 中快速配置 RSocket 服务端:
spring:
rsocket:
server:
port: 7000
transport: tcp # 可选 tcp / websocket四、连接设置(TCP / WebSocket)
RSocket 支持两种传输层:TCP 和 WebSocket。TCP 提供最高性能,适合内网微服务间通信;WebSocket 可以穿透 HTTP 代理和防火墙,适合浏览器或公网场景。
4.1 TCP 传输
服务端配置(application.yaml):
spring:
rsocket:
server:
port: 7000
transport: tcp编程式:
@Configuration
public class RSocketServerConfig {
@Bean
public RSocketMessageHandler rSocketMessageHandler() {
RSocketMessageHandler handler = new RSocketMessageHandler();
handler.setRSocketStrategy(new RSocketStrategies());
return handler;
}
@Bean
public ServerRSocketConnector serverRSocketConnector(
RSocketMessageHandler handler) {
return ServerRSocketConnector.create(handler);
}
}客户端连接:
// TCP 直连
RSocketRequester requester = builder
.dataMimeType(MediaType.APPLICATION_JSON)
.metadataMimeType(parseMetadata(WellKnownMimeType.MESSAGE_RSOCKET_ROUTING))
.connectTcp("localhost", 7000)
.block();4.2 WebSocket 传输
服务端配置:
spring:
rsocket:
server:
port: 8080
transport: websocket
mapping-path: /rsocket # WebSocket 端点路径配置类方式:
@Configuration
public class WebSocketRSocketConfig {
@Bean
public RSocketMessageHandler rsocketMessageHandler() {
RSocketMessageHandler handler = new RSocketMessageHandler();
handler.setRSocketStrategy(RSocketStrategies.builder()
.metadataExtractorRegistry(registry -> {
registry.metadataToExtract(
MimeTypeUtils.parseMimeType("message/x.rsocket.routing.v0"),
String.class,
"route"
);
})
.build());
return handler;
}
@Bean
public HandlerMapping rsocketWebSocketMapping(
RSocketMessageHandler rsocketMessageHandler) {
return new SimpleUrlHandlerMapping() {{
setUrlMap(Map.of("/rsocket", rsocketMessageHandler));
setOrder(10);
}};
}
}客户端连接:
// WebSocket 连接
RSocketRequester requester = builder
.dataMimeType(MediaType.APPLICATION_JSON)
.metadataMimeType(parseMetadata(WellKnownMimeType.MESSAGE_RSOCKET_ROUTING))
.connectWebSocket(URI.create("ws://localhost:8080/rsocket"))
.block();4.3 TCP 与 WebSocket 选型建议
| 考量因素 | TCP | WebSocket |
|---|---|---|
| 吞吐量 | 更高(无 HTTP 头部开销) | 中等(HTTP 升级开销) |
| 防火墙穿透 | 困难(需开放端口) | 容易(基于 80/443) |
| 浏览器支持 | 不支持 | 支持 |
| SSL/TLS | RSocket TLS | WSS |
| 多路复用 | 原生 | 原生 |
| 适用场景 | 微服务内部调用 | 浏览器客户端、公网 |
五、路由配置
RSocket 的路由通过元数据传递,不在协议帧头部或 payload 中硬编码。Spring RSocket 使用 MESSAGE_RSOCKET_ROUTING(message/x.rsocket.routing.v0)作为元数据 MIME 类型。
5.1 路由元数据解析
服务端元数据提取器配置:
@Configuration
public class RSocketStrategyConfig {
@Bean
public RSocketStrategies rSocketStrategies() {
return RSocketStrategies.builder()
.metadataExtractorRegistry(registry -> {
registry.metadataToExtract(
MimeTypeUtils.parseMimeType(
"message/x.rsocket.routing.v0"),
String.class,
"route"
);
})
.build();
}
}5.2 路由通配符与变量
@MessageMapping 支持 Ant 风格通配符:
@Controller
public class GenericRSocketController {
// 精确匹配:stock.price.aapl
@MessageMapping("stock.price.{symbol}")
public Mono<StockPrice> getStockPrice(@DestinationVariable String symbol) {
return stockService.getPrice(symbol);
}
// 单层通配符:stock.*.latest
@MessageMapping("stock.*.latest")
public Flux<StockPrice> getLatestPrices() {
return stockService.getAllLatestPrices();
}
// 多层通配符:stock.**(匹配 stock/ 下的任意路由)
@MessageMapping("stock.**")
public Flux<StockPrice> allStockRoutes() {
return stockService.getAllPrices();
}
}5.3 客户端路由构建
// 方式一:字符串路由
requester.route("stock.price.AAPL").retrieveMono(StockPrice.class);
// 方式二:路由构建器(适合动态参数)
requester
.metadata(metadata -> metadata
.route("stock.price")
.metadata("AAPL", MimeTypeUtils.parseMimeType("application/x-java-object")))
.retrieveMono(StockPrice.class);
// 方式三:复合元数据(同时传递路由和安全 Token)
requester
.route("order.create")
.metadata(authToken, AUTH_MIME_TYPE)
.data(orderRequest)
.retrieveMono(OrderResponse.class);5.4 路由命名规范建议
领域.资源.操作 → stock.price.query
领域.资源.事件类型 → trade.execution.stream
版本.领域.资源 → v1.portfolio.snapshot
# 避免使用 HTTP 思维
# 错误:/api/v1/stock/price
# 正确:stock.price.v1RSocket 路由推荐使用点号分隔的扁平命名空间,而非 HTTP 的路径层级风格。路由名称是字符串匹配,不支持参数化路径(但支持 @DestinationVariable)。
六、安全(SSL/TLS、基本鉴权)
6.1 SSL/TLS 配置
服务端——配置 SSL 证书:
spring:
rsocket:
server:
port: 7001
transport: tcp
ssl:
enabled: true
key-store: classpath:server.p12
key-store-type: PKCS12
key-store-password: changeit
key-alias: rsocket-server编程式:
@Bean
public RSocketServer rSocketServer() {
SslContextBuilder sslContextBuilder = SslContextBuilder
.forServer(
new File("certs/server.crt"),
new File("certs/server.key"))
.protocols("TLSv1.3");
return RSocketServer.create(responderSocketAcceptor())
.transport(
TcpServerTransport.create("0.0.0.0", 7001)
.configure(config -> config.sslContext(sslContextBuilder)));
}客户端连接:
// 客户端——信任所有证书(仅开发环境)
SslContext sslContext = SslContextBuilder
.forClient()
.trustManager(InsecureTrustManagerFactory.INSTANCE)
.build();
RSocketRequester requester = RSocketRequester.builder()
.rsocketConnector(connector -> connector
.tcpClient(TcpClient.create()
.secure(ssl -> ssl.sslContext(sslContext))))
.connectTcp("localhost", 7001)
.block();6.2 基本鉴权(Bearer Token / Basic Auth)
服务端——JWT 鉴权拦截器:
@Component
public class AuthenticationRSocketInterceptor implements RSocketInterceptor {
private final JwtTokenValidator tokenValidator;
public AuthenticationRSocketInterceptor(JwtTokenValidator tokenValidator) {
this.tokenValidator = tokenValidator;
}
@Override
public Mono<Void> handle(ConnectionSetupPayload setupPayload,
RSocket rSocket) {
// 从连接建立元数据中提取 Token
String token = setupPayload.getMetadata()
.get("auth.token")
.map(buf -> {
byte[] bytes = new byte[buf.readableByteCount()];
buf.readBytes(bytes);
return new String(bytes, StandardCharsets.UTF_8);
})
.orElseThrow(() -> new IllegalArgumentException("缺失认证 Token"));
return tokenValidator.validate(token)
.flatMap(valid -> valid
? Mono.empty()
: Mono.error(new SecurityException("Token 验证失败")));
}
}服务端——集成 Spring Security:
spring:
rsocket:
server:
port: 7000
security:
rsocket:
authentication:
basic:
enabled: true@Configuration
@EnableRSocketSecurity
public class RSocketSecurityConfig {
@Bean
SecurityPolicyMarker securityPolicy(
RSocketSecurity rsocket) {
return rsocket
.authorizePayload(authorize ->
authorize
.route("stock.price.*").permitAll()
.route("order.**").authenticated()
.anyExchange().authenticated()
)
.simpleAuthentication(auth -> auth
.addUser(new UserDetails(
"user",
passwordEncoder().encode("password"),
List.of(new SimpleGrantedAuthority("ROLE_USER"))
))
)
.build();
}
@Bean
PasswordEncoder passwordEncoder() {
return new BCryptPasswordEncoder();
}
}客户端——携带 Token 连接:
@Bean
RSocketRequester authenticatedRequester(RSocketRequester.Builder builder) {
return builder
.setupPayload(payload -> {
// 通过 Setup 元数据传递 Token
payload.setMetadata(
"auth.token",
"eyJhbGciOiJIUzI1NiJ9...",
MimeTypeUtils.parseMimeType("message/x.rsocket.authentication.v0")
);
})
.connectTcp("localhost", 7000)
.block();
}七、与 WebSocket / gRPC 对比
7.1 综合对比表格
| 维度 | RSocket | WebSocket | gRPC |
|---|---|---|---|
| 协议 | 自定义二进制协议(RSocket 帧) | RFC 6455(WebSocket 帧) | HTTP/2(Protobuf) |
| 传输 | TCP / WebSocket | TCP(HTTP 升级) | HTTP/2 |
| 交互模式 | 4 种(RR/FF/RS/CH) | 1 种(双向消息) | 2 种(Unary / Streaming) |
| 数据序列化 | 应用层自行决定(JSON/Protobuf/等) | 应用层自行决定 | Protobuf(强制) |
| 元数据分离 | 原生支持(帧级元数据头部) | 不支持(需应用层编码) | 不支持(依赖 HTTP Headers) |
| 背压 | 原生支持(Reactive Streams 语义) | 不支持(需应用层实现) | 部分支持(HTTP/2 流控) |
| 多路复用 | 原生(单连接多流 ID) | 不支持(多路复用需多连接) | 原生(HTTP/2 流 ID) |
| 会话恢复 | 原生支持(Resumption) | 不支持 | 不支持 |
| 浏览器支持 | 需 JS 库(rsocket-js) | 原生 | 需 gRPC-Web + Envoy |
| 工具生态 | 中等 | 丰富 | 丰富 |
| 学习曲线 | 中等 | 低 | 中等 |
| 典型场景 | 微服务间异步通信、实时流、IoT | 网页实时聊天、游戏 | 内部 RPC、跨语言服务 |
7.2 选型决策树
需要双向流?
├── 是 → 需要浏览器支持?
│ ├── 是 → 需要元数据分离 / 会话恢复?
│ │ ├── 是 → RSocket over WebSocket
│ │ └── 否 → WebSocket
│ └── 否 → 需要背压 / 会话恢复?
│ ├── 是 → RSocket over TCP
│ └── 否 → gRPC(强类型 + 高性能)
└── 否 → 需要背压?
├── 是 → RSocket Request-Stream
└── 否 → 强类型约束 + 跨语言?
├── 是 → gRPC
└── 否 → RSocket / REST
追求最低延迟 + 最高吞吐?
├── RSocket over TCP(内网微服务)
└── gRPC over HTTP/2(跨语言、标准化)7.3 RSocket 的核心优势场景
- 微服务间异步流式通信:低延迟、内置背压、单连接多路复用,适合高吞吐事件驱动架构
- 实时数据推送:Request-Stream / Channel 模式天然支持,无需额外 SSE 或 WebSocket 封装
- 物联网(IoT):设备长时间连接、网络不稳定场景下,RSocket 会话恢复机制大幅降低重连开销
- 金融交易:双向流撮合、实时行情推送,对延迟和吞吐要求极高
八、实战:金融行情系统
本节构建一个完整的实时行情推送系统,涵盖 RSocket 的 Request-Stream 和 Channel 两种模式。
8.1 系统概览
┌─────────────┐ RSocket TCP ┌──────────────────┐
│ 行情客户端 │ ──────────────────────→ │ 行情服务端 │
│ StockClient │ Request-Stream │ StockServer │
│ (订阅实时报价) │ ←────────────────────── (推送实时报价) │
│ │ │ │
│ │ Channel(双向流) │ │
│ │ ←────────────────────── │ │
│ │ ──────────────────────→ │ │
└─────────────┘ └──────────────────┘
│
┌───┴───┐
│ 行情数据源 │
│(模拟器) │
└───────┘8.2 公共数据模型
// StockPrice.java — 股票报价
public class StockPrice {
private String symbol; // 股票代码
private BigDecimal bid; // 买入价
private BigDecimal ask; // 卖出价
private BigDecimal last; // 最新价
private long volume; // 成交量
private Instant timestamp; // 时间戳
// 无参构造(JSON 反序列化需要)
public StockPrice() {}
public StockPrice(String symbol, BigDecimal bid, BigDecimal ask,
BigDecimal last, long volume, Instant timestamp) {
this.symbol = symbol;
this.bid = bid;
this.ask = ask;
this.last = last;
this.volume = volume;
this.timestamp = timestamp;
}
// Getters / Setters(省略,实际需生成)
public String getSymbol() { return symbol; }
public void setSymbol(String symbol) { this.symbol = symbol; }
public BigDecimal getBid() { return bid; }
public void setBid(BigDecimal bid) { this.bid = bid; }
public BigDecimal getAsk() { return ask; }
public void setAsk(BigDecimal ask) { this.ask = ask; }
public BigDecimal getLast() { return last; }
public void setLast(BigDecimal last) { this.last = last; }
public long getVolume() { return volume; }
public void setVolume(long volume) { this.volume = volume; }
public Instant getTimestamp() { return timestamp; }
public void setTimestamp(Instant timestamp) { this.timestamp = timestamp; }
@Override
public String toString() {
return String.format("%s | 买: %.2f 卖: %.2f 最新: %.2f 量: %d @ %s",
symbol, bid, ask, last, volume, timestamp);
}
}// SubscriptionRequest.java — 订阅请求
public class SubscriptionRequest {
private String symbol; // 股票代码
private int intervalMs; // 推送间隔(毫秒)
public SubscriptionRequest() {}
public SubscriptionRequest(String symbol, int intervalMs) {
this.symbol = symbol;
this.intervalMs = intervalMs;
}
public String getSymbol() { return symbol; }
public void setSymbol(String symbol) { this.symbol = symbol; }
public int getIntervalMs() { return intervalMs; }
public void setIntervalMs(int intervalMs) { this.intervalMs = intervalMs; }
}8.3 服务端完整实现
build.gradle(依赖配置):
dependencies {
implementation 'org.springframework.boot:spring-boot-starter-rsocket'
implementation 'org.springframework.boot:spring-boot-starter-webflux'
implementation 'org.projectlombok:lombok'
annotationProcessor 'org.projectlombok:lombok'
testImplementation 'org.springframework.boot:spring-boot-starter-test'
testImplementation 'io.projectreactor:reactor-test'
}application.yaml:
spring:
application:
name: stock-quote-server
rsocket:
server:
port: 7000
transport: tcp
# 模拟行情配置
stock:
symbols: AAPL,GOOGL,MSFT,AMZN,TSLA
tick-interval-ms: 1000
price-range:
AAPL: 150-200
GOOGL: 2800-3000
MSFT: 300-400
AMZN: 3300-3600
TSLA: 700-900StockPriceSimulator.java —— 行情模拟器:
import reactor.core.publisher.Flux;
import reactor.core.publisher.Sinks;
import reactor.core.scheduler.Schedulers;
import jakarta.annotation.PostConstruct;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.time.Instant;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ThreadLocalRandom;
import java.util.concurrent.TimeUnit;
@Component
public class StockPriceSimulator {
private final Map<String, PriceRange> stockConfig = new HashMap<>();
private final Map<String, Sinks.Many<StockPrice>> sinks = new ConcurrentHashMap<>();
@Value("${stock.tick-interval-ms:1000}")
private int tickIntervalMs;
@PostConstruct
public void init() {
// 初始化各股票行情源,每秒推送一次价格变动
stockConfig.put("AAPL", new PriceRange(BigDecimal.valueOf(150), BigDecimal.valueOf(200)));
stockConfig.put("GOOGL", new PriceRange(BigDecimal.valueOf(2800), BigDecimal.valueOf(3000)));
stockConfig.put("MSFT", new PriceRange(BigDecimal.valueOf(300), BigDecimal.valueOf(400)));
stockConfig.put("AMZN", new PriceRange(BigDecimal.valueOf(3300), BigDecimal.valueOf(3600)));
stockConfig.put("TSLA", new PriceRange(BigDecimal.valueOf(700), BigDecimal.valueOf(900)));
for (String symbol : stockConfig.keySet()) {
Sinks.Many<StockPrice> sink = Sinks.many().multicast().onBackpressureBuffer();
sinks.put(symbol, sink);
PriceRange range = stockConfig.get(symbol);
Flux.interval(java.time.Duration.ofMillis(tickIntervalMs))
.map(tick -> generateTick(symbol, range))
.doOnNext(price -> {
Sinks.EmitResult result = sink.tryEmitNext(price);
if (result != Sinks.EmitResult.OK) {
System.err.printf("丢行情: %s %s%n", symbol, result);
}
})
.subscribe();
}
}
private StockPrice generateTick(String symbol, PriceRange range) {
ThreadLocalRandom rnd = ThreadLocalRandom.current();
double last = range.min.doubleValue()
+ rnd.nextDouble() * range.max.subtract(range.min).doubleValue();
double spread = last * 0.001; // 0.1% 价差
return new StockPrice(
symbol,
BigDecimal.valueOf(last - spread).setScale(2, RoundingMode.HALF_UP),
BigDecimal.valueOf(last + spread).setScale(2, RoundingMode.HALF_UP),
BigDecimal.valueOf(last).setScale(2, RoundingMode.HALF_UP),
rnd.nextLong(100, 10000),
Instant.now()
);
}
public Flux<StockPrice> streamPrice(String symbol) {
Sinks.Many<StockPrice> sink = sinks.get(symbol.toUpperCase());
if (sink == null) {
return Flux.error(new IllegalArgumentException("未知股票代码: " + symbol));
}
return sink.asFlux();
}
public Flux<StockPrice> streamAll() {
return Flux.fromIterable(stockConfig.keySet())
.flatMap(this::streamPrice);
}
private record PriceRange(BigDecimal min, BigDecimal max) {}
}StockQuoteService.java —— 行情服务:
import reactor.core.publisher.Flux;
import org.springframework.stereotype.Service;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@Service
public class StockQuoteService {
private final StockPriceSimulator simulator;
// 维护每个客户端请求的背压上下文
private final Map<String, Long> requestCounts = new ConcurrentHashMap<>();
public StockQuoteService(StockPriceSimulator simulator) {
this.simulator = simulator;
}
/**
* Request-Stream: 客户端订阅单只股票,持续推送报价
*/
public Flux<StockPrice> subscribe(String symbol) {
return simulator.streamPrice(symbol)
.doOnSubscribe(sub -> System.out.printf("[订阅] 客户端订阅 %s%n", symbol))
.doOnCancel(() -> System.out.printf("[取消订阅] 客户端取消 %s%n", symbol));
}
/**
* Channel: 客户端发送可变订阅列表,服务端返回合并行情
* 支持动态增删订阅的股票
*/
public Flux<StockPrice> channelSubscribe(Flux<SubscriptionRequest> requestStream) {
return requestStream
.flatMap(req -> {
System.out.printf("[Channel] 订阅请求: %s 间隔=%dms%n",
req.getSymbol(), req.getIntervalMs());
return simulator.streamPrice(req.getSymbol())
.delayElements(java.time.Duration.ofMillis(req.getIntervalMs()));
})
.doOnCancel(() -> System.out.println("[Channel] 客户端断开双向流"));
}
}StockRSocketController.java —— RSocket 处理器:
import org.springframework.messaging.handler.annotation.DestinationVariable;
import org.springframework.messaging.handler.annotation.MessageMapping;
import org.springframework.stereotype.Controller;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
@Controller
public class StockRSocketController {
private final StockQuoteService quoteService;
public StockRSocketController(StockQuoteService quoteService) {
this.quoteService = quoteService;
}
/**
* Request-Response: 查询单只股票最新报价快照
* 路由: stock.quote.{symbol}
*/
@MessageMapping("stock.quote.{symbol}")
public Mono<StockPrice> getQuote(@DestinationVariable String symbol) {
return quoteService.subscribe(symbol)
.next() // 取第一条作为快照
.timeout(java.time.Duration.ofSeconds(5));
}
/**
* Fire-and-Forget: 客户端发送订阅操作日志,无需响应
* 路由: stock.subscribe.log
*/
@MessageMapping("stock.subscribe.log")
public Mono<Void> logSubscription(SubscriptionRequest req) {
System.out.printf("[日志] 订阅操作: %s%n", req.getSymbol());
return Mono.empty();
}
/**
* Request-Stream: 实时行情推送
* 客户端发送股票代码,服务端持续推送 StockPrice
* 路由: stock.price.stream
*/
@MessageMapping("stock.price.stream")
public Flux<StockPrice> streamPrice(String symbol) {
return quoteService.subscribe(symbol);
}
/**
* Channel: 双向流批量订阅
* 客户端可动态追加/减少订阅标的
* 路由: stock.price.channel
*/
@MessageMapping("stock.price.channel")
public Flux<StockPrice> priceChannel(Flux<SubscriptionRequest> requests) {
return quoteService.channelSubscribe(requests);
}
}StockQuoteServerApplication.java —— 启动类:
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@SpringBootApplication
public class StockQuoteServerApplication {
public static void main(String[] args) {
SpringApplication.run(StockQuoteServerApplication.class, args);
System.out.println("=== 行情服务已启动,RSocket TCP 端口: 7000 ===");
}
}8.4 客户端完整实现
StockQuoteClientApplication.java —— 客户端启动配置:
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.builder.SpringApplicationBuilder;
import org.springframework.context.annotation.Bean;
import org.springframework.messaging.rsocket.RSocketRequester;
import org.springframework.util.MimeTypeUtils;
@SpringBootApplication
public class StockQuoteClientApplication {
public static void main(String[] args) {
new SpringApplicationBuilder(StockQuoteClientApplication.class)
.run(args);
}
@Bean
RSocketRequester rSocketRequester(RSocketRequester.Builder builder) {
return builder
.dataMimeType(org.springframework.util.MimeTypeUtils.APPLICATION_JSON)
.metadataMimeType(MimeTypeUtils
.parseMimeType("message/x.rsocket.routing.v0"))
.connectTcp("localhost", 7000)
.block();
}
}StockClientRunner.java —— 客户端运行器:
import org.springframework.boot.CommandLineRunner;
import org.springframework.messaging.rsocket.RSocketRequester;
import org.springframework.stereotype.Component;
import reactor.core.Disposable;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.math.BigDecimal;
import java.time.Duration;
import java.time.Instant;
import java.util.List;
import java.util.Scanner;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.atomic.AtomicInteger;
@Component
public class StockClientRunner implements CommandLineRunner {
private final RSocketRequester requester;
public StockClientRunner(RSocketRequester requester) {
this.requester = requester;
}
@Override
public void run(String... args) throws Exception {
System.out.println("=== 金融行情 RSocket 客户端 ===");
System.out.println("1. Request-Response — 查询单只报价");
System.out.println("2. Fire-and-Forget — 发送订阅日志");
System.out.println("3. Request-Stream — 订阅实时行情");
System.out.println("4. Channel — 双向流批量订阅");
System.out.print("请选择模式 (1-4): ");
Scanner scanner = new Scanner(System.in);
String choice = scanner.nextLine();
switch (choice.trim()) {
case "1" -> requestResponse();
case "2" -> fireAndForget();
case "3" -> requestStream();
case "4" -> channel();
default -> System.out.println("无效选择");
}
// 等待异步完成
Thread.sleep(30_000);
System.exit(0);
}
// ========== Request-Response ==========
private void requestResponse() {
System.out.println("\n--- Request-Response: 查询 AAPL 最新报价 ---");
Mono<StockPrice> result = requester
.route("stock.quote.AAPL")
.retrieveMono(StockPrice.class);
result.subscribe(
price -> System.out.println("查询结果: " + price),
error -> System.err.println("查询失败: " + error.getMessage())
);
}
// ========== Fire-and-Forget ==========
private void fireAndForget() {
System.out.println("\n--- Fire-and-Forget: 发送订阅日志 ---");
Mono<Void> result = requester
.route("stock.subscribe.log")
.data(new SubscriptionRequest("AAPL", 1000))
.send();
result.doOnSuccess(v -> System.out.println("日志发送成功(无响应确认)"))
.subscribe();
}
// ========== Request-Stream ==========
private void requestStream() {
System.out.println("\n--- Request-Stream: 订阅 GOOGL 实时行情(10秒后自动取消)---");
AtomicInteger counter = new AtomicInteger(0);
Disposable disposable = requester
.route("stock.price.stream")
.data("GOOGL")
.retrieveFlux(StockPrice.class)
.doOnCancel(() -> System.out.println("\n=== 订阅已取消 ==="))
.subscribe(
price -> System.out.printf("[#%02d] %s%n",
counter.incrementAndGet(), price),
error -> System.err.println("流异常: " + error)
);
// 10 秒后自动取消
Mono.delay(Duration.ofSeconds(10))
.doOnNext(v -> {
System.out.println("\n取消订阅...");
disposable.dispose();
})
.subscribe();
}
// ========== Channel ==========
private void channel() {
System.out.println("\n--- Channel: 双向流分批订阅多只股票 ---");
List<SubscriptionRequest> subscriptions = List.of(
new SubscriptionRequest("AAPL", 500), // 500ms 推送一次
new SubscriptionRequest("MSFT", 800), // 800ms 推送一次
new SubscriptionRequest("AMZN", 1200) // 1200ms 推送一次
);
// 一次发送所有订阅请求,然后持续接收行情
requester
.route("stock.price.channel")
.data(Flux.fromIterable(subscriptions)) // 服务端收到 Flux 并展开
.retrieveFlux(StockPrice.class)
.take(Duration.ofSeconds(8)) // 8 秒后停止
.subscribe(
price -> System.out.printf("[Channel] %s%n", price),
error -> System.err.println("Channel 异常: " + error),
() -> System.out.println("\n=== Channel 完成 ===")
);
}
}8.5 运行与测试
启动服务端:
# 终端 1:启动行情服务
./gradlew bootRun --args='--spring.profiles.active=server'
# 或直接运行主类
java -jar stock-quote-server.jar启动客户端:
# 终端 2:启动客户端
./gradlew bootRun --args='--spring.profiles.active=client'
# 选择 3,输出类似:
=== 金融行情 RSocket 客户端 ===
1. Request-Response
2. Fire-and-Forget
3. Request-Stream
4. Channel
# Request-Stream 输出示例:
[#01] GOOGL | 买: 2912.35 卖: 2918.45 最新: 2915.40 量: 4521 @ 2025-06-15T10:30:01.123Z
[#02] GOOGL | 买: 2913.10 卖: 2919.20 最新: 2916.15 量: 3812 @ 2025-06-15T10:30:02.124Z
[#03] GOOGL | 买: 2911.80 卖: 2917.90 最新: 2914.85 量: 5234 @ 2025-06-15T10:30:03.125Z8.6 压测验证要点
# 性能测试脚本(使用 rsocket-cli 工具)
# 安装: npm install -g rsocket-cli
# 测试 Request-Stream
rsocket-cli --stream --data "AAPL" \
--route "stock.price.stream" \
--server tcp://localhost:7000
# 测试连接数
# 使用 wrk2 或 ghz 进行 HTTP 基准测试
# RSocket 单连接可承载 10 万+ 流,注意调整 Netty 参数Netty 调优建议:
# application.yaml — Netty 线程模型优化
spring:
rsocket:
server:
port: 7000
transport: tcp
netty:
leak-detection: advanced// 编程式自定义 Netty 参数
@Bean
public RSocketServer rSocketServerCustomized() {
return RSocketServer.create(socketAcceptor())
.transport(TcpServerTransport.create(7000)
.configure(server -> server
.runOn(EventLoopGroup.builder()
.numThreads(Runtime.getRuntime().availableProcessors() * 2)
.build())
.channelOption(ChannelOption.CONNECT_TIMEOUT_MILLIS, 5000)
.channelOption(ChannelOption.SO_BACKLOG, 128)
.childOption(ChannelOption.TCP_NODELAY, true)
.childOption(ChannelOption.SO_KEEPALIVE, true)
));
}
// 占位——避免编译错误,实际需导入对应类
private RSocketResponderSocketAcceptor socketAcceptor() {
return null;
}九、常见问题与最佳实践
9.1 常见问题(FAQ)
Q1:RSocket 能否与浏览器直接通信?
RSocket 本身是二进制协议,浏览器不能原生支持。有两种方案:
- RSocket over WebSocket:使用
rsocket-js库(npm 包rsocket-websocket-client) - 通过 WebFlux 桥接:服务端提供 REST/SSE 端点,内部转发 RSocket 请求
Q2:RSocket 如何处理服务发现?
RSocket 不内置服务发现。推荐集成方式:
- Spring Cloud + Eureka/Nacos:客户端先通过服务发现获取目标地址,再建立 RSocket 连接
- RSocket Broker:部署 RSocket Broker(如 RSocket-Java Broker),客户端连接 Broker 进行路由转发
Q3:如何确保消息可靠性?
- Fire-and-Forget:依赖 TCP 可靠性,服务端宕机时消息可能丢失。需要可靠性时使用 Request-Response
- Request-Stream:内置背压 + 请求 n(Demand),消费者崩溃时生产者停止生产
- 会话恢复:启用 RSocket Resumption,断线重连后从断点继续
Q4:RSocket 连接池如何管理?
RSocket 单连接即可多路复用,通常不需要连接池。特殊场景可配置:
// 连接池——为不同目标主机建立独立连接
Map<String, RSocketRequester> pool = new ConcurrentHashMap<>();
public RSocketRequester getConnection(String host, int port) {
return pool.computeIfAbsent(host + ":" + port, k ->
RSocketRequester.builder()
.connectTcp(host, port)
.block()
);
}9.2 最佳实践
- 路由命名:使用点号分层命名,避免 HTTP 路径风格(
user.profile.get而非/api/user/profile) - MIME 类型配置:元数据使用
message/x.rsocket.routing.v0,数据推荐 Protobuf(压缩率高)或 JSON(开发调试友好) - 背压管理:在 Request-Stream 模式中合理使用
limitRate()/take()控制消费速率 - 错误处理:在
@MessageMapping方法中统一捕获异常,返回Mono.error()而非抛出运行时异常 - 连接保活:配置心跳(KeepAlive)间隔,TCP 默认 20 秒,可调整:
spring:
rsocket:
server:
port: 7000
transport: tcp
keep-alive-interval: 10s
keep-alive-max-lifetime: 60s- 监控与指标:集成 Micrometer 采集 RSocket 连接数、帧速率、延迟分布等指标:
@Bean
public RSocketServerCustomizer metricsCustomizer(MeterRegistry registry) {
return server -> server.interceptors(registryInterceptor(registry));
}
private RSocketInterceptorRegistry registryInterceptor(MeterRegistry registry) {
// 实际集成需引入 micrometer-rsocket 适配器
System.out.println("Micrometer RSocket metrics 已注册");
return null;
}十、总结
RSocket 作为新一代响应式网络协议,为微服务间通信提供了 HTTP 和 gRPC 之外的第三种选择。其核心优势在于:
| 优势 | 说明 |
|---|---|
| 原生背压 | 基于 Reactive Streams 规范,上下游速率匹配无需应用层编码 |
| 四种交互模式 | 覆盖请求-响应、即发即忘、流式订阅、双向通道全场景 |
| 单连接多路复用 | 避免连接风暴,降低系统资源开销 |
| 元数据与数据分离 | 路由、安全、追踪信息与业务载荷解耦 |
| 会话恢复 | 适用于不稳定网络场景,减少重连开销 |
适用场景建议:
- 推选 RSocket:微服务间实时流通信、高吞吐事件驱动架构、IoT 设备双向通信
- 继续使用 gRPC:跨语言 RPC 强类型约束、已有 Protobuf 契约体系
- 继续使用 WebSocket:浏览器原生支持、简单双向消息场景
Spring RSocket 通过 @MessageMapping 注解和 RSocketRequester API,将 RSocket 的复杂性封装在框架内部,开发者可以像编写普通 Controller 一样构建响应式通信服务。对于追求低延迟、高吞吐的金融行情、物联网、实时协同等场景,RSocket 是一个值得投入的技术选项。