【问题标题】:Continuous state reduction with Flux使用 Flux 持续减少状态
【发布时间】:2019-03-26 22:53:48
【问题描述】:

假设我有两种事件类型(AB)和Fluxes 以某种方式生成它们:

Flux<A> aFlux = ...;
Flux<B> bFlux = ...;

还有一个保存当前状态的类型,用类型S表示:

class S {
  final int val;
}

我想创建以下内容:

final S sInitial = ...;

Flux<S> sFlux = Flux.merge(aFlux, bFlux)
  .scan((a, e) -> {
    if(e instanceof A) {
      return mapA(a, (A)e);
    } else if(e instanceof B) {
      return mapB(a, (B)e);
    } else {
      throw new RuntimeException("invalid event");
    }
  })
  .startWith(sInitial);

其中sCurr是sFlux最后输出的S的实例,以sInitial开头,mapA/mapB返回S类型的新值。 SsInitial 都是不可变的。

也就是说,我想:

  • 不断输出最新状态...
  • ...正在生成...
  • ...基于当前状态和接收到的事件...
  • ...按照映射器函数的规定

有没有办法以其他方式重新组织上述流流,尤其是为了避免使用instanceof

【问题讨论】:

    标签: java project-reactor


    【解决方案1】:

    您可以添加接口并为您的 A 和 B 类实现它

    interface ToSConvertible {
        S toS(S s);
    }
    

    现在你可以使用reactor.core.publisher.Flux#scan(A, java.util.function.BiFunction&lt;A,? super T,A&gt;) 方法:

    Flux<S> sFlux = Flux.merge(aFlux, bFlux)
            .scan(sInitial, (s, e) -> e.toS(s));
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2015-10-07
      • 1970-01-01
      • 2015-08-06
      • 1970-01-01
      • 2021-09-02
      • 2016-02-25
      • 2018-09-16
      相关资源
      最近更新 更多