Dav*_*ank 3 multithreading scala reactive-programming system.reactive rx-java
我是ReactiveX库的新手(我使用它的scala变体,RxScala).
我有一个Observable以高速率发出价值的东西.我想将函数应用于Observable(map)的所有值.我使用的函数在map计算上相当昂贵.
有没有办法让线程池map并行计算相位?
是的,有办法做到这一点.
我会将流缓冲到块中并使用cpus分配负载Schedulers.computation()(使用Executor基于大小等于可用处理器数量的线程池):
int chunkSize = 1000;
source
.buffer(chunkSize)
.flatMap(
list ->
Observable
.from(list)
.map(expensive)
.subscribeOn(Schedulers.computation()))
...
Run Code Online (Sandbox Code Playgroud)
如果map操作足够昂贵,那么你可能没有buffer:
source
.flatMap(
x ->
Observable
.just(x)
.map(expensive)
.subscribeOn(Schedulers.computation()))
Run Code Online (Sandbox Code Playgroud)