【问题标题】:Project Reactor's Flux get failed item during handling errorProject Reactor 的 Flux 在处理错误期间获取失败项目
【发布时间】:2018-09-01 19:07:46
【问题描述】:

我需要处理Flux 流错误,为此我需要知道究竟是什么项目失败了。似乎方法doOnError 应该适合处理错误,但是这样我只能得到异常而不是失败的项目。有什么方法可以同时获取失败项和异常?

private void testFluxIterableFlow() {
    Flux.fromIterable(Arrays.asList(1, 2, 3, 4, 5))
            .map(this::process)
            .doOnError(ex -> {
                ...
            })
            .doOnNext(processedValue -> ...)
            .subscribe();
}

private String process(Integer value) {
    if (value == 4) {
        throw new RuntimeException("error...");
    }
    return "processed " + value;
}

在本例中,我需要在错误处理程序中接收失败项目4 和异常消息。

【问题讨论】:

    标签: java reactive-programming project-reactor


    【解决方案1】:

    这里有点晚了,但万一其他人需要它...您可能想要做的是改用 onErrorContinue。请注意,此方法将从通量中删除错误项并继续处理,因此请务必牢记这一点。如果你想要错误+项目并停止一切,我相信如果你从这个方法重新抛出异常,并且没有其他下游 onErrorContinue,它应该终止通量。这是我使用 kafka 反应式 api 的示例

    此外,传递给 onErrorContinue 的项目将不是通量中的原始项目,而是在反应链中的步骤中尝试处理的对象。例如,如果链在尝试转换 B -> C 时失败,但在转换 A -> B 之后,您将在处理程序中获得的对象将是 B。

    @EventListener(ApplicationReadyEvent.class)
    public void process() {
        receiver.receiveAutoAck()
                .concatMap(Function.identity())
                .flatMap(this::deserialzeRecord)
                .flatMap(this::processEvent)
                .flatMap(responseSender::send)
                .onErrorContinue(this::handleRecordError)
                .subscribe();
    }
    
    private void handleRecordError(Throwable throwable, Object record) {
        log.error("Received error processing record={}", record, throwable);
    }
    

    如果你想终止flux并放弃其余的项目,就这样做

    private void handleRecordError(Throwable throwable, Object record) {
        log.error("Received error processing record={}", record, throwable);
        throw throwable;
    }
    

    【讨论】:

      【解决方案2】:

      如果某个值导致异常,则应将该值视为异常原因。

      由您决定是否值得添加新的异常类型,

      @RequiredArgsConstructor
      public class WrongValueInStreamException extends RuntimeException {
          @Getter
          private final Object wrongValue;
      }
      
      public class StreamProcessor {
          public String process(Integer value) {
              if (value == 4) {
                  throw new WrongValueInStreamException(4);
              }
              return "processed " + value;
          }
      }
      

      但只要异常传达有用的相关信息,这是一种很好的做法:

      .doOnError(WrongValueInStreamException.class::isInstance, e -> {
          final Object value = ((WrongValueInStreamException) e).getWrongValue();
          // use 'value'
      })
      

      【讨论】:

      • 谢谢,它有效,但对我来说,这就像一个解决方法。
      • @VasiliySarzhynskyi 我同意,但我找不到任何适用于此案例的 API
      猜你喜欢
      • 2020-10-31
      • 2021-06-12
      • 2017-05-18
      • 2017-07-29
      • 2018-07-04
      • 2020-09-14
      • 1970-01-01
      • 2019-01-07
      • 2019-08-18
      相关资源
      最近更新 更多