RxJava实现Observable并行执行仍在主线程串行运行问题求助
问题原因
你代码的串行执行问题出在getPositions方法的逻辑位置:
getPositions方法中的System.out.println、Thread.sleep(500)是方法被调用时,直接在上游发射数据的主线程同步执行的,你把耗时逻辑写在了返回Observable的代码外层,并没有包进要异步执行的任务体里- 你声明的
subscribeOn(Schedulers.computation())只对它上游的Single.fromCallable包裹的代码生效,方法外层的逻辑不受这个线程调度影响
修复代码
把所有需要异步执行的逻辑都放进fromCallable内部即可:
private static Observable<String> getPositions(List<String> id) { // 所有异步逻辑都移入fromCallable包裹的任务体中 return Single.fromCallable(() -> { System.out.println("thread: " + Thread.currentThread().getName()); Thread.sleep(500); if (id.contains("3")) { return Arrays.asList("a", "aa", "aaa"); } else if (id.contains("2")) { return Arrays.asList("b", "bb", "bbb"); } return Arrays.asList("c", "cc", "ccc"); }).flatMapObservable(Observable::fromIterable).subscribeOn(Schedulers.computation()); } public static void main(String[] args) throws InterruptedException { Observable<String> a = Observable.fromIterable(Arrays.asList("1", "2", "3")); a.buffer(2).flatMap(buff -> getPositions(buff), 4).toList().subscribe(val -> System.out.println(val)); Thread.sleep(1500); }
效果说明
修复后运行会看到两个打印的线程都是computation调度池的线程,且两次任务几乎同时执行,总耗时只有500ms左右,符合并行执行的预期。
内容的提问来源于stack exchange,提问作者sorror
相关产品推荐
相关产品推荐

