【问题标题】:Handling errors in FluxSink that continuously emits events处理 FluxSink 中不断发出事件的错误
【发布时间】:2017-08-23 16:04:45
【问题描述】:

我们需要监听来自队列的消息并对传入的对象执行一系列操作。

我们计划在这个用例中使用 spring reactor 框架,我们可以在 Flux 中将传入的对象(来自队列)作为事件发出,并在 Flux 上设置一系列映射操作,使对象流动通过以完成处理。

为了实现这一点,我们正在尝试使用 FluxSink,如下所示:

Flux<ObjectFromQueue> flux = Flux.create(fluxSink -> { 
    messageListener.setFluxSink(fluxSink); 
    //messageListener has a 'process' method that will be invoked as soon as there is a new message on the Queue.
}, FluxSink.OverflowStrategy.BUFFER);

ConnectableFlux<ObjectFromQueue> connectableFlux = flux.publish();
connectableFlux.map(o -> {
        return handler1.handle(o);
    }).map(o -> {
        return handler2.handle(o);
    }).doOnError(t -> {
        errorHandler.handle(t)
    }).subscribe()
connectableFlux.connect();

MessageListener 实现了一个 Listener 接口(一个自定义框架类)并覆盖了 'process(ObjectFromQueue)' 方法,一旦队列上有消息并且在 'process' 方法内部,我们将调用该方法

fluxSink.next(objectFromQueue);

通过这种配置,我们可以实现我们的要求,即一旦队列上有消息,MessageListener 就会收到消息,并且消息会通过配置的 handler1 和 handler2 操作符 但是以防出现任何错误在处理消息期间,由于 onError 是一个终端事件,因此通量停止运行(即发送到预配置的运算符)

解决此问题的最佳方法是什么?还是不推荐这种配置 Flux 的方式?如果不推荐,建议采用什么方法来实现此要求?

【问题讨论】:

    标签: project-reactor


    【解决方案1】:

    现在你有两种选择:

    • 使用 concatMap 代替 map 来获得一个可以应用错误处理的按值子发布者:

      .concatMap(v -> Mono.just(handler1.handle(v))
                          .doOnError(errorHandler::handle)
                          .onErrorResume(Mono.empty())
      )
      
      • 在上面的示例中,带有空 Mono 的 onErrorResume 归结为删除错误值。
      • 由于您仍在 concatMap 中,因此可以访问v,您还可以应用返回默认值的错误处理程序(但永远不会 null
    • 使用handle:它提供了一个sink,如果应该删除该值,您可以跳过使用它。

      • 通常您会执行 try/catch 并调用错误处理程序而不是 catch 块中的接收器。
      • 您可以更改错误处理程序类,使其接受一个值、一个处理程序和一个接收器,并且它会隐藏 try/catch。比如:

        errorHandler.safeHandle(v, sink, handler1);
        
        public abstract class ErrorHandler<T> {
        
            public <V> V safeHandle(T value, SynchronousSink<V> sink, Handler<T,V> handler) {
              try {
                V mapped = handler.handle(value);
                sink.next(mapped);
              } catch (Throwable t) {
                //process error, eg. log or increment metrics
                processError(value, t);
              }
            }
        
            public abstract void processError(T valueCause, Throwable error);
        }
        

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2022-12-14
      • 2015-07-07
      • 1970-01-01
      • 2019-10-20
      • 2013-05-25
      • 1970-01-01
      • 2019-01-08
      • 2019-05-04
      相关资源
      最近更新 更多