如何使用现有的 Flux 运算符使 Flux 将传入值返回到多个列表中,并且返回之间的延迟最小?
问问题
1180 次
2 回答
1
这可以通过一组重要的组合运算符来实现。
import java.time.Duration;
import java.util.*;
import reactor.core.publisher.*;
public class DelayedBuffer {
public static void main(String[] args) {
Flux.just(1, 2, 3, 6, 7, 10)
.flatMap(v -> Mono.delayMillis(v * 1000)
.doOnNext(w -> System.out.println("T=" + v))
.map(w -> v)
)
.compose(f -> delayedBufferAfterFirst(f, Duration.ofSeconds(2)))
.doOnNext(System.out::println)
.blockLast();
}
public static <T> Flux<List<T>> delayedBufferAfterFirst(Flux<T> source, Duration d) {
return source
.publish(f -> {
return f.take(1).collectList()
.concatWith(f.buffer(d).take(1))
.repeatWhen(r -> r.takeUntilOther(f.ignoreElements()));
});
}
}
(但请注意,由于涉及时间,预期的发射模式可能与自定义运算符更好地匹配。)
于 2017-01-30T14:28:35.280 回答
0
我以为buffer(Duration)
会满足您的需求,但事实并非如此。
编辑:留下这个以防万一有你完全相同需求的人想使用该运算符。这种缓冲区变体将序列拆分为连续的时间窗口(每个时间窗口产生一个buffer
)。也就是说,新的delay
开始于前一个的末尾,而不是每当发出新的延迟元素时。
于 2017-01-27T11:59:00.073 回答