【问题标题】:rxjava drops downstream resultsrxjava 丢弃下游结果
【发布时间】:2017-09-15 22:59:25
【问题描述】:

我有一个 rx 方法来进行 api 调用,这个方法的调用者可能会在很短的时间内多次发生。所以 rx 方法是

public void apiCallWithRx() {
    apiService.makeApiCall()
        .subscribeOn(Schecdulars.io())
        .observeOn(AndroidSchedulers.mainTread())
        .subscribe(
           // onNext
           new onConsume(),
          // onError
           new onConsume(),

         );
} 

调用者方法可以在短时间内多次调用此 apiCallWithRx.. 但问题是,从第二次或任何特定时间调用时,我有时无法从下游得到响应。不会调用 onNext、onError 或 onComplete。 我想知道,这是因为缓冲还是背压.. 尝试使用 rxjava1 和 rxjava2,它们是相同的。

如果您有任何建议,我将不胜感激。

更新 1

我没有看到任何背压异常,所以它不可能是背压问题。

更新 2

请忽略细节,Rx 代码大部分时间都有效。我只是为了说明而省略了一些代码

更新 3

我在后台有一个BlockingQueue,所以这个rx方法实际上是在队列中有可用数据时调用的。可以随时将数据添加到队列中。而这个rx方法不是异步调用的,因为这个方法是在第一次响应之后才调用的,然后检查队列,如果有数据,那么我们发送第二个api请求。

【问题讨论】:

  • 在 RxJava 中使用 BlockingQueue 很容易出现死锁。您可能需要一个 UnicastSubject 来缓冲数据,直到下游可以使用它。
  • @akamokd apiCallWithRx()方法是从android的UI线程调用的,所以ExecutorService在后台线程从BlockingQueue获取数据,并将数据传递给UI线程,并触发apiCallWithRx()方法在 UI 线程上。每当我们从服务器获取先前请求数据的 api 响应时,它将检查阻塞队列以获取下一个请求数据。所以 api 调用和 BlockingQueue 是相当分离的,我认为这里没有死锁

标签: rx-java rx-java2


【解决方案1】:

Observable 默认是惰性的,你必须subscribe 才能运行Observable

【讨论】:

  • 是的,我们有那个代码,为了便于说明,它被省略了。
  • 您的代码使用了subscribe 的版本,它触发了可观察但忽略了所有事件。
  • 为了说明,省略了那个代码,真实的代码有,只是添加了onNext,onError
【解决方案2】:

我不知道为什么在 RxJava2 中 subscribeOnobserverOn 运算符可以在不指定必须执行的线程的情况下使用。

但至少在 RxJava1 中,如果您使用这些运算符,您会异步运行它,然后您必须等待执行才能获得结果。

检查这个测试

@Test
public void testObservableAsync() throws InterruptedException {
    Subscription subscription = Observable.from(numbers)
            .doOnNext(increaseTotalItemsEmitted())
            .subscribeOn(Schedulers.newThread())
            .subscribe(number -> System.out.println("Items emitted:" + total));
    System.out.println("I finish before the observable finish.  Items emitted:" + total);
    new TestSubscriber((Observer) subscription)
            .awaitTerminalEvent(100, TimeUnit.MILLISECONDS);
}

你可以在这里看到另一个例子https://github.com/politrons/reactive/blob/master/src/test/java/rx/observables/scheduler/ObservableAsynchronous.java

【讨论】:

  • 更新问题,我使用 BlockingQueue 来确保 rx 方法被顺序调用
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多