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