【问题标题】:Spring 5 reactive websockets: Clients not receiving same data from hot streamSpring 5 反应式 websockets:客户端没有从热流中接收到相同的数据
【发布时间】:2019-03-06 12:48:35
【问题描述】:

我的WebSocketHandler 实现中有这个:

@Override
public Mono<Void> handle(WebSocketSession session) {

    return session.send(
       session.receive()
              .flatMap(webSocketMessage -> {
                  int id = Integer.parseInt(webSocketMessage.getPayloadAsText());

                  Flux<EfficiencyData> flux = service.subscribeToEfficiencyData(id);
                  var publisher = flux
                      .<String>handle((o, sink) -> {
                         try {
                            sink.next(objectMapper.writeValueAsString(o));
                         } catch (JsonProcessingException e) {
                            e.printStackTrace();                               
                         }
                      })
                      .map(session::textMessage);

                  return publisher;
              })
    );
}

Flux&lt;EfficiencyData&gt; 当前在服务中生成用于测试,如下所示:

public Flux<EfficiencyData> subscribeToEfficiencyData(long weavingLoomId) {
    return Flux.interval(Duration.ofSeconds(1))
               .map(aLong -> {
                   longAdder.increment();
                   return new EfficiencyData(new MachineSpeed(
                           RotationSpeed.ofRpm(longAdder.intValue()),
                           RotationSpeed.ofRpm(0),
                           RotationSpeed.ofRpm(400)));
               }).publish().autoConnect();
}

我正在使用publish().autoConnect() 使其成为热门流。我创建了一个单元测试,它在返回的Flux 上启动 2 个线程:

flux.log().handle((s, sink) -> {
            LOGGER.info("{}", s.getMachineSpeed().getCurrent());
        }).subscribe();

在这种情况下,我看到两个线程每秒都打印出相同的值。

但是,当我打开 2 个浏览器选项卡时,我在两个网页中看不到相同的值。连接的 websocket 客户端越多,值之间的差异就越大(因此原始 Flux 中的每个值似乎都发送到不同的客户端,而不是发送到所有客户端)。

【问题讨论】:

    标签: java spring-webflux project-reactor


    【解决方案1】:

    感谢Brian Clozel on twitter,设法解决了这个问题。

    问题是对于每个连接的 websocket 客户端,我调用service.subscribeToEfficiencyData(id) 方法,每次调用它都会返回一个 new Flux。所以当然,这些独立的 Flux 不会在不同的 websocket 客户端之间共享。

    为了解决这个问题,我在构造函数中创建了 Flux 实例或我的服务的 PostConstruct 方法,因此 subscribeToEfficiencyData 每次都返回相同的 Flux 实例。

    请注意,Flux 上的 .publish().autoConnect() 仍然很重要,因为没有它,websocket 客户端将再次看到不同的值!

    【讨论】:

      猜你喜欢
      • 2020-05-01
      • 1970-01-01
      • 2018-04-22
      • 1970-01-01
      • 2013-08-17
      • 1970-01-01
      • 1970-01-01
      • 2021-09-30
      • 1970-01-01
      相关资源
      最近更新 更多