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

如何在Flutter中实现RSocket客户端的路由与消息接收?

问题

我这边有个Java WebFlux的RSocket服务端,会通过chat.receive路由主动推送消息。目前用Flutter的rsocket:^1.0.0库完成了基础连接和请求响应,但不知道怎么在Flutter里实现类似Java中@MessageMapping("chat.receive")的路由控制器,来接收服务端通过这个路由发送的消息。

相关代码

Java WebFlux服务端核心代码

public Mono<String> sendMessage(SendMessageRequest requestBody, RSocketRequester rSocketRequester) {
    Long userId = clientManager.getUserIdBySocket(rSocketRequester);

    Chat chat = Chat.builder()
            .chatroomId(requestBody.getChatroomId())
            .createdDate(LocalDateTime.now())
            .lastModifiedDate(LocalDateTime.now())
            .userId(userId)
            .build();

    return chatRepository.save(chat)
            .flatMap(entity ->  chatReadRepository.findByUserIdAndChatroomId(userId, entity.getChatroomId()))
            .flatMapMany(entity -> chatMemberRepository.findAllByChatroomIdAndUserIdNot(entity.getChatroomId(), userId))
            .map(ChatMember::getUserId)
            .map(clientManager::getSocketByUserId)
            .flatMap(socketOptional ->
                socketOptional.<org.reactivestreams.Publisher<String>>map(socketRequester -> socketRequester.route("chat.receive")
                    .data(requestBody)
                    .send()
                    .thenReturn("Success!"))
                    .orElseGet(() -> Mono.just("fail!"))
            )
            .collectList().thenReturn("Success!");
}

服务端推送消息核心片段

socketRequester -> socketRequester.route("chat.receive")
                    .data(requestBody)
                    .send()
                    .thenReturn("Success!"))
                    .orElseGet(() -> Mono.just("fail!"))

Flutter客户端现有代码

void main() async {
  String jwt = "Bearer eyJ0eXAiOiJKV1QiLCJhbGciOiJIUzI1NiJ9.eyJzdWIiOiIxIiwicm9sZXMiOlsiVVNFUiJdLCJpYXQiOjE2NTk1OTg1OTgsImV4cCI6MTY1OTYwODU5OH0.bJHn4IKm6DtnGAQAxyruRb-LJSgKt-72-L9g7JqBtHw";
  var rSocket = await RSocketConnector.create()
      .setupPayload(routeAndDataPayload("socket.acceptor", jwt))
      .connect('tcp://192.168.219.101:8081');

  var payload = await rSocket.requestResponse!(routeAndDataPayload("healthcheck", "data!!!"));

  print(payload.getDataUtf8());

  rSocket.close();

  runApp(const MyApp());
}

解决方法

在Flutter的rsocket库中,要接收服务端指定路由的推送消息,需要通过RSocketConnector的acceptor方法注册消息处理器,对应服务端使用的Fire-and-Forget交互模式(也就是send()调用)。修改后的代码如下:

void main() async {
  String jwt = "Bearer eyJ0eXAiOiJKV1QiLCJhbGciOiJIUzI1NiJ9.eyJzdWIiOiIxIiwicm9sZXMiOlsiVVNFUiJdLCJpYXQiOjE2NTk1OTg1OTgsImV4cCI6MTY1OTYwODU5OH0.bJHn4IKm6DtnGAQAxyruRb-LJSgKt-72-L9g7JqBtHw";
  var rSocket = await RSocketConnector.create()
      .setupPayload(routeAndDataPayload("socket.acceptor", jwt))
      // 注册消息处理器,处理服务端推送的消息
      .acceptor((setup, sendingSocket) async {
        return RSocket(
          // 处理Fire-and-Forget类型的消息,对应服务端的send()
          fireAndForget: (payload) {
            // 从元数据中获取路由信息
            String? route = payload.metadata?.getRoutingMetadata();
            if (route == "chat.receive") {
              // 解析消息内容
              String messageContent = payload.getDataUtf8();
              print("收到聊天消息: $messageContent");
              // 这里可以添加UI更新、消息存储等逻辑
            }
          },
          // 复用原有RSocket的其他方法实现
          requestResponse: sendingSocket.requestResponse,
          requestStream: sendingSocket.requestStream,
          requestChannel: sendingSocket.requestChannel,
          metadataPush: sendingSocket.metadataPush,
        );
      })
      .connect('tcp://192.168.219.101:8081');

  var payload = await rSocket.requestResponse!(routeAndDataPayload("healthcheck", "data!!!"));
  print(payload.getDataUtf8());

  // 注意:不要立即关闭连接,否则无法接收后续推送的消息
  // rSocket.close();

  runApp(const MyApp());
}

关键说明

  1. 交互模式匹配:服务端用send()发送消息,对应RSocket的Fire-and-Forget模式,所以需要实现RSocket中的fireAndForget方法。
  2. 路由判断:通过payload.metadata?.getRoutingMetadata()获取服务端指定的路由,判断是否为chat.receive。
  3. 连接保持:不要调用rSocket.close(),否则连接断开后无法接收服务端后续推送的消息。
  4. 消息处理:在fireAndForget方法内可以解析消息内容,执行更新UI、存储消息等业务逻辑。

内容的提问来源于stack exchange,提问作者Too

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 14:18:32