【问题标题】:Connecting a Flux/Publisher/Subscriber to a continual generator将 Flux/Publisher/Subscriber 连接到持续生成器
【发布时间】:2019-01-23 12:24:39
【问题描述】:

我有一个不断生成和存储新数据值的类(使用线程池)。我想为客户端代码(“订阅者”)提供一种访问(连接到)新数据值序列的方法。但是,如果我的班级没有客户,或者所有客户都已完成从序列中读取,我希望它继续生成和存储新值而不会停止。当客户端连接到该序列时,它会接收新生成的值,但不会接收过去生成的值。哪个 Project Reactor 类(或多个类)适合这样做?

我想我需要使用Flux 来表示新值的序列,但是要使用哪个Flux 类(或工厂方法)?

【问题讨论】:

    标签: java reactive-programming project-reactor


    【解决方案1】:

    使用 DirectProcessor

    据我所知,无论是否有订阅者,都需要订阅上游的能力。

    这可以在DirectProcessor 的支持下实现。由于ProcessorPublisherSubscriber 的组合,它可以“运行”上游并持续监听传入的信号。同时,DirectProcessor 启用消息多路分解,或简单地将消息广播给所有可用的下游订阅者(如果他们正在侦听)。

    例如,让我们考虑以下代码示例:

    Flux<Long> intervalFlux = Flux.interval(Duration.ofMillis(500)).log("upstream");
    DirectProcessor processor = DirectProcessor.create();
    
    intervalFlux.subscribe(processor);
    
    Thread.sleep(2000);
    
    Disposable downstream1 = processor.log("downstream1")
                                      .subscribe();
    
    Thread.sleep(1000);
    
    downstream1.dispose();
    
    Thread.sleep(1000);
    
    Disposable downstream2 = processor.log("downstream2")
                                      .subscribe();
    Thread.sleep(2000);
    

    如我们所见,我们使用处理器订阅了上游,因此间隔Flux 开始生成数据。然后我们订阅了processor 并等待了 1 秒,因此下游 1 应该观察到两个事件,log("upstream") 操作员通常会记录 6 个事件。在那之后,我们取消了订阅,所以downstream1 订阅者应该停止观察任何事件,但log("upstream") 应该仍然观察间隔。然后,又一次暂停后,我们用另一个 downstrea2 订阅者订阅了流,该订阅者应该观察另外四个事件。

    上述代码的一般输出如下:

    2019-01-23 15:09:04,246 INFO upstream [main] onSubscribe(FluxInterval.IntervalRunnable)
    2019-01-23 15:09:04,249 INFO upstream [main] request(unbounded)
    2019-01-23 15:09:04,757 INFO upstream [parallel-1] onNext(0)
    2019-01-23 15:09:05,252 INFO upstream [parallel-1] onNext(1)
    2019-01-23 15:09:05,751 INFO upstream [parallel-1] onNext(2)
    2019-01-23 15:09:06,252 INFO upstream [parallel-1] onNext(3)
    2019-01-23 15:09:06,258 INFO downstream1 [main] onSubscribe(DirectProcessor.DirectInner)
    2019-01-23 15:09:06,258 INFO downstream1 [main] request(unbounded)
    2019-01-23 15:09:06,754 INFO upstream [parallel-1] onNext(4)
    2019-01-23 15:09:06,755 INFO downstream1 [parallel-1] onNext(4)
    2019-01-23 15:09:07,254 INFO upstream [parallel-1] onNext(5)
    2019-01-23 15:09:07,254 INFO downstream1 [parallel-1] onNext(5)
    2019-01-23 15:09:07,263 INFO downstream1 [main] cancel()
    2019-01-23 15:09:07,755 INFO upstream [parallel-1] onNext(6)
    2019-01-23 15:09:08,255 INFO upstream [parallel-1] onNext(7)
    2019-01-23 15:09:08,265 INFO downstream2 [main] onSubscribe(DirectProcessor.DirectInner)
    2019-01-23 15:09:08,265 INFO downstream2 [main] request(unbounded)
    2019-01-23 15:09:08,755 INFO upstream [parallel-1] onNext(8)
    2019-01-23 15:09:08,756 INFO downstream2 [parallel-1] onNext(8)
    2019-01-23 15:09:09,255 INFO upstream [parallel-1] onNext(9)
    2019-01-23 15:09:09,256 INFO downstream2 [parallel-1] onNext(9)
    2019-01-23 15:09:09,751 INFO upstream [parallel-1] onNext(10)
    2019-01-23 15:09:09,751 INFO downstream2 [parallel-1] onNext(10)
    2019-01-23 15:09:10,255 INFO upstream [parallel-1] onNext(11)
    2019-01-23 15:09:10,255 INFO downstream2 [parallel-1] onNext(11)
    

    正如我们所见,DirectProcessor 启用了所需的行为,因此它可能很适合那里。

    注意

    DirectProcessor 不支持背压,所以如果背压很重要,可以使用limitRate operator 运算符。

    另见

    https://projectreactor.io/docs/core/release/reference/#_direct_processor

    【讨论】:

      猜你喜欢
      • 2015-12-11
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-04-12
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多