【发布时间】:2019-11-03 22:03:03
【问题描述】:
如何在不使用 reduce() 或 last() 折叠 Flux 的情况下获得最后一项?这是我的用例:
1) 我有基于状态生成Flux<T> 的生成器。
2) 当内部Flux 完成时,它会改变影响我在生成器中发出的下一个Flux 对象的状态。
大概是这个样子
static class State {
int secret = 2;
int iteration = 0;
}
Random rand = new Random(1024);
Flux<Integer> stream = Flux.<Flux<Integer>, State>generate(State::new, (state, sink) -> {
System.out.println(String.format("Generate: %d", state.secret));
Flux<Integer> inner = Flux.range(1, rand.nextInt(10));
sink.next(inner.doOnComplete(() -> {
// How do I get last item of `inner` here ?
// For example I'd like to decrement `state.secret` by last value of `inner`
}));
return state;
}).flatMap(Function.identity());
UPD:我没有标记我的答案,因为黑客被证明是不可靠的。 .generate() 可能会在之前的 Flux 被完全消耗之前被调用,因此导致 last 中的值不正确。
【问题讨论】:
-
您可以在您的内部流量上调用
share()以允许多个订阅者,这样您就可以毫无问题地调用last()吗? -
@MichaelBerry 会不会因为源现在很热而增加流开始时丢失项目的风险?
-
可能,是的。如果你可以
cache()整个事情是可行的,那么它并不重要,但如果不是(并且缺少元素很重要),那么这不是我推荐的方法。 -
因为我现在没有东西可以测试它。你想试试流的物化吗?因为它会将其转换为信号+上下文它可能具有“完成”信号中的最后一个值。
标签: java project-reactor reactive-streams