【发布时间】:2016-08-19 17:37:12
【问题描述】:
我正在尝试制作一个反应式 RabbitMQ 侦听器,它允许我们处理具有多个订阅者的每条消息。我们只想在所有订阅者成功完成后ack 消息。
这是我目前的设置:
Observable
.fromCallable(() -> {
// Set up connection
return consumer;
})
.flatMap(consumer -> Observable
.fromCallable(consumer::nextDelivery)
.doOnError(throwable -> {
try {
consumer.getChannel().getConnection().close();
} catch (IOException ignored) { }
})
.repeat())
.retryWhen(observable -> observable.delay(3, TimeUnit.SECONDS))
.publish()
.refCount();
这将建立一次连接,与所有订阅者共享所有消息,并在 3 秒后重新连接,如果由于例如兔子变得不可用。
我还需要做的是ack 或nack 消息。由于我们所有的消息处理程序都是幂等的,因此如果任何处理程序发生故障,我可以将消息重新排队,以确保每个处理程序都成功完成。
有没有办法判断 任何 订阅者是否失败?我目前正在考虑订阅这样的内容:
public void subscribe(Action1 action) {
deliveries
.flatMap(delivery -> Observable
.just(delivery)
.doOnNext(action)
.doOnError(throwable -> {
// nack
})
.doOnCompleted(() -> {
// ack
})
)
.subscribe();
}
但这显然acks 或nacks 在第一次失败或成功时。有什么方法可以merge 特定消息的所有订阅者,然后检查错误或完成情况?
我还尝试过使用AtomicInteger 来计算所有订阅者,然后计算成功/失败,但显然,只要有人在处理过程中订阅或取消订阅,就没有简单的方法可以在不阻塞整个处理步骤的情况下进行同步。
我也可以给每个订阅者一个Observable<Delivery> 并让他们返回一个错误或完成,类似于retryWhen(作为一种回复渠道),但我无法生成所需数量的预先观察并在之后合并它们。
有什么想法吗?感谢阅读!
【问题讨论】: