如何在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()); }
关键说明
- 交互模式匹配:服务端用
send()发送消息,对应RSocket的Fire-and-Forget模式,所以需要实现RSocket中的fireAndForget方法。 - 路由判断:通过
payload.metadata?.getRoutingMetadata()获取服务端指定的路由,判断是否为chat.receive。 - 连接保持:不要调用
rSocket.close(),否则连接断开后无法接收服务端后续推送的消息。 - 消息处理:在
fireAndForget方法内可以解析消息内容,执行更新UI、存储消息等业务逻辑。
内容的提问来源于stack exchange,提问作者Too
相关产品推荐
相关产品推荐

