WebFlux RSocket服务端推送时rsocket-js客户端Responder收不到消息如何解决
RSocket服务端主动推送及客户端Responder配置问题排查方案
问题1:rsocket-js的Responder能否响应服务端的请求?
可以。RSocket是对等双向通信协议,客户端配置的Responder就是专门用于处理服务端主动向客户端发起的各类请求的,支持fireAndForget、requestResponse、requestStream、requestChannel四种交互模型。
问题2:当前Responder无法收到消息的原因及修复方案
你当前的问题同时存在服务端实现错误和客户端Responder配置缺失两方面问题,按以下步骤修复即可:
第一步:修正Spring Boot服务端实现
你当前的服务端有两个核心错误:
- 注入的
RSocketRequester不是绑定到具体客户端连接的实例,无法直接用于向客户端发请求 - 没有实现客户端连接管理逻辑,无法获取在线客户端的请求发送实例
1.1 新增客户端连接管理器
@Controller public class RSocketConnectionManager { // 线程安全的列表存储所有在线客户端的RSocketRequester实例 private final List<RSocketRequester> onlineClients = Collections.synchronizedList(new ArrayList<>()); // 监听客户端连接事件,保存连接实例 @ConnectMapping public void onClientConnect(RSocketRequester requester) { // 监听连接关闭事件,移除失效实例 requester.rsocket().onClose() .doFinally(signal -> onlineClients.remove(requester)) .subscribe(); onlineClients.add(requester); } // 对外提供获取所有在线客户端的方法 public List<RSocketRequester> getOnlineClients() { return new ArrayList<>(onlineClients); } }
1.2 修正服务端主动推送逻辑
@Service @RequiredArgsConstructor public class RsocketService { private final RSocketConnectionManager connectionManager; public void serverToClientPush(String eventContent){ // 遍历所有在线客户端推送事件 connectionManager.getOnlineClients().forEach(requester -> { // 指定发往客户端的路由为 client.event.push,使用fireAndForget单向推送 requester.route("client.event.push") .data("事件内容:" + eventContent + " 生成时间:" + LocalDateTime.now()) .send() // send()对应fireAndForget交互模型,无需客户端返回 .subscribe(null, err -> System.err.println("推送失败:" + err.getMessage())); }); } }
1.3 (可选)补全客户端发起请求的处理逻辑
你当前客户端代码主动发起了到request.stream路由的requestStream请求,需要补全服务端Controller逻辑才能收到返回:
@MessageMapping("request.stream") public Flux<String> requestStream(@Payload String requestData) { // 示例:每2秒返回一条消息,共返回10条 return Flux.interval(Duration.ofSeconds(2)) .map(i -> "服务端响应:" + i + " 收到客户端数据:" + requestData) .take(10); }
第二步:修正客户端Responder配置
你需要自己实现Responder的对应方法,并且解析路由元数据匹配指定路由的请求:
2.1 自定义Responder实现
import { Flowable } from 'rsocket-flowable'; class CustomResponder { // 解析路由元数据的工具方法,匹配message/x.rsocket.routing.v0格式 #parseRoute(metadata) { if (!metadata) return null; const routeLength = metadata[0]; return metadata.slice(1, 1 + routeLength).toString('utf8'); } // 处理服务端发起的fireAndForget请求(单向推送) fireAndForget(payload) { const route = this.#parseRoute(payload.metadata); // 匹配指定路由的请求 if (route === 'client.event.push') { const eventData = payload.data; console.log('收到服务端主动推送事件:', eventData); // 这里可以写更新React状态、渲染UI的逻辑 } } // 如果需要处理服务端发起的requestStream请求,实现该方法即可 requestStream(payload, responder) { const route = this.#parseRoute(payload.metadata); if (route === 'your.client.route') { // 返回流数据给服务端 return Flowable.just('客户端返回1', '客户端返回2', '客户端返回3'); } return Flowable.error(new Error('路由不存在')); } }
2.2 替换客户端初始化的Responder
const client = new RSocketClient({ // 替换为你自定义的Responder实例 responder: new CustomResponder(), transport: new RSocketWebSocketClient( { url: 'ws://localhost:7000/rsocket', wsCreator: (url) => new WebSocket(url), debug: true, } ), setup: { dataMimeType: "text/plain", metadataMimeType: 'message/x.rsocket.routing.v0', keepAlive: 600000, lifetime: 180000, } });
补充说明
如果你的核心需求只是服务端有事件时主动推给客户端,更简单的实现方案是客户端主动发起requestStream请求订阅服务端的事件流,服务端将该订阅的FluxSink保存,有事件时直接发送即可,不需要实现客户端Responder,逻辑更简单也更符合常规推送场景的实践。
内容的提问来源于stack exchange,提问作者Vani Koratekere
相关产品推荐
相关产品推荐

