【问题标题】:How to have multiple subscribers to Flux that run on different execution contexts / threads如何让 Flux 的多个订阅者在不同的执行上下文/线程上运行
【发布时间】:2020-03-14 21:30:43
【问题描述】:

我正在开发一个用于物联网实时数据可视化的 Spring Boot WebFlux 应用程序。

我有一个Flux,它模拟来自设备的数据,我希望在每个事件都建立 websocket 连接时:

  • 必须通过 websocket 发送实时可视化(使用响应式WebSocketHandler
  • 必须根据给定条件进行检查,以便通过 HTTP REST 调用 (RestTemplate) 发送通知

从我的日志看来,两个订阅者(websocket 处理程序和通知程序)获得了两个具有完全不同值的不同流(在日志下方)。

我还尝试了在 MySource 类中的 map 之后链接 share 方法的变体,在这种情况下,看起来虽然我只有一个 Flux,但只有一个线程,所以一切都在阻塞(我可以看到 REST 调用阻止了通过 websocket 发送)。

这里发生了什么?我怎样才能让两个订阅者在不同的执行上下文(不同的线程)中运行,从而完全相互独立?

下面是相关代码sn-ps和logs。

提前谢谢大家!

更新: 为清楚起见,我必须指定 MyEvents 具有随机生成的值,因此由于 @NikolaB 的回答,我通过使用 ConnectableFlux / @ 解决了一个问题987654329@ 保证具有相同的Flux,但我仍然希望为两个订阅者提供单独的执行上下文。

public class MyWebSocketHandler implements WebSocketHandler {

   @Autowired
   public MySource mySource;

   @Autowired
   public Notifier notifier;

   public Mono<Void> handle(WebSocketSession webSocketSession) {
            Flux<MyEvent> events = mySource.events();
            events.subscribe(event -> notifier.sendNotification(event));
            return webSocketSession.send(events.map(this::toJson).map(webSocketSession::textMessage));
   }

   private String toJson(MyEvent event) {
       log.info("websocket toJson " + event.getValue());
       ...
   }
}
public class MySource {
   public Flux<MyEvent> events() {
      return Flux.interval(...).map(i -> new MyEvent(*Random Generate Value*);
   }
}
public class Notifier {

   public void sendNotification (MyEvent event) {
      log.info("notifier sendNotification " + event.getValue());
      if (condition met)
         restTemplate.exchange(...)
   }
}
2019-11-19 11:58:55.375 INFO [     parallel-3] i.a.m.websocket.MyWebSocketHandler  : websocket toJson 4.09
2019-11-19 11:58:55.375 INFO [     parallel-1] i.a.m.notifier.Notifier : notifier sendNotification 4.86
2019-11-19 11:58:57.366 INFO [     parallel-1] i.a.m.notifier.Notifier : notifier sendNotification 4.24
2019-11-19 11:58:57.374 INFO [     parallel-3] i.a.m.websocket.MyWebSocketHandler  : websocket toJson 4.11
2019-11-19 11:58:59.365 INFO [     parallel-1] i.a.m.notifier.Notifier : notifier sendNotification 4.61
2019-11-19 11:58:59.374 INFO [     parallel-3] i.a.m.websocket.MyWebSocketHandler  : websocket toJson 4.03
2019-11-19 11:59:01.365 INFO [     parallel-1] i.a.m.notifier.Notifier : notifier sendNotification 4.88
2019-11-19 11:59:01.375 INFO [     parallel-3] i.a.m.websocket.MyWebSocketHandler  : websocket toJson 4.29
2019-11-19 11:59:03.364 INFO [     parallel-1] i.a.m.notifier.Notifier : notifier sendNotification 4.37

【问题讨论】:

    标签: java resttemplate spring-webflux project-reactor spring-websocket


    【解决方案1】:

    这里有几个问题,首先RestTemplate 是同步/阻塞HTTP 客户端,所以你应该使用WebClient 这是反应式的,也可以创建ConnectableFluxFlux 可以有多个订阅者)你需要在map 运算符之前分享它并创建新的Flux-es,它是从连接的一个创建的。

    例子:

    Flux<MyEvent> connectedFlux = mySource.events().share();
    Flux.from(connectedFlux).subscribe(event -> notifier.sendNotification(event));
    return webSocketSession.send(Flux.from(connectedFlux).map(this::toJson).map(webSocketSession::textMessage));
    

    此外,sendNotification 方法应该返回 Mono&lt;Void&gt;,因为响应式方法应该总是返回 MonoFluxtypes。

    要启动独立执行,您可以Zip 这两个Monos。

    编辑

    首先如上所述,将WebClient 用于传出HTTP 调用,这是响应式HTTP 客户端并重新编写Notifier 类:

    public class Notifier {
    
       public Mono<Void> sendNotification (MyEvent event) {
          log.info("notifier sendNotification " + event.getValue());
          return Mono.just(event)
                     .filter(e -> /* your condition */)
                     .flatMap(e -> WebClient.builder().baseUrl("XXX")...)
                     .then();
       }
    
    }
    

    现在看看执行上下文是否不同。

    【讨论】:

    • 我设法使用了ConnectableFlux / share,所以现在可以保证单个 Flux 源,但问题是 sendNotification 是由执行 websocket 内容的同一线程调用的。我想要实现的是订阅者执行上下文的分离。可能吗?我觉得我必须在某处添加subscribeOn,但我不确定......
    • 我已经编辑了我的问题主体以更好地澄清这一点。
    • @vortex.alex 是的,这是可能的,但不是 subscribeOnpublishOn 因为 WebClient 和响应式 WebSocket 都有自己的非阻塞线程,而您的 subscribeOnpublishOn 将是被它们的内部执行定义覆盖。在我的答案中添加了更多代码。
    • 我试过了,但要么我没有得到WebClient 的任何回复,要么我最终在exchangeretrieve 之后放了一个block,但它不起作用,因为我有以下错误:block()/blockFirst()/blockLast() are blocking, which is not supported in thread parallel-2。我是否应该在WebClient 上明确订阅以触发请求(例如使用block)?或者,exchange 返回的 Mono 可能必须通过 sendNotification 方法返回,但那该怎么办呢?你提到的zip方法有什么关系吗?
    • 绝对你必须调用subscribe 来启动响应式执行,所以notifier.sendNotification(event).subscribe() 将启动执行。您不必在 Controllers 中返回的反应式方法/操作符链上调用 subscribe,因为 Spring 隐式调用它,但对于所有其他执行,您需要显式调用它。 zipWith 运算符可用于独立的反应式执行,当您需要下游多个独立反应式方法执行的结果以进行进一步操作时。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-08-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多