【问题标题】:Check if any subscriber throws an exception in RxJava检查是否有订阅者在 RxJava 中抛出异常
【发布时间】: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 秒后重新连接,如果由于例如兔子变得不可用。

我还需要做的是acknack 消息。由于我们所有的消息处理程序都是幂等的,因此如果任何处理程序发生故障,我可以将消息重新排队,以确保每个处理程序都成功完成。

有没有办法判断 任何 订阅者是否失败?我目前正在考虑订阅这样的内容:

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(作为一种回复渠道),但我无法生成所需数量的预先观察并在之后合并它们。

有什么想法吗?感谢阅读!

【问题讨论】:

    标签: java rx-java reactivex


    【解决方案1】:

    您可以使用 onErrorResumeNext 来控制从您的管道传播的异常并设置 nack,然后将您的 onComplete 用作 ack

    这里是一个例子

    Observable.just(null).map(Object::toString)
                         .doOnError(failure -> System.out.println("Error:" + failure.getCause()))
                         .retryWhen(errors -> errors.doOnNext(o -> count++)
                         .flatMap(t -> count > 3 ? Observable.error(t) : Observable.just(null).delay(100, TimeUnit.MILLISECONDS)),
                                                         Schedulers.newThread())
                         .onErrorResumeNext(t -> {
                                                  System.out.println("Error after all retries:" + t.getCause());
                                                  return Observable.just("I save the world!");
                                              })
                         .subscribe(s -> System.out.println(s));
    

    【讨论】:

    • 谢谢。这是否考虑到来自多个订阅者的潜在错误?
    • 我真的不知道如何应用它。我将onErrorResumeNext 放在refCount 之前还是之后?我怎么知道所有订阅者都成功地处理了消息以便ack呢?我不希望订阅者关闭,因为会有来自rabbit的后续消息。
    【解决方案2】:

    您想使用Observable.mergeDelayError,然后使用.onError* 方法之一。

    只有在所有可观察对象都完成/出错后,第一个才会传播错误;第二个将允许您在处理完成后处理错误

    编辑:要获得计数,请计算成功次数:

    Object message = ...;
    List<Action1<?>> actions = ...;
    Observable.from(actions)
     .map(action ->
        Observable.defer(() -> Observable.just(action.call(message))
        .ignoreEmements()
        .cast(Integer.class)
        .switchIfEmpty(1)
        .onErrorReturnResumeNext(e->Observable.empty())
     )
     .compose(Observable::merge)
     .count();
    

    这有点令人费解,但可以更清楚地说明:拨打电话,忽略错误,计算成功。

    【讨论】:

    • 是的,这听起来是对的。我仍在努力寻找一种方法来产生正确数量的 observables。我怎么知道有多少订阅者会收到消息?只有 refCount 有准确的订阅者号码 AFAIK,这会阻止每条消息这样做。
    • 感谢您的编辑。如果有人在处理流时订阅,它不会收到消息,但计数会包含它并且不同步?我没有订阅者列表,我让RxJava 处理它,但是看看你的示例,我可能必须自己跟踪订阅者并明确分派给他们以完成这项工作。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-12-11
    • 1970-01-01
    • 2016-12-29
    • 2021-03-16
    • 2011-09-21
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多