【发布时间】:2020-01-07 20:09:11
【问题描述】:
我有无限的热数据流。我即将对流中的每个元素执行一个操作,每个元素都返回一个 Mono,它将在有限的时间后(以一种或另一种方式)完成。
这些操作可能会引发错误。如果是这样,我想重新订阅热通量而不会丢失任何内容,重试抛出错误时正在处理的元素(即任何未成功完成的元素)。
我在这里做什么?我可以容忍对相同元素的重复操作,但不能完全从流中丢失元素。
我尝试使用 ReplayProcessor 来处理这个问题,但是如果不重复很多可能已经成功的操作(使用非常保守的超时),或者由于丢失元素,我看不到让它工作的方法覆盖缓冲区中旧元素的新元素(如下所示)。
测试用例:
@Test
public void fluxTest() {
List<String> strings = new ArrayList<>();
strings.add("one");
strings.add("two");
strings.add("three");
strings.add("four");
ConnectableFlux<String> flux = Flux.fromIterable(strings).publish();
//Goes boom after three uses of its method, otherwise
//returns a mono. completing after a little time
DangerousClass dangerousClass = new DangerousClass(3);
ReplayProcessor<String> replay = ReplayProcessor.create(3);
flux.subscribe(replay);
replay.flatMap(dangerousClass::doThis)
.retry(1)
.doOnNext(s -> LOG.info("Completed {}", s))
.subscribe();
flux.connect();
flux.blockLast();
}
public class DangerousClass {
Logger LOG = LoggerFactory.getLogger(DangerousClass.class);
private int boomCount;
private AtomicInteger count;
public DangerousClass(int boomCount) {
this.boomCount = boomCount;
this.count = new AtomicInteger(0);
}
public Mono<String> doThis(String s) {
return Mono.fromSupplier(() -> {
LOG.info("doing dangerous {}", s);
if (count.getAndIncrement() == boomCount) {
LOG.error("Throwing exception from {}", s);
throw new RuntimeException("Boom!");
}
return s;
}).delayElement(Duration.ofMillis(600));
}
}
打印出来:
doing dangerous one
doing dangerous two
doing dangerous three
doing dangerous four
Throwing exception from four
doing dangerous two
doing dangerous three
doing dangerous four
Completed four
Completed two
Completed three
一个永远不会完成。
【问题讨论】:
-
我想出了一个强硬的解决方案,它使用 WorkQueueProcessor 而不是重放处理器和 SynchronizedCollection - 在执行危险任务之前,将每个元素添加到集合中。在平面图的下游,现在已经确保成功,每个元素都被删除。发生错误时,元素会从集合中移至处理器以进行重试。这看起来既笨重又阻塞,我希望有一些更“被动”的方式来解决这个问题。
标签: java project-reactor reactor