【发布时间】: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