【问题标题】:How can I make this rxjava zip to run in parallel?如何使这个 rxjava zip 并行运行?
【发布时间】:2018-06-20 00:44:17
【问题描述】:

我有一个 sleep 方法来模拟长时间运行的进程。

private void sleep() {
    try {
        Thread.sleep(2000);
    } catch (InterruptedException e) {
        e.printStackTrace();
    }
}

然后我有一个方法返回一个 Observable,其中包含参数中给出的 2 个字符串的列表。它在返回字符串之前调用睡眠。

private Observable<List<String>> getStrings(final String str1, final String str2) {
    return Observable.fromCallable(new Callable<List<String>>() {
        @Override
        public List<String> call() {
            sleep();
            List<String> strings = new ArrayList<>();
            strings.add(str1);
            strings.add(str2);
            return strings;
        }
    });
}

然后我在 Observalb.zip 中调用 getStrings 三次,我希望这三个调用并行运行,所以执行的总时间应该在 2 秒 或最多 3 秒内因为睡眠只有2秒。但是,它总共需要 6 秒。 如何让它并行运行,以便在 2 秒内完成?

Observable
.zip(getStrings("One", "Two"), getStrings("Three", "Four"), getStrings("Five", "Six"), mergeStringLists())
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe(new Observer<List<String>>() {
    @Override
    public void onCompleted() {

    }

    @Override
    public void onError(Throwable e) {

    }

    @Override
    public void onNext(List<String> strings) {
        //Display the strings
    }
});

mergeStringLists 方法

private Func3<List<String>, List<String>, List<String>, List<String>> mergeStringLists() {
    return new Func3<List<String>, List<String>, List<String>, List<String>>() {
        @Override
        public List<String> call(List<String> strings, List<String> strings2, List<String> strings3) {
            Log.d(TAG, "...");

            for (String s : strings2) {
                strings.add(s);
            }

            for (String s : strings3) {
                strings.add(s);
            }

            return strings;
        }
    };
}

【问题讨论】:

  • 你试过用 Observable.combineLatest 代替 Observable.zip 吗?

标签: java android rx-java reactive-programming rx-android


【解决方案1】:

这是因为订阅您的 zipped observable 发生在同一个 io 线程中。

你为什么不试试这个呢:

Observable
    .zip(
        getStrings("One", "Two")
            .subscribeOn(Schedulers.newThread()),
        getStrings("Three", "Four")
            .subscribeOn(Schedulers.newThread()),
        getStrings("Five", "Six")
            .subscribeOn(Schedulers.newThread()),
        mergeStringLists())
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe(new Observer<List<String>>() {
        @Override
        public void onCompleted() {

        }

        @Override
        public void onError(Throwable e) {

        }

        @Override
        public void onNext(List<String> strings) {
            //Display the strings
        }
    });

如果有帮助请告诉我

【讨论】:

  • Schedulers.io() 默认情况下是一个按需增长的线程池。
  • @TassosBassoukos 你的意思是 Schedulers.io() 会根据需要自动创建新线程吗?
  • @Bartek 您的解决方案有效,您知道除了您的解决方案之外是否还有其他解决方案可以让 zip 并行运行?
  • @s-hunter 确实,per the docs
  • @TassosBassoukos,在这种情况下,为什么我的代码是顺序运行而不是并行运行的?
【解决方案2】:

这里有一个我以异步方式使用 Zip 的示例,以防你好奇

/**
 * Since every observable into the zip is created to                 subscribeOn a diferent thread, it´s means all of them will run in parallel.
 * By default Rx is not async, only if you explicitly use subscribeOn.
 */
@Test
public void testAsyncZip() {
    scheduler = Schedulers.newThread();
    scheduler1 = Schedulers.newThread();
    scheduler2 = Schedulers.newThread();
    long start = System.currentTimeMillis();
    Observable.zip(obAsyncString(), obAsyncString1(), obAsyncString2(), (s, s2, s3) -> s.concat(s2).concat(s3))
            .subscribe(result -> showResult("Async in:", start, result));
}

public Observable<String> obAsyncString() {
    return Observable.just("")
            .observeOn(scheduler)
            .doOnNext(val -> System.out.println("Thread " + Thread.currentThread().getName()))
            .map(val -> "Hello");
}

public Observable<String> obAsyncString1() {
    return Observable.just("")
            .observeOn(scheduler1)
            .doOnNext(val -> System.out.println("Thread " + Thread.currentThread().getName()))
            .map(val -> " World");
}

public Observable<String> obAsyncString2() {
    return Observable.just("")
            .observeOn(scheduler2)
            .doOnNext(val -> System.out.println("Thread " + Thread.currentThread().getName()))
            .map(val -> "!");
}

您可以在这里查看更多示例https://github.com/politrons/reactive

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-05-11
    • 1970-01-01
    • 2019-07-21
    • 1970-01-01
    • 1970-01-01
    • 2022-10-05
    相关资源
    最近更新 更多