You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

WebFlux RSocket服务端推送时rsocket-js客户端Responder收不到消息如何解决

RSocket服务端主动推送及客户端Responder配置问题排查方案

问题1:rsocket-js的Responder能否响应服务端的请求?

可以。RSocket是对等双向通信协议,客户端配置的Responder就是专门用于处理服务端主动向客户端发起的各类请求的,支持fireAndForget、requestResponse、requestStream、requestChannel四种交互模型。

问题2:当前Responder无法收到消息的原因及修复方案

你当前的问题同时存在服务端实现错误和客户端Responder配置缺失两方面问题,按以下步骤修复即可:

第一步:修正Spring Boot服务端实现

你当前的服务端有两个核心错误:

  1. 注入的RSocketRequester不是绑定到具体客户端连接的实例,无法直接用于向客户端发请求
  2. 没有实现客户端连接管理逻辑,无法获取在线客户端的请求发送实例
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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.07 06:09:02