【发布时间】:2015-02-03 10:04:02
【问题描述】:
我需要并行获取一些文件。 get 操作本身是 IO 密集型的,可以从并行执行中受益匪浅。
使用 RxJava,我可以通过使用 Async.toAsync 包装我的函数来实现这一点。
我想知道使用subscribeOn() 或observeOn() 是否有更简洁的方法?我无法弄清楚。尝试了不同的方式,但任何方式都只能使用一个线程,并且处理是按顺序进行的。
import rx.Observable;
import rx.Scheduler;
import rx.functions.Func1;
import rx.schedulers.Schedulers;
import rx.util.async.Async;
import java.util.List;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
public class ParallelMultiGet {
public List<String> readContents(List<String> paths) {
Func1<String, String> getFunction = new Func1<String, String>() {
@Override
public String call(String path) {
return get(path);
}
};
Scheduler scheduler = Schedulers.from(
Executors.newFixedThreadPool(
Math.min(paths.size(), 50)));
Future<List<String>> result = Observable
.from(paths)
.flatMap(
Async.toAsync(getFunction, scheduler))
.toList()
.toBlocking()
.toFuture();
try {
return result.get(30, TimeUnit.MINUTES);
} catch (Exception e) {
// For example if Func1.call above throws an exception it ends up in here
throw new IllegalStateException("Couldn't read paths", e);
}
}
private String get(String path) {
// this would be the slow operation, waiting for IO
return "content";
}
}
即使这已经很不错了,因为我不需要构建自己的循环来提交期货并将它们中的值组合到结果列表中。但也许不必这么冗长?
【问题讨论】:
-
你不需要异步;只是一个普通的
ExecutorService有.invokeAll() -
这里确实值得一提。结果是 List
,因此仍需要至少对其进行迭代并将值映射到列表中。 invokeAll()的整体复杂性较低,但这个问题主要是关于如何使用 RxJava 来完成。很高兴找到一种方便的方法链接方式,并且能够使用Observable.from等而不是“手动”迭代paths来为invokeAll()创建任务。
标签: java asynchronous reactive-programming observable rx-java