【问题标题】:can we publish a stream more then once?我们可以多次发布流吗?
【发布时间】:2018-08-09 23:21:53
【问题描述】:

以下所有代码均不打印任何内容。为什么?

ConnectableFlux<Integer> publish = Flux.just(1)
        .publish();

ConnectableFlux<Integer> publish1 = Flux.just(2)
        .flatMap(x -> publish)
        .publish();

publish1.subscribe(System.out::println, System.out::println, System.out::println);
publish1.connect();

ConnectableFlux<Integer> publish1 = Flux.just(2)
        .publish()
        .publish();

publish1.subscribe(System.out::println, System.out::println, System.out::println);
publish1.connect();

ConnectableFlux<Integer> publish1 = Flux.just(2)
        .publish()
        .doOnNext(System.out::println)
        .publish();

publish1.subscribe(System.out::println, System.out::println, System.out::println);
publish1.connect();

【问题讨论】:

  • 一个流用完就被消费了。
  • connectableFlux 就是这样工作的。

标签: java project-reactor


【解决方案1】:

不要忘记为每个ConnectableFlux 提供一个.connection

在所有这些示例中,都缺少.connection 声明。

对于第一种情况,要使其正常工作,我们必须先 .connect 到第一个 publish ConnectableFlux :

ConnectableFlux<Integer> publish = Flux.just(1)
        .publish();

ConnectableFlux<Integer> publish1 = Flux.just(2)
        .flatMap(x -> publish)
        .publish();

publish1.subscribe(System.out::println, System.out::println, System.out::println);
publish1.connect();
publish.connect();

对于以下两个示例,我们有类似的东西。当我们使用Flux.just(...).publish().publish() 时,我们创建了两个ConnectableFlux。这里的问题是第一个被删除了。如果必须有后续的.publishing(这很不合逻辑),我们可以使用以下技术来避免擦除以前的ConnectableFluxes:

ConnectableFlux<Integer> publish1 = Flux.just(2)
        .publish()
        .autoConnect() // or .autoConnect(0)
        .doOnNext(System.out::println)
        .publish();

publish1.subscribe(System.out::println, System.out::println, System.out::println);
publish1.connect();

在该示例中,我们使用.autoConnect() 运算符,在.autoConnect(0) 的情况下,它只是ConnectableFlux#connect 和return this; 语句的组合。在.autoConnect(&gt;0) 的情况下,使用了一些对初始源的惰性订阅,这听起来像“当且仅当我们获得 N 个订阅者时才连接到初始源”

【讨论】:

  • 太棒了!谢谢阿吉安奥莱:)
猜你喜欢
  • 2014-06-12
  • 2016-03-27
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2012-05-04
  • 1970-01-01
  • 2013-05-18
  • 1970-01-01
相关资源
最近更新 更多