【问题标题】:RxJava: How to get all results AND errors from an ObservableRxJava:如何从 Observable 中获取所有结果和错误
【发布时间】:2016-10-18 09:01:18
【问题描述】:

我正在做一个涉及 Hystrix 的项目,我决定使用 RxJava。现在,剩下的就忘了 Hystrix,因为我认为主要问题是我完全搞砸了正确编写 Observable 代码。

需要: 我需要一种方法来返回一个表示多个可观察对象的可观察对象,每个对象都运行一个用户任务。我希望 Observable 能够返回任务的所有结果,甚至是错误。

问题: 可观察的流死于错误。如果我有三个任务,而第二个任务抛出异常,那么即使第三个任务会成功,我也永远不会收到。

我的代码:

public <T> Observable<T> observeManagedAsync(String groupName,List<EspTask<T>> tasks) {
    return Observable
            .from(tasks)
            .flatMap(task -> {
                try {
                    return new MyCommand(task.getTaskId(),groupName,task).toObservable().subscribeOn(this.schedulerFactory.get(groupName));
                } catch(Exception ex) {
                    return Observable.error(ex);
                }
            });
}

鉴于 MyCommand 是一个扩展 HystrixObservableCommand 的类,它返回一个 Observable,因此不应该考虑我所看到的问题。

尝试 1:

如上所述使用Observable.flatMap

  • 很好:每个命令都安排在自己的线程上,并且任务异步运行。
  • 不好:在第一个命令异常时,Observable 完成了发出先前成功的结果并发出异常。任何进行中的命令都会被忽略。

尝试 2:

使用Observable.concatMapDelayError 而不是flatMap

  • 不好:由于某种原因,任务同步运行。为什么??
  • 很好:我得到了所有成功的结果。
  • ~Good: OnError 得到一个复合异常,其中包含抛出的异常列表。

任何帮助都将不胜感激,并且可能会让我因为自己没有想到而感到非常尴尬。

附加代码

此测试使用Observable.flatMap 成功,但使用Observable.concatMapDelayError 时失败,因为任务没有异步运行:

java.lang.AssertionError:执行时间超过 350ms 限制:608

@Test
public void shouldRunManagedAsyncTasksConcurrently() throws Exception {
    Observable<String> testObserver = executor.observeManagedAsync("asyncThreadPool",getTimedTasks()); 
    TestSubscriber<String> testSubscriber = new TestSubscriber<>();
    long startTime = System.currentTimeMillis();
    testObserver.doOnError(throwable -> {
        System.out.println("error: " + throwable.getMessage());
    }).subscribe(testSubscriber);
    System.out.println("Test execution time: "+(System.currentTimeMillis()-startTime));
    testSubscriber.awaitTerminalEvent();
    long execTime = (System.currentTimeMillis()-startTime);
    System.out.println("Test execution time: "+execTime);
    testSubscriber.assertCompleted();
    System.out.println("Errors: "+testSubscriber.getOnErrorEvents());
    System.out.println("Results: "+testSubscriber.getOnNextEvents());
    testSubscriber.assertNoErrors();
    assertTrue("Execution time ran under the 300ms limit: "+execTime,execTime>=300);
    assertTrue("Execution time ran over the 350ms limit: "+execTime,execTime<=350);
    testSubscriber.assertValueCount(3);
    assertThat(testSubscriber.getOnNextEvents(),containsInAnyOrder("hello","wait","world"));
    verify(this.mockSchedulerFactory, times(3)).get("asyncThreadPool");
}

上述单元测试的任务:

protected List<EspTask<String>> getTimedTasks() {
    EspTask longTask = new EspTask("helloTask") {
        @Override
        public Object doCall() throws Exception {
            Thread.currentThread().sleep(100);
            return "hello";
        }
    };
    EspTask longerTask = new EspTask("waitTask") {
        @Override
        public Object doCall() throws Exception {
            Thread.currentThread().sleep(150);
            return "wait";
        }

    };
    EspTask longestTask = new EspTask("worldTask") {
        @Override
        public Object doCall() throws Exception {
            Thread.currentThread().sleep(300);
            return "world";
        }
    };
    return Arrays.asList(longTask, longerTask, longestTask);
}

【问题讨论】:

    标签: java multithreading rx-java reactive-programming observable


    【解决方案1】:

    您可以使用Observable.onErrorReturn(),并返回特殊值(例如null),然后过滤下游的非特殊值。请记住,源 observable 将在出错时完成。此外,根据用例Observable.onErrorResumeNext()methods 也很有用。如果您对错误通知感兴趣,请使用Observable.materialize(),这会将项目和onError()onComplete() 转换为通知,然后可以通过Notification.getKind() 进行过滤

    编辑。 上面提到的所有运算符都应该在.toObservable().subscribeOn(this.schedulerFactory.get(groupName)); 之后添加,假设 try/catch 不存在。

    【讨论】:

      【解决方案2】:

      你想用mergeDelayError

      public <T> Observable<T> observeManagedAsync(String groupName,List<EspTask<T>> tasks) {
          return Observable.mergeDelayError(Observable
              .from(tasks)
              .map(task -> {
                  try {
                      return new MyCommand(task.getTaskId(),groupName,task).toObservable().subscribeOn(this.schedulerFactory.get(groupName));
                  } catch(Exception ex) {
                      return Observable.error(ex);
                  }
              }));
      }
      

      请注意,您的 MyCommand 构造函数不应抛出任何异常;这可以让你的代码写得更简洁:

      public <T> Observable<T> observeManagedAsync(String groupName,List<EspTask<T>> tasks) {
          return from(tasks)
                 .map(task -> new MyCommand(task.getTaskId(), groupName, task)
                              .toObservable()
                              .subscribeOn(this.schedulerFactory.get(groupName)))
                 .compose(Observable::mergeDelayError);
      

      }

      请记住,这仍然会调用 onError 最多一次;如果您需要显式处理所有错误,请使用 Either&lt;CommandResult, Throwable&gt; 作为返回类型(或处理错误并返回一个空的 observable)。

      【讨论】:

        【解决方案3】:

        使用.materialize() 允许所有排放和错误作为包装通知通过,然后根据需要处理它们:

             .flatMap(task -> {
                    try {
                        return new MyCommand(task.getTaskId(),groupName,task)
                            .toObservable()
                            .subscribeOn(this.schedulerFactory.get(groupName))
                            .materialize();
                    } catch(Exception ex) {
                        return Observable.error(ex).materialize();
                    }
                });
        

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 1970-01-01
          • 2017-10-15
          • 2019-03-20
          • 2017-11-05
          • 1970-01-01
          • 1970-01-01
          • 2011-06-26
          • 2018-03-24
          相关资源
          最近更新 更多