【问题标题】:How can I aggregate elements on a flux by group / how to reduce groupwise?如何按组聚合通量上的元素/如何按组减少?
【发布时间】:2019-06-18 01:31:30
【问题描述】:

假设您有一系列具有以下结构的对象:

class Element {
  String key;
  int count;
}

现在想象这些元素以预定义的排序顺序流动,总是以一组键的形式出现,例如

{ key = "firstKey",  count=123}
{ key = "firstKey",  count=1  }
{ key = "secondKey", count=4  }
{ key = "thirdKey",  count=98 }
{ key = "thirdKey",  count=5  }
 .....

我想要做的是创建一个通量,它为每个不同的 key 返回一个元素,并为每个键组求和 count。 所以基本上就像每个组的经典减少,但使用 reduce 运算符不起作用,因为它只返回一个元素,我想为每个不同的键获得一个元素的通量。

使用bufferUntil 可能有效,但有一个缺点,即我必须保持一个状态来检查key 与前一个相比是否发生了变化。

使用groupBy 有点过头了,因为我知道一旦找到新密钥,每个组都结束了,所以我不想在该事件之后保留任何缓存。

是否可以使用Flux 进行这样的聚合,而无需在流之外保持状态?

【问题讨论】:

    标签: java flux project-reactor


    【解决方案1】:

    如果不自己跟踪状态,目前(从 3.2.5 开始)这是不可能的。 distinctUntilChanged 可以满足最小状态的要求,但不会发出状态,只是根据所述状态它认为“不同”的值。

    解决这个问题的最简单的方法是使用windowUntilcompose + 一个AtomicReference 用于每个订阅者的状态:

    Flux<Tuple2<T, Integer>> sourceFlux = ...; //assuming key/count represented as `Tuple2`
    Flux<Tuple2<T, Integer>> aggregated = sourceFlux.compose(source -> {
        //having this state inside a compose means it will not be shared by multiple subscribers
        AtomicReference<T> last = new AtomicReference<>(null);
    
        return source
          //use "last seen" state so split into windows, much like a `groupBy` but with earlier closing
          .windowUntil(i -> !i.getT1().equals(last.getAndSet(i.getT1())), true)
          //reduce each window
          .flatMap(window -> window.reduce((i1, i2) -> Tuples.of(i1.getT1(), i1.getT2() + i2.getT2()))
    });
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-09-21
      • 2019-11-22
      • 1970-01-01
      • 2022-01-06
      • 2017-03-25
      • 1970-01-01
      • 2013-01-22
      • 1970-01-01
      相关资源
      最近更新 更多