【问题标题】:Why Parallel execution not happening for multiple RXJava Observables?为什么多个 RXJava Observables 不会发生并行执行?
【发布时间】:2017-11-21 04:42:47
【问题描述】:

我是 RxJava 的新手,正在尝试从链接执行多个 Observable 的并行执行示例: RxJava Fetching Observables In Parallel

虽然上面链接中提供的示例是并行执行 observable,但是当我在 forEach 方法中添加 Thread.sleep(TIME_IN_MILLISECONDS) 时,系统开始一次执行一个 Observable。请帮助我理解为什么 Thread.sleep 会停止 Observables 的并行执行。

以下是导致多个可观察对象同步执行的修改示例:

import rx.Observable;
import rx.Subscriber;
import rx.schedulers.Schedulers;

public class ParallelExecution {

    public static void main(String[] args) {
        System.out.println("------------ mergingAsync");
        mergingAsync();
    }

    private static void mergingAsync() {
        Observable.merge(getDataAsync(1), getDataAsync(2)).toBlocking()
        .forEach(x -> { try{Thread.sleep(4000);}catch(Exception ex){}; 
        System.out.println(x + " " + Thread.currentThread().getId());});
    }

    // artificial representations of IO work
    static Observable<Integer> getDataAsync(int i) {
        return getDataSync(i).subscribeOn(Schedulers.io());
    }

    static Observable<Integer> getDataSync(int i) {
        return Observable.create((Subscriber<? super Integer> s) -> {
            // simulate latency
            try {
                Thread.sleep(1000);
            } catch (Exception e) {
                e.printStackTrace();
            }
            s.onNext(i);
            s.onCompleted();
        });
    }
}

在上面的例子中,我们使用了 observable 的 subscribeOn 方法并提供了一个 ThreadPool(Schedules.io) 用于执行,因此每个 Observable 的订阅将发生在单独的线程上。

Thread.sleep 可能会锁定线程之间的任何共享对象,但我仍然不清楚。请帮忙。

【问题讨论】:

  • 你怎么知道源没有并行运行?
  • 我通过打印当前正在执行的线程 ID 知道这一点。

标签: java jakarta-ee rx-java reactive-programming reactivex


【解决方案1】:

实际上,您的示例并行执行确实发生了,您只是看错了,执行工作的位置和发出通知的位置之间存在差异。

如果你将线程 id 的日志放在Observable.create,你会注意到每个 Observable 同时在不同的线程上执行。但通知是连续发生的。这种行为正如 Observable 合约的一部分所期望的那样,即 observables 必须串行(而不是并行)向观察者发出通知。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2014-12-02
    • 1970-01-01
    • 2010-11-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-08-02
    相关资源
    最近更新 更多