代码之家  ›  专栏  ›  技术社区  ›  Mritunjay

WebFlux WebSocketClient,如何在同一会话中发送多个请求[设计客户端库]

  •  8
  • Mritunjay  · 技术社区  · 7 年前

    TL;DR;

    我们正在尝试使用SpringWebFluxWebSocket实现来设计WebSocket服务器。服务器具有常见的HTTP服务器操作,例如 create/fetch/update/fetchall .使用websockets,我们试图公开一个端点,这样客户机就可以利用一个连接进行所有类型的操作,因为websockets就是为此目的而设计的。WebFlux和WebSockets的设计是否正确?

    长版

    我们正在启动一个项目,该项目将使用来自 spring-webflux .我们需要构建一个反应式客户端库,消费者可以使用它连接到服务器。

    在服务器上, 我们得到一个请求,读取一条消息,保存它并返回一个静态响应:

    public Mono<Void> handle(WebSocketSession webSocketSession) {
        Flux<WebSocketMessage> response = webSocketSession.receive()
                .map(WebSocketMessage::retain)
                .concatMap(webSocketMessage -> Mono.just(webSocketMessage)
                        .map(parseBinaryToEvent) //logic to get domain object
                        .flatMap(e -> service.save(e))
                        .thenReturn(webSocketSession.textMessage(SAVE_SUCCESSFUL))
                );
    
        return webSocketSession.send(response);
    }
    

    在客户身上 ,我们想在有人打电话时打个电话 save 方法并从返回响应 server .

    public Mono<String> save(Event message) {
        new ReactorNettyWebSocketClient().execute(uri, session -> {
          session
                  .send(Mono.just(session.binaryMessage(formatEventToMessage)))
                  .then(session.receive()
                          .map(WebSocketMessage::getPayloadAsText)
                          .doOnNext(System.out::println).then()); //how to return this to client
        });
        return null;
    }
    

    我们不确定如何着手设计这个。理想情况下,我们认为

    1) client.execute 只应调用一次,并以某种方式保持 session . 同一会话应用于在后续调用中发送数据。

    2)如何返回我们进入的服务器的响应 session.receive ?

    3)万一 fetch 当响应很大(不仅仅是静态字符串,而是事件列表)时, 接待处 ?

    我们正在做一些研究,但是我们无法在网上找到WebFlux WebSocket客户端文档/实现的适当资源。关于如何前进的任何提示。

    2 回复  |  直到 7 年前
        1
  •  7
  •   Oleh Dokuka    7 年前

    拜托!使用 RSocket !

    它是绝对正确的设计,它值得为所有可能的操作节省资源,并且只为每个客户机使用一个连接。

    但是,不要实现一个轮子并使用提供所有这些通信的协议。

    • RSOCK有一个 请求-响应 这个模型允许您今天进行最常见的客户机-服务器交互。
    • RSOCK有一个 请求流 通信模型,这样您就可以满足所有的需要,并异步地返回一个事件流,重用同一个连接。rsocket将所有逻辑流映射到网络连接和网络连接,因此您自己不会感到这样做的痛苦。
    • rsocket有更多的交互模型,比如 火忘 溪流 如果是这样的话 以两种方式发送数据流。

    如何在弹簧中使用RSocket

    这样做的一个选择是使用RSotoJava实现RSoT协议。RSOCK Java是建立在项目反应堆之上的,因此它自然适合于春季WebF通量生态系统。

    不幸的是,没有与春季生态系统的特色集成。幸运的是,我花了几个小时来提供 RSocket Spring Boot Starter 它将SpringWebFlux与RSocket集成,并将WebSocket RSocket服务器与WebFlux HTTP服务器一起公开。

    为什么RSocket是更好的方法?

    基本上,rsocket隐藏了自己实现相同方法的复杂性。对于RSocket,我们不必关心交互模型定义作为自定义协议和在爪哇中的实现。rsocket为我们提供数据到特定的逻辑通道。它提供了一个将消息发送到同一个WS-Connection的内置客户端,因此我们不必为此发明自定义实现。

    让它变得更好 RSocket-RPC

    因为rsocket只是一个协议,它不提供任何消息格式,所以这个挑战是针对业务逻辑的。但是,有一个rsocket-rpc项目,它提供了一个协议缓冲区作为消息格式,并重用了与grpc相同的代码生成技术。因此,使用rsocket-rpc,我们可以很容易地为客户机和服务器构建一个API,并且完全不考虑传输和协议抽象。

    同一个RSocket弹簧引导集成提供了 example 以及rsocket-rpc的用法。

    好吧,它还没说服我,我还想有一个自定义的WebSocket服务器

    因此,为了达到这个目的,你必须自己去实现这个地狱。我以前已经做过一次,但是我不能指出那个项目,因为它是一个企业项目。 不过,我可以共享一些代码示例,这些示例可以帮助您构建适当的客户机和服务器。

    服务器端

    处理程序和打开逻辑订阅服务器映射

    必须考虑的第一点是,一个物理连接中的所有逻辑流都应存储在某个位置:

    class MyWebSocketRouter implements WebSocketHandler {
    
      final Map<String, EnumMap<ActionMessage.Type, ChannelHandler>> channelsMapping;
    
    
      @Override
      public Mono<Void> handle(WebSocketSession session) {
        final Map<String, Disposable> channelsIdsToDisposableMap = new HashMap<>();
        ...
      }
    }
    

    上面的示例中有两个地图。第一个是您的路由映射,它允许您根据传入的消息参数来标识路由。第二个是为请求流用例创建的(在我的例子中,它是活动订阅的映射),因此您可以发送一个消息框架来创建订阅,或者订阅特定的操作并保留该订阅,这样一旦执行了取消订阅操作,如果存在订阅,您将被取消订阅。

    使用处理器进行消息多路复用

    为了从所有逻辑流发回消息,必须将消息多路复用到一个流。例如,使用reactor,可以使用 UnicastProcessor :

    @Override
    public Mono<Void> handle(WebSocketSession session) {
      final UnicastProcessor<ResponseMessage<?>> funIn = UnicastProcessor.create(Queues.<ResponseMessage<?>>unboundedMultiproducer().get());
      ...
    
      return Mono
        .subscriberContext()
        .flatMap(context -> Flux.merge(
          session
            .receive()
            ...
            .cast(ActionMessage.class)
            .publishOn(Schedulers.parallel())
            .doOnNext(am -> {
              switch (am.type) {
                case CREATE:
                case UPDATE:
                case CANCEL: {
                  ...
                }
                case SUBSCRIBE: {
                  Flux<ResponseMessage<?>> flux = Flux
                    .from(
                      channelsMapping.get(am.getChannelId())
                                     .get(ActionMessage.Type.SUBSCRIBE)
                                     .handle(am) // returns Publisher<>
                    );
    
                  if (flux != null) {
                    channelsIdsToDisposableMap.compute(
                      am.getChannelId() + am.getSymbol(), // you can generate a uniq uuid on the client side if needed
                      (cid, disposable) -> {
                        ...
    
                        return flux
                          .subscriberContext(context)
                          .subscribe(
                            funIn::onNext, // send message to a Processor manually
                            e -> {
                              funIn.onNext(
                                new ResponseMessage<>( // send errors as a messages to Processor here
                                  0,
                                  e.getMessage(),
                                  ...
                                  ResponseMessage.Type.ERROR
                                )
                              );
                            }
                          );
                      }
                    );
                  }
    
                  return;
                }
                case UNSABSCRIBE: {
                  Disposable disposable = channelsIdsToDisposableMap.get(am.getChannelId() + am.getSymbol());
    
                  if (disposable != null) {
                    disposable.dispose();
                  }
                }
              }
            })
            .then(Mono.empty()),
    
            funIn
                ...
                .map(p -> new WebSocketMessage(WebSocketMessage.Type.TEXT, p))
                .as(session::send)
          ).then()
        );
    }
    

    正如我们从上面的示例中看到的,这里有很多东西:

    1. 信息应包括路线信息
    2. 消息应包含与之相关的唯一流ID。
    3. 用于消息复用的独立处理器,其中错误也应为消息
    4. 每个通道都应该存储在某个地方,在这种情况下,我们都有一个简单的用例,其中每个消息都可以提供 Flux 或者只是一个 Mono (在Mono的情况下,它可以在服务器端实现得更简单,因此您不必保留唯一的流ID)。
    5. 此示例不包括消息编码解码,因此此挑战留给您。

    客户端

    客户也不是那么简单:

    句柄会话

    为了处理连接,我们必须分配两个处理器,以便进一步使用它们来多路复用和解复用消息:

    UnicastProcessor<> outgoing = ...
    UnicastPorcessor<> incoming = ...
    (session) -> {
      return Flux.merge(
         session.receive()
                .subscribeWith(incoming)
                .then(Mono.empty()),
         session.send(outgoing)
      ).then();
    }
    

    将所有逻辑流保留在某个位置

    所有创建的流 单声道 通量 应该存储在某个地方,以便我们能够区分流消息与哪个相关:

    Map<String, MonoSink> monoSinksMap = ...;
    Map<String, FluxSink> fluxSinksMap = ...;
    

    由于Monosink,我们必须保留两个映射,而FluxSink没有相同的父接口。

    消息路由

    在上面的示例中,我们只是考虑了客户端的初始部分。现在我们必须构建一个消息路由机制:

    ...
    .subscribeWith(incoming)
    .doOnNext(message -> {
        if (monoSinkMap.containsKey(message.getStreamId())) {
            MonoSink sink = monoSinkMap.get(message.getStreamId());
            monoSinkMap.remove(message.getStreamId());
            if (message.getType() == SUCCESS) {
                sink.success(message.getData());
            }
            else {
                sink.error(message.getCause());
            }
        } else if (fluxSinkMap.containsKey(message.getStreamId())) {
            FluxSink sink = fluxSinkMap.get(message.getStreamId());
            if (message.getType() == NEXT) {
                sink.next(message.getData());
            }
            else if (message.getType() == COMPLETE) {
                fluxSinkMap.remove(message.getStreamId());
                sink.next(message.getData());
                sink.complete();
            }
            else {
                fluxSinkMap.remove(message.getStreamId());
                sink.error(message.getCause());
            }
        }
    })
    

    上面的代码示例显示了如何路由传入消息。

    多路传输请求

    最后一部分是消息复用。为此,我们将讨论可能的发送方类impl:

    class Sender {
        UnicastProcessor<> outgoing = ...
        UnicastPorcessor<> incoming = ...
    
        Map<String, MonoSink> monoSinksMap = ...;
        Map<String, FluxSink> fluxSinksMap = ...;
    
        public Sender () {
    

    //在此处创建WebSocket连接并放置前面提到的代码 }

        Mono<R> sendForMono(T data) {
            //generate message with unique 
            return Mono.<R>create(sink -> {
                monoSinksMap.put(streamId, sink);
                outgoing.onNext(message); // send message to server only when subscribed to Mono
            });
        }
    
         Flux<R> sendForFlux(T data) {
             return Flux.<R>create(sink -> {
                fluxSinksMap.put(streamId, sink);
                outgoing.onNext(message); // send message to server only when subscribed to Flux
            });
         }
    }
    

    自定义实现的总结

    1. 中坚分子
    2. 没有实施反压力支持,因此这可能是另一个挑战。
    3. 很容易射到自己的脚

    外卖

    1. 请使用rsocket,不要自己发明协议,这很难!!!!
    2. 从关键人物那里了解更多关于rsocket的信息- https://www.youtube.com/watch?v=WVnAbv65uCU
    3. 从我的一次谈话中了解更多关于rsocket的信息- https://www.youtube.com/watch?v=XKMyj6arY2A
    4. 在rsocket之上构建了一个名为proteus的特色框架-您可能对此感兴趣- https://www.netifi.com/
    5. 从rsocket协议的核心开发人员了解更多关于proteus的信息- https://www.google.com/url?sa=t&source=web&rct=j&url=https://m.youtube.com/watch%3Fv%3D_rqQtkIeNIQ&ved=2ahUKEwjpyLTpsLzfAhXDDiwKHUUUA8gQt9IBMAR6BAgNEB8&usg=AOvVaw0B_VdOj42gjr0YrzLLUX1E
        2
  •  2
  •   Ricard Kollcaku    7 年前

    不确定这是不是你的问题?? 我看到你正在发送一个静态的流量响应(这是一个接近的流) 您需要一个打开的流来向该会话发送消息,例如,您可以创建一个处理器。

    public class SocketMessageComponent {
    private DirectProcessor<String> emitterProcessor;
    private Flux<String> subscriber;
    
    public SocketMessageComponent() {
        emitterProcessor = DirectProcessor.create();
        subscriber = emitterProcessor.share();
    }
    
    public Flux<String> getSubscriber() {
        return subscriber;
    }
    
    public void sendMessage(String mesage) {
        emitterProcessor.onNext(mesage);
    }
    

    }

    然后你可以发送

     public Mono<Void> handle(WebSocketSession webSocketSession) {
        this.webSocketSession = webSocketSession;
        return webSocketSession.send(socketMessageComponent.getSubscriber()
                .map(webSocketSession::textMessage))
                .and(webSocketSession.receive()
                        .map(WebSocketMessage::getPayloadAsText).log());
    }