【问题标题】:How to handle exceptions thrown in subscriptions to processors in Project Reactor如何处理在 Project Reactor 中订阅处理器时引发的异常
【发布时间】:2019-05-03 10:07:37
【问题描述】:

考虑以下测试:

@Test
void test() {
    DirectProcessor<Object> objectTopicProcessor = DirectProcessor.create();

    Runnable r = mock(Runnable.class);

    objectTopicProcessor.subscribe(next -> {throw new RuntimeException("eee");});
    objectTopicProcessor.subscribe(next -> r.run());

    assertThrows(RuntimeException.class, () -> objectTopicProcessor.onNext("")); // exception is thrown

    verify(r).run(); // it's not run
}

想象一下,我构建了一个 API,将处理器公开给客户端。 当某人有多个订阅并且其中一个抛出异常时,不会执行其他调用。此外,异常会从objectTopicProcessor.onNext("") 传播和抛出。我想阻止这种行为。

我知道客户可以将他的代码包装在订阅内的 try-catch 中,但是还有其他方法吗?例如,有时可能会发生 NullPointer,或者客户端可能会忘记检查异常。对于 API,强制客户端尝试捕获所有异常也很不方便。

处理此类情况的最佳策略是什么?

【问题讨论】:

    标签: java spring project-reactor


    【解决方案1】:

    在本例中,传递给subscribe 方法的代码默认在主线程上执行。它首先遇到异常并立即失败,而不执行第二个subscribe 块。
    为了实现并行,使用.publishOn(scheduler)方法:

    @Test
    void test() {
        DirectProcessor<Object> processor = DirectProcessor.create();
        Flux<Object> flux = processor.publishOn(Schedulers.parallel());
        Runnable r = mock(Runnable.class);
        flux.subscribe(next -> {throw new RuntimeException("eee");});
        flux.subscribe(next -> r.run());
    
        processor.onNext(""); // onNext no longer throws an exception
    
        verify(r, timeout(1000)).run();
    }
    

    【讨论】:

      猜你喜欢
      • 2017-05-18
      • 2022-01-09
      • 2017-07-29
      • 2015-07-17
      • 1970-01-01
      • 1970-01-01
      • 2022-12-10
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多